Skip to main content

faucet_core/file_format/
container.rs

1//! Multi-file decoding of the self-describing container formats — Avro and
2//! ORC (#719).
3//!
4//! A connector reading a prefix or a directory decodes many files into one
5//! stream, so their shapes must agree. [`ContainerDecoder`] holds that rule in
6//! one place, keyed on the first file it decodes:
7//!
8//! - **Avro** — the configured reader schema, or else the first file's writer
9//!   schema, is the reader schema for every file. A later file written with an
10//!   evolved schema is resolved against it (added fields take their defaults,
11//!   removed ones are dropped); a file that cannot be resolved fails with an
12//!   error naming both files.
13//! - **ORC** — the first file's (projected) Arrow schema is the reference; a
14//!   file whose schema differs fails with an error naming both files.
15//!
16//! Decoding is synchronous. Call it from a blocking context for large files;
17//! both [`FileInput`] variants stream rather than materialize the decoded
18//! rows.
19
20#![cfg_attr(
21    not(any(feature = "file-format-avro", feature = "file-format-orc")),
22    allow(unreachable_code, unused_variables, dead_code)
23)]
24
25use super::{FileFormat, FormatOptions};
26use crate::error::FaucetError;
27use serde_json::Value;
28
29/// Where one container file's bytes come from.
30#[derive(Debug)]
31pub enum FileInput {
32    /// A local file — decoded incrementally (Avro block by block, ORC stripe
33    /// by stripe), so memory is bounded by the chunk, not the file.
34    File(std::fs::File),
35    /// A whole object already in memory (object stores).
36    Bytes(Vec<u8>),
37}
38
39#[derive(Debug)]
40enum Anchor {
41    #[cfg(feature = "file-format-avro")]
42    Avro(apache_avro::Schema),
43    #[cfg(feature = "file-format-orc")]
44    Orc(arrow::datatypes::SchemaRef),
45}
46
47/// Decodes a sequence of Avro or ORC files against one shape.
48#[derive(Debug)]
49pub struct ContainerDecoder {
50    format: FileFormat,
51    #[cfg(feature = "file-format-orc")]
52    opts: FormatOptions,
53    #[cfg(feature = "file-format-avro")]
54    configured: bool,
55    anchor: Option<(String, Anchor)>,
56}
57
58impl ContainerDecoder {
59    /// A decoder for `format`, which must be [`FileFormat::Avro`] or
60    /// [`FileFormat::Orc`] in a build with the matching feature.
61    pub fn new(format: FileFormat, opts: &FormatOptions) -> Result<Self, FaucetError> {
62        let anchor: Option<(String, Anchor)> = match format {
63            #[cfg(feature = "file-format-avro")]
64            FileFormat::Avro => opts
65                .avro
66                .parsed_schema()?
67                .map(|s| ("avro.schema".to_string(), Anchor::Avro(s))),
68            #[cfg(feature = "file-format-orc")]
69            FileFormat::Orc => None,
70            #[cfg(not(feature = "file-format-avro"))]
71            FileFormat::Avro => return Err(super::missing_feature(format, "file-format-avro")),
72            #[cfg(not(feature = "file-format-orc"))]
73            FileFormat::Orc => return Err(super::missing_feature(format, "file-format-orc")),
74            other => {
75                return Err(FaucetError::Config(format!(
76                    "ContainerDecoder handles avro and orc, not `{}`",
77                    other.as_str()
78                )));
79            }
80        };
81        Ok(Self {
82            format,
83            #[cfg(feature = "file-format-orc")]
84            opts: opts.clone(),
85            #[cfg(feature = "file-format-avro")]
86            configured: anchor.is_some(),
87            anchor,
88        })
89    }
90
91    /// The format this decoder reads.
92    pub fn format(&self) -> FileFormat {
93        self.format
94    }
95
96    /// Decode one whole file into its records.
97    pub fn decode_all(&mut self, name: &str, input: FileInput) -> Result<Vec<Value>, FaucetError> {
98        let mut out = Vec::new();
99        self.records(name, input, 0, &mut |c| {
100            out.extend(c);
101            Ok(())
102        })?;
103        Ok(out)
104    }
105
106    /// Decode one whole file into Arrow batches of at most `batch_size` rows.
107    #[cfg(feature = "arrow")]
108    pub fn decode_batches(
109        &mut self,
110        name: &str,
111        input: FileInput,
112        batch_size: usize,
113    ) -> Result<(arrow::datatypes::SchemaRef, Vec<arrow::array::RecordBatch>), FaucetError> {
114        let mut out = Vec::new();
115        let schema = self.batches(name, input, batch_size, &mut |b| {
116            out.push(b);
117            Ok(())
118        })?;
119        Ok((schema, out))
120    }
121
122    /// Decode one file's records in chunks of `chunk` (`0` = the whole file
123    /// as one chunk), calling `f` per chunk.
124    pub fn records(
125        &mut self,
126        name: &str,
127        input: FileInput,
128        chunk: usize,
129        f: &mut dyn FnMut(Vec<Value>) -> Result<(), FaucetError>,
130    ) -> Result<(), FaucetError> {
131        let chunk = if chunk == 0 { usize::MAX } else { chunk };
132        match self.format {
133            #[cfg(feature = "file-format-avro")]
134            FileFormat::Avro => {
135                let reader_schema = self.avro_reader_schema();
136                let result = with_reader(input, |r| {
137                    super::avro::read_records(r, reader_schema.as_ref(), chunk, f)
138                });
139                self.finish_avro(name, reader_schema, result)
140            }
141            #[cfg(feature = "file-format-orc")]
142            FileFormat::Orc => {
143                let batch = if chunk == usize::MAX { 0 } else { chunk };
144                let mut pending: Vec<Value> = Vec::new();
145                self.orc_read(name, input, batch, &mut |b| {
146                    let rows = crate::columnar::record_batch_to_values(&b)?;
147                    if chunk == usize::MAX {
148                        pending.extend(rows);
149                        Ok(())
150                    } else {
151                        f(rows)
152                    }
153                })?;
154                if !pending.is_empty() {
155                    f(pending)?;
156                }
157                Ok(())
158            }
159            _ => unreachable!("rejected in new()"),
160        }
161    }
162
163    /// Decode one file as Arrow batches of at most `batch_size` rows (`0` =
164    /// the reader's natural size), calling `f` per batch. Returns the file's
165    /// Arrow schema, identical for every file this decoder accepts.
166    #[cfg(feature = "arrow")]
167    pub fn batches(
168        &mut self,
169        name: &str,
170        input: FileInput,
171        batch_size: usize,
172        f: &mut dyn FnMut(arrow::array::RecordBatch) -> Result<(), FaucetError>,
173    ) -> Result<arrow::datatypes::SchemaRef, FaucetError> {
174        match self.format {
175            #[cfg(feature = "file-format-avro")]
176            FileFormat::Avro => {
177                let reader_schema = self.avro_reader_schema();
178                let result = with_reader(input, |r| {
179                    super::avro::read_batches(r, reader_schema.as_ref(), batch_size, f)
180                });
181                let arrow = result.as_ref().ok().map(|(_, a)| a.clone());
182                self.finish_avro(name, reader_schema, result.map(|(s, _)| s))?;
183                Ok(arrow.expect("finish_avro succeeded only on Ok"))
184            }
185            #[cfg(feature = "file-format-orc")]
186            FileFormat::Orc => self.orc_read(name, input, batch_size, f),
187            _ => unreachable!("rejected in new()"),
188        }
189    }
190
191    #[cfg(feature = "file-format-avro")]
192    fn avro_reader_schema(&self) -> Option<apache_avro::Schema> {
193        match &self.anchor {
194            Some((_, Anchor::Avro(s))) => Some(s.clone()),
195            _ => None,
196        }
197    }
198
199    #[cfg(feature = "file-format-avro")]
200    fn finish_avro(
201        &mut self,
202        name: &str,
203        reader_schema: Option<apache_avro::Schema>,
204        result: Result<apache_avro::Schema, FaucetError>,
205    ) -> Result<(), FaucetError> {
206        match (result, &self.anchor) {
207            (Ok(schema), None) => {
208                self.anchor = Some((name.to_string(), Anchor::Avro(schema)));
209                Ok(())
210            }
211            (Ok(_), Some(_)) => Ok(()),
212            (Err(e), Some((first, _))) if reader_schema.is_some() => {
213                let against = if self.configured {
214                    "the configured `avro.schema`".to_string()
215                } else {
216                    format!("'{first}' (the first file's schema)")
217                };
218                Err(FaucetError::Source(format!(
219                    "avro schema of '{name}' cannot be resolved against {against}: {e}"
220                )))
221            }
222            (Err(e), _) => Err(FaucetError::Source(format!("'{name}': {e}"))),
223        }
224    }
225
226    #[cfg(feature = "file-format-orc")]
227    fn orc_read(
228        &mut self,
229        name: &str,
230        input: FileInput,
231        batch_size: usize,
232        f: &mut dyn FnMut(arrow::array::RecordBatch) -> Result<(), FaucetError>,
233    ) -> Result<arrow::datatypes::SchemaRef, FaucetError> {
234        let input = match input {
235            FileInput::File(file) => super::orc::OrcInput::File(file),
236            FileInput::Bytes(b) => super::orc::OrcInput::Bytes(bytes::Bytes::from(b)),
237        };
238        let reference = match &self.anchor {
239            Some((first, Anchor::Orc(s))) => Some((first.clone(), s.clone())),
240            _ => None,
241        };
242        let mut check = |schema: &arrow::datatypes::SchemaRef| -> Result<(), FaucetError> {
243            match &reference {
244                Some((first, reference)) if !same_shape(reference, schema) => {
245                    Err(schema_conflict(name, first, reference, schema))
246                }
247                _ => Ok(()),
248            }
249        };
250        let schema =
251            super::orc::read_batches_checked(input, &self.opts.orc, batch_size, &mut check, f)
252                .map_err(|e| prefix(name, e))?;
253        if self.anchor.is_none() {
254            self.anchor = Some((name.to_string(), Anchor::Orc(schema.clone())));
255        }
256        Ok(schema)
257    }
258}
259
260/// The columnar page stream for a list of container objects: fetch each with
261/// an ordered `concurrency`-wide look-ahead, decode in listing order, and
262/// yield one [`ColumnarPage`](crate::columnar::ColumnarPage) per non-empty
263/// batch. The shared body of every object-store source's `stream_batches`
264/// for Avro and ORC.
265#[cfg(feature = "arrow")]
266pub fn columnar_pages<'a, F, Fut>(
267    names: Vec<String>,
268    concurrency: usize,
269    mut decoder: ContainerDecoder,
270    batch_size: usize,
271    fetch: F,
272) -> std::pin::Pin<
273    Box<dyn futures::Stream<Item = Result<crate::columnar::ColumnarPage, FaucetError>> + Send + 'a>,
274>
275where
276    F: Fn(String) -> Fut + Send + Sync + 'a,
277    Fut: std::future::Future<Output = Result<Vec<u8>, FaucetError>> + Send + 'a,
278{
279    use futures::StreamExt as _;
280    Box::pin(async_stream::try_stream! {
281        let fetch = &fetch;
282        let mut fetched = futures::stream::iter(names)
283            .map(|name| async move {
284                let body = fetch(name.clone()).await;
285                (name, body)
286            })
287            .buffered(concurrency.max(1));
288        while let Some((name, body)) = fetched.next().await {
289            let (_, batches) = decoder.decode_batches(&name, FileInput::Bytes(body?), batch_size)?;
290            for batch in batches {
291                yield crate::columnar::ColumnarPage::new(batch, None);
292            }
293        }
294    })
295}
296
297#[cfg(feature = "file-format-avro")]
298fn with_reader<T>(
299    input: FileInput,
300    f: impl FnOnce(&mut dyn std::io::Read) -> Result<T, FaucetError>,
301) -> Result<T, FaucetError> {
302    match input {
303        FileInput::File(file) => f(&mut std::io::BufReader::new(file)),
304        FileInput::Bytes(b) => f(&mut &b[..]),
305    }
306}
307
308#[cfg(feature = "file-format-orc")]
309fn prefix(name: &str, e: FaucetError) -> FaucetError {
310    match e {
311        FaucetError::Source(m) if !m.starts_with('\'') => {
312            FaucetError::Source(format!("'{name}': {m}"))
313        }
314        other => other,
315    }
316}
317
318/// Whether two schemas carry the same columns: names and types in order.
319/// Nullability and metadata are not shape — a writer marking a column
320/// nullable in one file and not in the next still produces the same rows.
321#[cfg(feature = "arrow")]
322pub fn same_shape(a: &arrow::datatypes::Schema, b: &arrow::datatypes::Schema) -> bool {
323    a.fields().len() == b.fields().len()
324        && a.fields()
325            .iter()
326            .zip(b.fields().iter())
327            .all(|(x, y)| x.name() == y.name() && x.data_type() == y.data_type())
328}
329
330/// The error for two files whose shapes disagree, naming both and the first
331/// field that differs.
332#[cfg(feature = "arrow")]
333pub fn schema_conflict(
334    file: &str,
335    first: &str,
336    reference: &arrow::datatypes::Schema,
337    schema: &arrow::datatypes::Schema,
338) -> FaucetError {
339    let detail = reference
340        .fields()
341        .iter()
342        .zip(schema.fields().iter())
343        .find(|(a, b)| a.name() != b.name() || a.data_type() != b.data_type())
344        .map(|(a, b)| {
345            format!(
346                "field `{}` ({}) vs `{}` ({})",
347                a.name(),
348                a.data_type(),
349                b.name(),
350                b.data_type()
351            )
352        })
353        .unwrap_or_else(|| {
354            format!(
355                "{} vs {} fields",
356                reference.fields().len(),
357                schema.fields().len()
358            )
359        });
360    FaucetError::Source(format!(
361        "schema of '{file}' conflicts with '{first}' (the first file's schema): {detail}"
362    ))
363}
364
365#[cfg(test)]
366mod tests {
367    use super::*;
368    #[cfg(feature = "file-format-avro")]
369    use serde_json::json;
370
371    #[cfg(feature = "file-format-avro")]
372    fn avro(records: &[Value]) -> Vec<u8> {
373        super::super::avro::encode(records, &Default::default()).expect("encode")
374    }
375
376    #[cfg(all(feature = "file-format-avro", feature = "arrow"))]
377    #[tokio::test]
378    async fn columnar_pages_decode_in_listing_order() {
379        use futures::StreamExt as _;
380        let bodies: std::collections::HashMap<String, Vec<u8>> = [
381            ("a".to_string(), avro(&[json!({"id": 1}), json!({"id": 2})])),
382            ("b".to_string(), avro(&[json!({"id": 3})])),
383        ]
384        .into_iter()
385        .collect();
386        let decoder = ContainerDecoder::new(FileFormat::Avro, &FormatOptions::default()).unwrap();
387        let pages: Vec<_> = columnar_pages(
388            vec!["a".into(), "b".into(), "missing".into()],
389            2,
390            decoder,
391            1,
392            |n| {
393                let body = bodies.get(&n).cloned();
394                async move { body.ok_or_else(|| FaucetError::Source(format!("no {n}"))) }
395            },
396        )
397        .collect()
398        .await;
399        assert_eq!(pages.len(), 4);
400        assert_eq!(
401            pages[..3]
402                .iter()
403                .map(|p| p.as_ref().unwrap().num_rows())
404                .sum::<usize>(),
405            3
406        );
407        assert!(pages[3].as_ref().is_err());
408    }
409
410    #[test]
411    fn only_container_formats_are_accepted() {
412        let err =
413            ContainerDecoder::new(FileFormat::Csv, &FormatOptions::default()).expect_err("csv");
414        assert!(err.to_string().contains("avro and orc"), "{err}");
415    }
416
417    #[cfg(feature = "file-format-avro")]
418    #[test]
419    fn later_avro_files_resolve_against_the_first() {
420        let mut d = ContainerDecoder::new(FileFormat::Avro, &FormatOptions::default()).unwrap();
421        assert_eq!(d.format(), FileFormat::Avro);
422        let mut rows = Vec::new();
423        d.records(
424            "a.avro",
425            FileInput::Bytes(avro(&[json!({"id": 1, "x": "a"})])),
426            0,
427            &mut |c| {
428                rows.extend(c);
429                Ok(())
430            },
431        )
432        .unwrap();
433        // Same field names, one dropped: resolvable, projected onto the first shape.
434        let wider = avro(&[json!({"id": 2, "x": "b", "extra": true})]);
435        d.records("b.avro", FileInput::Bytes(wider), 1, &mut |c| {
436            rows.extend(c);
437            Ok(())
438        })
439        .unwrap();
440        assert_eq!(
441            rows,
442            vec![json!({"id": 1, "x": "a"}), json!({"id": 2, "x": "b"})]
443        );
444        let err = d
445            .records(
446                "c.avro",
447                FileInput::Bytes(avro(&[json!({"id": "text"})])),
448                0,
449                &mut |_| Ok(()),
450            )
451            .expect_err("conflict");
452        let msg = err.to_string();
453        assert!(msg.contains("c.avro") && msg.contains("a.avro"), "{msg}");
454    }
455
456    #[cfg(feature = "file-format-avro")]
457    #[test]
458    fn a_configured_reader_schema_is_named_in_conflicts() {
459        let opts = FormatOptions {
460            avro: super::super::AvroOptions {
461                schema: Some(
462                    json!({"type": "record", "name": "faucet_record", "fields": [
463                        {"name": "id", "type": "long"}
464                    ]}),
465                ),
466                ..Default::default()
467            },
468            ..Default::default()
469        };
470        let mut d = ContainerDecoder::new(FileFormat::Avro, &opts).unwrap();
471        let err = d
472            .records(
473                "x.avro",
474                FileInput::Bytes(avro(&[json!({"other": 1})])),
475                0,
476                &mut |_| Ok(()),
477            )
478            .expect_err("unresolvable");
479        assert!(err.to_string().contains("configured"), "{err}");
480        let err = d
481            .records("y.avro", FileInput::Bytes(b"junk".to_vec()), 0, &mut |_| {
482                Ok(())
483            })
484            .expect_err("junk");
485        assert!(err.to_string().contains("y.avro"), "{err}");
486    }
487
488    #[cfg(all(feature = "file-format-avro", feature = "arrow"))]
489    #[test]
490    fn avro_batches_and_local_files() {
491        let dir = tempfile::tempdir().unwrap();
492        let path = dir.path().join("a.avro");
493        std::fs::write(&path, avro(&[json!({"id": 1}), json!({"id": 2})])).unwrap();
494        let mut d = ContainerDecoder::new(FileFormat::Avro, &FormatOptions::default()).unwrap();
495        let mut n = 0;
496        let schema = d
497            .batches(
498                "a.avro",
499                FileInput::File(std::fs::File::open(&path).unwrap()),
500                1,
501                &mut |b| {
502                    n += b.num_rows();
503                    Ok(())
504                },
505            )
506            .unwrap();
507        assert_eq!(n, 2);
508        assert_eq!(schema.field(0).name(), "id");
509    }
510
511    #[cfg(feature = "file-format-orc")]
512    #[test]
513    fn orc_files_must_share_a_schema() {
514        const FIXTURE: &[u8] = include_bytes!("../../tests/fixtures/orc/people.orc");
515        let mut d = ContainerDecoder::new(FileFormat::Orc, &FormatOptions::default()).unwrap();
516        let mut rows = Vec::new();
517        d.records("a.orc", FileInput::Bytes(FIXTURE.to_vec()), 0, &mut |c| {
518            rows.extend(c);
519            Ok(())
520        })
521        .unwrap();
522        assert_eq!(rows.len(), 3);
523        let mut chunks = 0;
524        d.records("b.orc", FileInput::Bytes(FIXTURE.to_vec()), 2, &mut |_| {
525            chunks += 1;
526            Ok(())
527        })
528        .unwrap();
529        assert_eq!(chunks, 2);
530        let mut projected = ContainerDecoder::new(
531            FileFormat::Orc,
532            &FormatOptions {
533                orc: super::super::OrcOptions {
534                    columns: Some(vec!["id".into()]),
535                },
536                ..Default::default()
537            },
538        )
539        .unwrap();
540        let s = projected
541            .batches("p.orc", FileInput::Bytes(FIXTURE.to_vec()), 0, &mut |_| {
542                Ok(())
543            })
544            .unwrap();
545        // A second decoder anchored on the full schema refuses the projected one.
546        let full = d.anchor.as_ref().map(|(_, a)| match a {
547            Anchor::Orc(s) => s.clone(),
548            #[allow(unreachable_patterns)]
549            _ => unreachable!(),
550        });
551        let err = schema_conflict("p.orc", "a.orc", &full.unwrap(), &s);
552        assert!(err.to_string().contains("p.orc") && err.to_string().contains("a.orc"));
553        d.anchor = Some(("a.orc".into(), Anchor::Orc(s)));
554        let err = d
555            .records("c.orc", FileInput::Bytes(FIXTURE.to_vec()), 0, &mut |_| {
556                Ok(())
557            })
558            .expect_err("conflict");
559        assert!(err.to_string().contains("c.orc"), "{err}");
560        let err = d
561            .records("d.orc", FileInput::Bytes(b"junk".to_vec()), 0, &mut |_| {
562                Ok(())
563            })
564            .expect_err("junk");
565        assert!(err.to_string().contains("d.orc"), "{err}");
566    }
567
568    #[cfg(feature = "arrow")]
569    #[test]
570    fn schema_conflict_reports_a_field_count_difference() {
571        use arrow::datatypes::{DataType, Field, Schema};
572        let a = Schema::new(vec![Field::new("x", DataType::Int64, true)]);
573        let b = Schema::new(vec![
574            Field::new("x", DataType::Int64, true),
575            Field::new("y", DataType::Int64, true),
576        ]);
577        let msg = schema_conflict("b", "a", &a, &b).to_string();
578        assert!(msg.contains("1 vs 2"), "{msg}");
579    }
580}