MCPcopy Create free account
hub / github.com/apache/datafusion / infer_schema

Method infer_schema

datafusion/datasource-json/src/file_format.rs:255–332  ·  view source on GitHub ↗
(
        &self,
        _state: &dyn Session,
        store: &Arc<dyn ObjectStore>,
        objects: &[ObjectMeta],
    )

Source from the content-addressed store, hash-verified

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 {

Callers

nothing calls this directly

Calls 7

newFunction · 0.85
convert_readMethod · 0.80
readerMethod · 0.80
getMethod · 0.45
as_refMethod · 0.45
pushMethod · 0.45

Tested by

no test coverage detected