1#![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#[derive(Debug)]
31pub enum FileInput {
32 File(std::fs::File),
35 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#[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 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 pub fn format(&self) -> FileFormat {
93 self.format
94 }
95
96 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 #[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 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 #[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#[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#[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#[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 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 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}