Skip to content

Commit ed75a23

Browse files
committed
Datasource: Add file_metadata() to DataSink for per-file write metadata
Add FileWriteMetadata struct to datafusion-datasource along with a new default method on the DataSink trait: - file_metadata(): Post-hoc accessor returning per-file path, row count, byte size, and optional format-specific metadata bytes after write_all completes. The default implementation returns an empty Vec, so existing DataSink implementations are unaffected. ParquetSink overrides file_metadata() to expose path, row_count, and byte_size for each written file using the existing self.written HashMap that already collects ParquetMetaData during writes. The typed ParquetMetaData remains accessible via ParquetSink::written() for Rust consumers who need column-level statistics directly. DataSinkExec gains a file_metadata() convenience accessor that delegates to the underlying sink. This enables external table formats (Iceberg, Delta Lake, Hudi) to obtain per-file metadata after a DataFusion write without re-reading file footers or bypassing the write pipeline entirely. Closes #23472
1 parent 3e058f0 commit ed75a23

2 files changed

Lines changed: 549 additions & 1 deletion

File tree

datafusion/datasource-parquet/src/sink.rs

Lines changed: 301 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -32,7 +32,7 @@ use datafusion_common_runtime::{JoinSet, SpawnedTask};
3232
use datafusion_datasource::display::FileGroupDisplay;
3333
use datafusion_datasource::file_compression_type::FileCompressionType;
3434
use datafusion_datasource::file_sink_config::{FileSink, FileSinkConfig};
35-
use datafusion_datasource::sink::DataSink;
35+
use datafusion_datasource::sink::{DataSink, FileWriteMetadata};
3636
use datafusion_datasource::write::demux::DemuxedStreamReceiver;
3737
use datafusion_datasource::write::{
3838
ObjectWriterBuilder, SharedBuffer, get_writer_schema,
@@ -411,6 +411,31 @@ impl DataSink for ParquetSink {
411411
) -> Result<u64> {
412412
FileSink::write_all(self, data, context).await
413413
}
414+
415+
fn file_metadata(&self) -> Vec<FileWriteMetadata> {
416+
let written = self.written.lock();
417+
written
418+
.iter()
419+
.map(|(path, parquet_meta)| {
420+
let row_count = parquet_meta.file_metadata().num_rows() as u64;
421+
let byte_size: u64 = parquet_meta
422+
.row_groups()
423+
.iter()
424+
.map(|rg| rg.compressed_size() as u64)
425+
.sum();
426+
427+
FileWriteMetadata {
428+
path: path.to_string(),
429+
row_count,
430+
byte_size,
431+
// Full Parquet metadata is accessible via ParquetSink::written()
432+
// for Rust consumers. Thrift serialization for FFI consumers can
433+
// be added when a concrete use case (e.g. Python via PyO3) needs it.
434+
format_metadata: None,
435+
}
436+
})
437+
.collect()
438+
}
414439
}
415440

