1use std::collections::HashMap;
8use std::path::PathBuf;
9use std::pin::Pin;
10use std::sync::Arc;
11
12use async_trait::async_trait;
13use faucet_core::{FaucetError, Stream, StreamPage};
14use futures::{StreamExt, TryStreamExt, stream};
15use object_store::ObjectStore;
16use object_store::aws::AmazonS3Builder;
17use object_store::path::Path as ObjectPath;
18use parquet::arrow::ProjectionMask;
19use parquet::arrow::async_reader::{ParquetObjectReader, ParquetRecordBatchStreamBuilder};
20use serde_json::Value;
21
22use crate::config::{ParquetLocation, ParquetS3Config, ParquetSourceConfig};
23use crate::convert::record_batch_to_json;
24
25pub struct ParquetSource {
27 config: ParquetSourceConfig,
28 s3_store: Option<Arc<dyn ObjectStore>>,
31}
32
33impl ParquetSource {
34 pub async fn new(config: ParquetSourceConfig) -> Result<Self, FaucetError> {
40 if config.concurrency == 0 {
44 return Err(FaucetError::Config(
45 "parquet source: concurrency must be > 0".into(),
46 ));
47 }
48
49 let s3_store = match &config.source {
50 ParquetLocation::S3(s3) => Some(build_s3_store(s3)?),
51 _ => None,
52 };
53
54 Ok(Self { config, s3_store })
55 }
56
57 async fn resolve_files(
62 &self,
63 context: &HashMap<String, Value>,
64 ) -> Result<Vec<FileTarget>, FaucetError> {
65 match &self.config.source {
66 ParquetLocation::LocalPath { path } => {
67 let resolved = substitute(path, context);
68 Ok(vec![FileTarget::Local(PathBuf::from(resolved))])
69 }
70 ParquetLocation::Glob { pattern } => {
71 let resolved = substitute(pattern, context);
72 expand_glob(&resolved)
73 }
74 ParquetLocation::S3(s3) => self.resolve_s3_files(s3, context).await,
75 }
76 }
77
78 async fn resolve_s3_files(
79 &self,
80 s3: &ParquetS3Config,
81 context: &HashMap<String, Value>,
82 ) -> Result<Vec<FileTarget>, FaucetError> {
83 match (&s3.key, &s3.prefix) {
84 (Some(_), Some(_)) => Err(FaucetError::Config(
85 "parquet source: S3 config cannot set both `key` and `prefix`".into(),
86 )),
87 (None, None) => Err(FaucetError::Config(
88 "parquet source: S3 config requires one of `key` or `prefix`".into(),
89 )),
90 (Some(key), None) => {
91 let key = substitute(key, context);
92 Ok(vec![FileTarget::S3(ObjectPath::from(key))])
93 }
94 (None, Some(prefix)) => {
95 let prefix = substitute(prefix, context);
96 let store = self.s3_store.as_ref().ok_or_else(|| {
97 FaucetError::Source("parquet source: S3 store not initialised".into())
98 })?;
99 list_s3_prefix(store.as_ref(), &prefix).await
100 }
101 }
102 }
103
104 async fn read_file(&self, target: &FileTarget) -> Result<FileOutput, FaucetError> {
108 let display = target.display();
109 match target {
110 FileTarget::Local(path) => {
111 let file = tokio::fs::File::open(path).await.map_err(|e| {
112 FaucetError::Source(format!("failed to open parquet file '{display}': {e}"))
113 })?;
114 self.decode(file, &display).await
115 }
116 FileTarget::S3(path) => {
117 let store = self.s3_store.as_ref().ok_or_else(|| {
118 FaucetError::Source("parquet source: S3 store not initialised".into())
119 })?;
120 let reader = ParquetObjectReader::new(store.clone(), path.clone());
121 self.decode(reader, &display).await
122 }
123 }
124 }
125
126 async fn decode<R>(&self, reader: R, display: &str) -> Result<FileOutput, FaucetError>
127 where
128 R: parquet::arrow::async_reader::AsyncFileReader + Send + Unpin + 'static,
129 {
130 let (mut batches, arrow_schema) = self.build_batch_stream(reader, display).await?;
131
132 let mut rows: Vec<Value> = Vec::new();
133 while let Some(batch) = batches.next().await {
134 let batch = batch.map_err(|e| {
135 FaucetError::Source(format!("parquet decode error in '{display}': {e}"))
136 })?;
137 let batch_rows = record_batch_to_json(&batch)?;
138 rows.extend(batch_rows);
139 }
140
141 Ok(FileOutput {
142 path: display.to_string(),
143 rows,
144 arrow_schema,
145 })
146 }
147
148 async fn build_batch_stream<R>(
158 &self,
159 reader: R,
160 display: &str,
161 ) -> Result<(BatchStream, arrow::datatypes::SchemaRef), FaucetError>
162 where
163 R: parquet::arrow::async_reader::AsyncFileReader + Send + Unpin + 'static,
164 {
165 let mut builder = ParquetRecordBatchStreamBuilder::new(reader)
166 .await
167 .map_err(|e| {
168 FaucetError::Source(format!(
169 "failed to read parquet metadata for '{display}': {e}"
170 ))
171 })?;
172
173 if self.config.batch_size > 0 {
178 builder = builder.with_batch_size(self.config.batch_size);
179 }
180
181 if let Some(cols) = self.config.columns.as_deref() {
182 let parquet_schema = builder.parquet_schema();
183 validate_projection(cols, parquet_schema, display)?;
184 let mask = ProjectionMask::columns(parquet_schema, cols.iter().map(String::as_str));
185 builder = builder.with_projection(mask);
186 }
187
188 let arrow_schema = builder.schema().clone();
189
190 let stream = builder.build().map_err(|e| {
191 FaucetError::Source(format!(
192 "failed to build parquet stream for '{display}': {e}"
193 ))
194 })?;
195
196 Ok((Box::pin(stream), arrow_schema))
197 }
198
199 async fn open_target_stream(
203 &self,
204 target: &FileTarget,
205 ) -> Result<(BatchStream, arrow::datatypes::SchemaRef, String), FaucetError> {
206 let display = target.display();
207 match target {
208 FileTarget::Local(path) => {
209 let file = tokio::fs::File::open(path).await.map_err(|e| {
210 FaucetError::Source(format!("failed to open parquet file '{display}': {e}"))
211 })?;
212 let (stream, schema) = self.build_batch_stream(file, &display).await?;
213 Ok((stream, schema, display))
214 }
215 FileTarget::S3(path) => {
216 let store = self.s3_store.as_ref().ok_or_else(|| {
217 FaucetError::Source("parquet source: S3 store not initialised".into())
218 })?;
219 let reader = ParquetObjectReader::new(store.clone(), path.clone());
220 let (stream, schema) = self.build_batch_stream(reader, &display).await?;
221 Ok((stream, schema, display))
222 }
223 }
224 }
225}
226
227type BatchStream =
230 Pin<Box<dyn futures::Stream<Item = parquet::errors::Result<arrow::array::RecordBatch>> + Send>>;
231
232#[async_trait]
233impl faucet_core::Source for ParquetSource {
234 async fn fetch_with_context(
235 &self,
236 context: &HashMap<String, Value>,
237 ) -> Result<Vec<Value>, FaucetError> {
238 let targets = self.resolve_files(context).await?;
239
240 tracing::info!(files = targets.len(), "Parquet source resolved files");
241
242 if targets.is_empty() {
243 return Ok(Vec::new());
244 }
245
246 let concurrency = self.config.concurrency.max(1);
247
248 let outputs: Vec<FileOutput> = stream::iter(targets)
249 .map(|target| async move {
250 let out = self.read_file(&target).await?;
251 tracing::debug!(file = %out.path, rows = out.rows.len(), "Parquet file decoded");
252 Ok::<FileOutput, FaucetError>(out)
253 })
254 .buffer_unordered(concurrency)
255 .try_collect()
256 .await?;
257
258 if outputs.len() > 1 {
259 let first = &outputs[0];
260 for other in &outputs[1..] {
261 if first.arrow_schema != other.arrow_schema {
262 return Err(FaucetError::Source(schema_mismatch_message(first, other)));
263 }
264 }
265 }
266
267 let total: usize = outputs.iter().map(|o| o.rows.len()).sum();
268 let mut all = Vec::with_capacity(total);
269 for out in outputs {
270 all.extend(out.rows);
271 }
272
273 tracing::info!(total_records = all.len(), "Parquet source fetch complete");
274 Ok(all)
275 }
276
277 fn stream_pages<'a>(
304 &'a self,
305 context: &'a HashMap<String, Value>,
306 _batch_size: usize,
307 ) -> Pin<Box<dyn Stream<Item = Result<StreamPage, FaucetError>> + Send + 'a>> {
308 Box::pin(async_stream::try_stream! {
309 let targets = self.resolve_files(context).await?;
310 tracing::info!(files = targets.len(), "Parquet source resolved files");
311
312 if targets.is_empty() {
313 return;
314 }
315
316 let mut reference: Option<(String, arrow::datatypes::SchemaRef)> = None;
325 for target in &targets {
326 let (_, arrow_schema, display) = self.open_target_stream(target).await?;
327 if let Some((first_path, first_schema)) = &reference {
328 if first_schema != &arrow_schema {
329 Err(FaucetError::Source(schema_mismatch_message_pair(
330 first_path,
331 first_schema,
332 &display,
333 &arrow_schema,
334 )))?;
335 }
336 } else {
337 reference = Some((display, arrow_schema));
338 }
339 }
340
341 let mut total_records = 0usize;
343 let mut total_pages = 0usize;
344 for target in &targets {
345 let (mut batches, _schema, display) = self.open_target_stream(target).await?;
346 while let Some(batch) = batches.next().await {
347 let batch = batch.map_err(|e| {
348 FaucetError::Source(format!(
349 "parquet decode error in '{display}': {e}"
350 ))
351 })?;
352 let rows = record_batch_to_json(&batch)?;
353 if rows.is_empty() {
354 continue;
355 }
356 total_records += rows.len();
357 total_pages += 1;
358 yield StreamPage { records: rows, bookmark: None };
359 }
360 }
361
362 tracing::info!(
363 pages = total_pages,
364 total_records,
365 batch_size = self.config.batch_size,
366 "Parquet source stream complete",
367 );
368 })
369 }
370
371 fn config_schema(&self) -> Value {
372 serde_json::to_value(faucet_core::schema_for!(ParquetSourceConfig))
373 .expect("schema serialization")
374 }
375
376 fn dataset_uri(&self) -> String {
377 use crate::config::ParquetLocation;
378 match &self.config.source {
379 ParquetLocation::LocalPath { path } => format!("file://{path}"),
380 ParquetLocation::Glob { pattern } => format!("file://{pattern}"),
381 ParquetLocation::S3(s3) => match (&s3.key, &s3.prefix) {
382 (Some(k), _) => format!("s3://{}/{}", s3.bucket, k),
383 (_, Some(p)) => format!("s3://{}/{}", s3.bucket, p),
384 _ => format!("s3://{}", s3.bucket),
385 },
386 }
387 }
388}
389
390struct FileOutput {
393 path: String,
394 rows: Vec<Value>,
395 arrow_schema: arrow::datatypes::SchemaRef,
396}
397
398#[derive(Debug, Clone)]
400enum FileTarget {
401 Local(PathBuf),
402 S3(ObjectPath),
403}
404
405impl FileTarget {
406 fn display(&self) -> String {
407 match self {
408 FileTarget::Local(p) => p.display().to_string(),
409 FileTarget::S3(p) => format!("s3://{p}"),
410 }
411 }
412}
413
414fn substitute(template: &str, context: &HashMap<String, Value>) -> String {
416 if context.is_empty() {
417 template.to_string()
418 } else {
419 faucet_core::util::substitute_context(template, context)
420 }
421}
422
423fn expand_glob(pattern: &str) -> Result<Vec<FileTarget>, FaucetError> {
425 let entries = glob::glob(pattern)
426 .map_err(|e| FaucetError::Config(format!("invalid glob '{pattern}': {e}")))?;
427
428 let mut paths = Vec::new();
429 for entry in entries {
430 let p = entry
431 .map_err(|e| FaucetError::Source(format!("glob entry error for '{pattern}': {e}")))?;
432 if p.is_file() {
433 paths.push(p);
434 }
435 }
436 paths.sort();
437 Ok(paths.into_iter().map(FileTarget::Local).collect())
438}
439
440async fn list_s3_prefix(
442 store: &dyn ObjectStore,
443 prefix: &str,
444) -> Result<Vec<FileTarget>, FaucetError> {
445 let prefix_path = if prefix.is_empty() {
446 None
447 } else {
448 Some(ObjectPath::from(prefix))
449 };
450
451 let mut listing = store.list(prefix_path.as_ref());
452 let mut keys = Vec::new();
453 while let Some(item) = listing.next().await {
454 let meta = item.map_err(|e| {
455 FaucetError::Source(format!("S3 list error for prefix '{prefix}': {e}"))
456 })?;
457 keys.push(meta.location);
458 }
459 keys.sort();
460 Ok(keys.into_iter().map(FileTarget::S3).collect())
461}
462
463fn build_s3_store(s3: &ParquetS3Config) -> Result<Arc<dyn ObjectStore>, FaucetError> {
465 if s3.bucket.trim().is_empty() {
466 return Err(FaucetError::Config(
467 "parquet source: S3 bucket must not be empty".into(),
468 ));
469 }
470
471 let mut builder = AmazonS3Builder::from_env().with_bucket_name(&s3.bucket);
472 if let Some(region) = &s3.region {
473 builder = builder.with_region(region);
474 }
475 if let Some(endpoint) = &s3.endpoint_url {
476 builder = builder.with_endpoint(endpoint);
477 if endpoint.starts_with("http://") {
478 builder = builder.with_allow_http(true);
479 }
480 }
481
482 let store = builder
483 .build()
484 .map_err(|e| FaucetError::Config(format!("failed to build S3 client: {e}")))?;
485 Ok(Arc::new(store))
486}
487
488fn validate_projection(
492 requested: &[String],
493 parquet_schema: &parquet::schema::types::SchemaDescriptor,
494 display: &str,
495) -> Result<(), FaucetError> {
496 let root = parquet_schema.root_schema();
497 let parquet::schema::types::Type::GroupType { fields, .. } = root else {
498 return Err(FaucetError::Source(format!(
499 "parquet root schema for '{display}' is not a group"
500 )));
501 };
502
503 let known: std::collections::HashSet<&str> = fields.iter().map(|f| f.name()).collect();
504
505 for name in requested {
506 if !known.contains(name.as_str()) {
507 return Err(FaucetError::Source(format!(
508 "parquet source: projected column '{name}' not found in file '{display}' \
509 (available: {})",
510 known.iter().copied().collect::<Vec<_>>().join(", ")
511 )));
512 }
513 }
514
515 Ok(())
516}
517
518fn schema_mismatch_message(first: &FileOutput, other: &FileOutput) -> String {
520 schema_mismatch_message_pair(
521 &first.path,
522 &first.arrow_schema,
523 &other.path,
524 &other.arrow_schema,
525 )
526}
527
528fn schema_mismatch_message_pair(
532 first_path: &str,
533 first_schema: &arrow::datatypes::SchemaRef,
534 other_path: &str,
535 other_schema: &arrow::datatypes::SchemaRef,
536) -> String {
537 let first_fields: Vec<String> = first_schema
538 .fields()
539 .iter()
540 .map(|f| format!("{}:{}", f.name(), f.data_type()))
541 .collect();
542 let other_fields: Vec<String> = other_schema
543 .fields()
544 .iter()
545 .map(|f| format!("{}:{}", f.name(), f.data_type()))
546 .collect();
547
548 let max_len = first_fields.len().max(other_fields.len());
550 let mut first_diff = None;
551 for i in 0..max_len {
552 let a = first_fields
553 .get(i)
554 .map(String::as_str)
555 .unwrap_or("<missing>");
556 let b = other_fields
557 .get(i)
558 .map(String::as_str)
559 .unwrap_or("<missing>");
560 if a != b {
561 first_diff = Some((i, a.to_string(), b.to_string()));
562 break;
563 }
564 }
565
566 let detail = match first_diff {
567 Some((i, a, b)) => format!(" (field #{i}: '{a}' vs '{b}')"),
568 None => String::new(),
569 };
570
571 format!("parquet source: schema mismatch between '{first_path}' and '{other_path}'{detail}")
572}
573
574#[cfg(test)]
575mod tests {
576 use super::*;
577 use crate::config::ParquetSourceConfig;
578 use faucet_core::Source;
579
580 #[test]
581 fn substitute_passes_through_when_context_empty() {
582 let ctx = HashMap::new();
583 assert_eq!(substitute("/tmp/{x}.parquet", &ctx), "/tmp/{x}.parquet");
584 }
585
586 #[test]
587 fn substitute_replaces_placeholders() {
588 let mut ctx = HashMap::new();
589 ctx.insert("region".to_string(), Value::String("us".into()));
590 assert_eq!(
591 substitute("data/{region}/x.parquet", &ctx),
592 "data/us/x.parquet"
593 );
594 }
595
596 #[tokio::test]
597 async fn accepts_zero_batch_size_as_sentinel() {
598 let cfg = ParquetSourceConfig::local("/tmp/x.parquet").batch_size(0);
602 let source = ParquetSource::new(cfg)
603 .await
604 .expect("batch_size=0 must be accepted as the no-batching sentinel");
605 assert_eq!(source.config.batch_size, 0);
606 }
607
608 #[tokio::test]
609 async fn rejects_zero_concurrency() {
610 let cfg = ParquetSourceConfig::local("/tmp/x.parquet").concurrency(0);
611 match ParquetSource::new(cfg).await {
612 Err(FaucetError::Config(msg)) => assert!(msg.contains("concurrency")),
613 other => panic!("expected Config error, got {:?}", other.err()),
614 }
615 }
616
617 #[tokio::test]
618 async fn rejects_s3_with_both_key_and_prefix() {
619 let mut s3 = ParquetS3Config::object("b", "k.parquet");
620 s3.prefix = Some("p/".into());
621 let cfg = ParquetSourceConfig::s3(s3);
622 let source = ParquetSource::new(cfg).await.unwrap();
623 let err = source.resolve_files(&HashMap::new()).await.unwrap_err();
624 assert!(matches!(err, FaucetError::Config(_)));
625 }
626
627 #[tokio::test]
628 async fn rejects_s3_with_neither_key_nor_prefix() {
629 let s3 = ParquetS3Config {
630 bucket: "b".into(),
631 key: None,
632 prefix: None,
633 region: None,
634 endpoint_url: None,
635 };
636 let cfg = ParquetSourceConfig::s3(s3);
637 let source = ParquetSource::new(cfg).await.unwrap();
638 let err = source.resolve_files(&HashMap::new()).await.unwrap_err();
639 assert!(matches!(err, FaucetError::Config(_)));
640 }
641
642 #[test]
643 fn empty_bucket_rejected() {
644 let s3 = ParquetS3Config::object("", "k.parquet");
645 let err = build_s3_store(&s3).unwrap_err();
646 assert!(matches!(err, FaucetError::Config(_)));
647 }
648
649 #[tokio::test]
650 async fn dataset_uri_local_path() {
651 let cfg = ParquetSourceConfig::local("/tmp/data.parquet");
652 let source = ParquetSource::new(cfg).await.unwrap();
653 assert_eq!(source.dataset_uri(), "file:///tmp/data.parquet");
654 }
655
656 #[tokio::test]
657 async fn dataset_uri_glob() {
658 let cfg = ParquetSourceConfig::glob("/tmp/data/*.parquet");
659 let source = ParquetSource::new(cfg).await.unwrap();
660 assert_eq!(source.dataset_uri(), "file:///tmp/data/*.parquet");
661 }
662
663 #[tokio::test]
664 async fn dataset_uri_s3_with_key() {
665 let s3 = ParquetS3Config::object("my-bucket", "path/to/file.parquet");
666 let cfg = ParquetSourceConfig::s3(s3);
667 let source = ParquetSource::new(cfg).await.unwrap();
668 assert_eq!(source.dataset_uri(), "s3://my-bucket/path/to/file.parquet");
669 }
670}