(
&self,
_state: &dyn Session,
store: &Arc<dyn ObjectStore>,
objects: &[ObjectMeta],
)
| 253 | } |
| 254 | |
| 255 | async fn infer_schema( |
| 256 | &self, |
| 257 | _state: &dyn Session, |
| 258 | store: &Arc<dyn ObjectStore>, |
| 259 | objects: &[ObjectMeta], |
| 260 | ) -> Result<SchemaRef> { |
| 261 | let mut schemas = Vec::new(); |
| 262 | let mut records_to_read = self |
| 263 | .options |
| 264 | .schema_infer_max_rec |
| 265 | .unwrap_or(DEFAULT_SCHEMA_INFER_MAX_RECORD); |
| 266 | let file_compression_type = FileCompressionType::from(self.options.compression); |
| 267 | let newline_delimited = self.options.newline_delimited; |
| 268 | |
| 269 | for object in objects { |
| 270 | // Early exit if we've read enough records |
| 271 | if records_to_read == 0 { |
| 272 | break; |
| 273 | } |
| 274 | |
| 275 | let r = store.as_ref().get(&object.location).await?; |
| 276 | |
| 277 | let (schema, records_consumed) = match r.payload { |
| 278 | #[cfg(not(target_arch = "wasm32"))] |
| 279 | GetResultPayload::File(file, _) => { |
| 280 | let decoder = file_compression_type.convert_read(file)?; |
| 281 | let reader = BufReader::new(decoder); |
| 282 | |
| 283 | if newline_delimited { |
| 284 | // NDJSON: use ValueIter directly |
| 285 | let iter = ValueIter::new(reader, None); |
| 286 | let mut count = 0; |
| 287 | let schema = |
| 288 | infer_json_schema_from_iterator(iter.take_while(|_| { |
| 289 | let should_take = count < records_to_read; |
| 290 | if should_take { |
| 291 | count += 1; |
| 292 | } |
| 293 | should_take |
| 294 | }))?; |
| 295 | (schema, count) |
| 296 | } else { |
| 297 | // JSON array format: use streaming converter |
| 298 | infer_schema_from_json_array(reader, records_to_read)? |
| 299 | } |
| 300 | } |
| 301 | GetResultPayload::Stream(_) => { |
| 302 | let data = r.bytes().await?; |
| 303 | let decoder = file_compression_type.convert_read(data.reader())?; |
| 304 | let reader = BufReader::new(decoder); |
| 305 | |
| 306 | if newline_delimited { |
| 307 | let iter = ValueIter::new(reader, None); |
| 308 | let mut count = 0; |
| 309 | let schema = |
| 310 | infer_json_schema_from_iterator(iter.take_while(|_| { |
| 311 | let should_take = count < records_to_read; |
| 312 | if should_take { |
nothing calls this directly
no test coverage detected