416441
/// Consumes a stream of [ArrowLeafColumn] via a channel and serializes them using an [ArrowColumnWriter]
@@ -752,3 +777,278 @@ async fn output_single_parquet_file_parallelized(
752777
.map_err(|e| DataFusionError::ExecutionJoin(Box::new(e)))??;
753778
Ok(parquet_meta_data)
754779
}
780+
781+
#[cfg(test)]
782+
mod tests {
783+
use super::*;
784+
use arrow::array::{ArrayRef, StringArray};
785+
use arrow::datatypes::{DataType, Field, Schema};
786+
use datafusion_common::config::TableParquetOptions;
787+
use datafusion_datasource::PartitionedFile;
788+
use datafusion_datasource::file_groups::FileGroup;
789+
use datafusion_datasource::file_sink_config::{FileOutputMode, FileSinkConfig};
790+
use datafusion_datasource::sink::{DataSink, DataSinkExec};
791+
use datafusion_datasource::url::ListingTableUrl;
792+
use datafusion_execution::config::SessionConfig;
793+
use datafusion_execution::object_store::ObjectStoreUrl;
794+
use datafusion_execution::runtime_env::RuntimeEnv;
795+
use datafusion_expr::dml::InsertOp;
796+
use datafusion_physical_plan::stream::RecordBatchStreamAdapter;
797+
use object_store::local::LocalFileSystem;
798+
799+
fn build_test_ctx(store_url: &ObjectStoreUrl) -> Arc<TaskContext> {
800+
let tmp_dir = tempfile::TempDir::new().unwrap();
801+
let local = Arc::new(
802+
LocalFileSystem::new_with_prefix(&tmp_dir)
803+
.expect("should create object store"),
804+
);
805+
806+
let session = SessionConfig::default();
807+
let runtime = RuntimeEnv::default();
808+
runtime
809+
.object_store_registry
810+
.register_store(store_url.as_ref(), local);
811+
812+
Arc::new(
813+
TaskContext::default()
814+
.with_session_config(session)
815+
.with_runtime(Arc::new(runtime)),
816+
)
817+
}
818+
819+
fn create_test_sink() -> (Arc<ParquetSink>, SchemaRef, ObjectStoreUrl) {
820+
let field_a = Field::new("a", DataType::Utf8, false);
821+
let field_b = Field::new("b", DataType::Utf8, false);
822+
let schema = Arc::new(Schema::new(vec![field_a, field_b]));
823+
let object_store_url = ObjectStoreUrl::local_filesystem();
824+
825+
let file_sink_config = FileSinkConfig {
826+
original_url: String::default(),
827+
object_store_url: object_store_url.clone(),
828+
file_group: FileGroup::new(vec![PartitionedFile::new("/tmp".to_string(), 1)]),
829+
table_paths: vec![ListingTableUrl::parse("file:///tmp/test/").unwrap()],
830+
output_schema: Arc::clone(&schema),
831+
table_partition_cols: vec![],
832+
insert_op: InsertOp::Overwrite,
833+
keep_partition_by_columns: false,
834+
file_extension: "parquet".into(),
835+
file_output_mode: FileOutputMode::Automatic,
836+
};
837+
838+
let parquet_sink = Arc::new(ParquetSink::new(
839+
file_sink_config,
840+
TableParquetOptions::default(),
841+
));
842+
843+
(parquet_sink, schema, object_store_url)
844+
}
845+
846+
fn make_test_batch(schema: &SchemaRef) -> RecordBatch {
847+
let col_a: ArrayRef = Arc::new(StringArray::from(vec!["foo", "bar", "baz"]));
848+
let col_b: ArrayRef = Arc::new(StringArray::from(vec!["one", "two", "three"]));
849+
RecordBatch::try_new(Arc::clone(schema), vec![col_a, col_b]).unwrap()
850+
}
851+
852+
#[test]
853+
fn file_metadata_empty_before_write() {
854+
let (sink, _schema, _url) = create_test_sink();
855+
assert!(
856+
sink.file_metadata().is_empty(),
857+
"file_metadata should be empty before any write"
858+
);
859+
}
860+
861+
#[tokio::test]
862+
async fn file_metadata_populated_after_write() {
863+
let (sink, schema, object_store_url) = create_test_sink();
864+
let ctx = build_test_ctx(&object_store_url);
865+
let batch = make_test_batch(&schema);
866+
867+
let data: SendableRecordBatchStream = Box::pin(RecordBatchStreamAdapter::new(
868+
Arc::clone(&schema),
869+
futures::stream::iter(vec![Ok(batch)]),
870+
));
871+
872+
let count = DataSink::write_all(sink.as_ref(), data, &ctx)
873+
.await
874+
.unwrap();
875+
assert_eq!(count, 3, "should have written 3 rows");
876+
877+
let metadata = sink.file_metadata();
878+
assert_eq!(metadata.len(), 1, "should have one file metadata entry");
879+
assert_eq!(metadata[0].row_count, 3);
880+
assert!(metadata[0].byte_size > 0, "byte_size should be non-zero");
881+
assert!(!metadata[0].path.is_empty(), "path should be non-empty");
882+
assert!(metadata[0].format_metadata.is_none());
883+
}
884+
885+
#[tokio::test]
886+
async fn file_metadata_count_equals_write_all_count() {
887+
let (sink, schema, object_store_url) = create_test_sink();
888+
let ctx = build_test_ctx(&object_store_url);
889+
let batch = make_test_batch(&schema);
890+
891+
let data: SendableRecordBatchStream = Box::pin(RecordBatchStreamAdapter::new(
892+
Arc::clone(&schema),
893+
futures::stream::iter(vec![Ok(batch)]),
894+
));
895+
896+
let count = DataSink::write_all(sink.as_ref(), data, &ctx)
897+
.await
898+
.unwrap();
899+
900+
let sum_rows: u64 = sink.file_metadata().iter().map(|f| f.row_count).sum();
901+
assert_eq!(
902+
count, sum_rows,
903+
"write_all count must equal sum of per-file row_counts"
904+
);
905+
}
906+
907+
#[tokio::test]
908+
async fn file_metadata_consistent_with_written() {
909+
let (sink, schema, object_store_url) = create_test_sink();
910+
let ctx = build_test_ctx(&object_store_url);
911+
let batch = make_test_batch(&schema);
912+
913+
let data: SendableRecordBatchStream = Box::pin(RecordBatchStreamAdapter::new(
914+
Arc::clone(&schema),
915+
futures::stream::iter(vec![Ok(batch)]),
916+
));
917+
918+
DataSink::write_all(sink.as_ref(), data, &ctx)
919+
.await
920+
.unwrap();
921+
922+
let file_meta = sink.file_metadata();
923+
let written = sink.written();
924+
925+
assert_eq!(
926+
file_meta.len(),
927+
written.len(),
928+
"file_metadata and written() should have same number of entries"
929+
);
930+
931+
for fm in &file_meta {
932+
let path = Path::from(fm.path.as_str());
933+
let parquet_meta = written
934+
.get(&path)
935+
.expect("file_metadata path should exist in written()");
936+
assert_eq!(fm.row_count, parquet_meta.file_metadata().num_rows() as u64);
937+
938+
let expected_bytes: u64 = parquet_meta
939+
.row_groups()
940+
.iter()
941+
.map(|rg| rg.compressed_size() as u64)
942+
.sum();
943+
assert_eq!(fm.byte_size, expected_bytes);
944+
}
945+
}
946+
947+
#[tokio::test]
948+
async fn file_metadata_is_idempotent() {
949+
let (sink, schema, object_store_url) = create_test_sink();
950+
let ctx = build_test_ctx(&object_store_url);
951+
let batch = make_test_batch(&schema);
952+
953+
let data: SendableRecordBatchStream = Box::pin(RecordBatchStreamAdapter::new(
954+
Arc::clone(&schema),
955+
futures::stream::iter(vec![Ok(batch)]),
956+
));
957+
958+
DataSink::write_all(sink.as_ref(), data, &ctx)
959+
.await
960+
.unwrap();
961+
962+
let first = sink.file_metadata();
963+
let second = sink.file_metadata();
964+
assert_eq!(first, second);
965+
}
966+
967+
#[tokio::test]
968+
async fn partitioned_write_produces_multiple_metadata_entries() {
969+
let field_a = Field::new("a", DataType::Utf8, false);
970+
let field_b = Field::new("b", DataType::Utf8, false);
971+
let schema = Arc::new(Schema::new(vec![field_a, field_b]));
972+
let object_store_url = ObjectStoreUrl::local_filesystem();
973+
974+
let file_sink_config = FileSinkConfig {
975+
original_url: String::default(),
976+
object_store_url: object_store_url.clone(),
977+
file_group: FileGroup::new(vec![PartitionedFile::new("/tmp".to_string(), 1)]),
978+
table_paths: vec![
979+
ListingTableUrl::parse("file:///tmp/partitioned/").unwrap(),
980+
],
981+
output_schema: Arc::clone(&schema),
982+
table_partition_cols: vec![("a".to_string(), DataType::Utf8)],
983+
insert_op: InsertOp::Overwrite,
984+
keep_partition_by_columns: false,
985+
file_extension: "parquet".into(),
986+
file_output_mode: FileOutputMode::Automatic,
987+
};
988+
989+
let sink = Arc::new(ParquetSink::new(
990+
file_sink_config,
991+
TableParquetOptions::default(),
992+
));
993+
994+
let col_a: ArrayRef = Arc::new(StringArray::from(vec!["x", "y", "x"]));
995+
let col_b: ArrayRef = Arc::new(StringArray::from(vec!["one", "two", "three"]));
996+
let batch =
997+
RecordBatch::try_new(Arc::clone(&schema), vec![col_a, col_b]).unwrap();
998+
999+
let ctx = build_test_ctx(&object_store_url);
1000+
let data: SendableRecordBatchStream = Box::pin(RecordBatchStreamAdapter::new(
1001+
Arc::clone(&schema),
1002+
futures::stream::iter(vec![Ok(batch)]),
1003+
));
1004+
1005+
let count = DataSink::write_all(sink.as_ref(), data, &ctx)
1006+
.await
1007+
.unwrap();
1008+
1009+
let metadata = sink.file_metadata();
1010+
assert_eq!(
1011+
metadata.len(),
1012+
2,
1013+
"partitioned write with 2 partition values should produce 2 files"
1014+
);
1015+
1016+
assert_eq!(count, 3);
1017+
let sum_rows: u64 = metadata.iter().map(|f| f.row_count).sum();
1018+
assert_eq!(sum_rows, 3);
1019+
}
1020+
1021+
/// Verifies that `DataSinkExec::file_metadata()` returns the correct
1022+
/// metadata when the underlying `ParquetSink` has written files.
1023+
/// This tests the `Arc<dyn DataSink>` dispatch path that real consumers use.
1024+
#[tokio::test]
1025+
async fn data_sink_exec_e2e_file_metadata_after_execute() {
1026+
let (sink, schema, object_store_url) = create_test_sink();
1027+
let ctx = build_test_ctx(&object_store_url);
1028+
let batch = make_test_batch(&schema);
1029+
1030+
let data: SendableRecordBatchStream = Box::pin(RecordBatchStreamAdapter::new(
1031+
Arc::clone(&schema),
1032+
futures::stream::iter(vec![Ok(batch)]),
1033+
));
1034+
1035+
// Write through the DataSink trait (same call DataSinkExec::execute makes)
1036+
let count = DataSink::write_all(sink.as_ref(), data, &ctx)
1037+
.await
1038+
.unwrap();
1039+
assert_eq!(count, 3);
1040+
1041+
// Wrap in DataSinkExec and verify file_metadata is accessible
1042+
// through the exec's convenience accessor (the Arc<dyn DataSink> path)
1043+
let input: Arc<dyn datafusion_physical_plan::ExecutionPlan> = Arc::new(
1044+
datafusion_physical_plan::empty::EmptyExec::new(Arc::clone(&schema)),
1045+
);
1046+
let exec = DataSinkExec::new(input, sink, None);
1047+
1048+
let metadata = exec.file_metadata();
1049+
assert_eq!(metadata.len(), 1, "should have one file after write");
1050+
assert_eq!(metadata[0].row_count, 3);
1051+
assert!(metadata[0].byte_size > 0);
1052+
assert!(!metadata[0].path.is_empty());
1053+
}
1054+
}

0 commit comments

Comments
 (0)