From 573e2b67618564f8b972cf26d4fc12d01116336c Mon Sep 17 00:00:00 2001 From: Burak Sen Date: Thu, 30 Jul 2026 16:23:05 +0300 Subject: [PATCH 1/4] fix(proto): prevent duplicate partition statistics on roundtrip --- .../proto/src/physical_plan/from_proto.rs | 21 ++++++++++++++++--- 1 file changed, 18 insertions(+), 3 deletions(-) diff --git a/datafusion/proto/src/physical_plan/from_proto.rs b/datafusion/proto/src/physical_plan/from_proto.rs index d5a1e0efac6b6..9dd053fd1cec3 100644 --- a/datafusion/proto/src/physical_plan/from_proto.rs +++ b/datafusion/proto/src/physical_plan/from_proto.rs @@ -660,7 +660,7 @@ impl TryFromProto<&protobuf::PartitionedFile> for PartitionedFile { pf = pf.with_range(file_range.start, file_range.end); } if let Some(proto_stats) = val.statistics.as_ref() { - pf = pf.with_statistics(Arc::new(proto_stats.try_into()?)); + pf.statistics = Some(Arc::new(proto_stats.try_into()?)); } Ok(pf) } @@ -807,8 +807,8 @@ impl datafusion_physical_expr_common::physical_expr::proto_decode::PhysicalExprD #[cfg(test)] mod tests { - use super::*; + use arrow::datatypes::{DataType, Field, Schema}; #[test] fn partitioned_file_path_roundtrip_percent_encoded() { @@ -833,7 +833,6 @@ mod tests { #[test] fn partitioned_file_arrow_schema_roundtrip() { - use arrow::datatypes::{DataType, Field, Schema}; use std::collections::HashMap; let arrow_schema = Arc::new(Schema::new_with_metadata( @@ -858,6 +857,22 @@ mod tests { ); } + #[test] + fn partitioned_file_statistics_roundtrip_with_partition_values() { + use datafusion_common::Statistics; + + let file_schema = Schema::new(vec![Field::new("a", DataType::Int32, true)]); + let pf = PartitionedFile::new("foo/bar.parquet", 1234) + .with_partition_values(vec![ScalarValue::from("2024-01-01")]) + .with_statistics(Arc::new(Statistics::new_unknown(&file_schema))); + assert_eq!(pf.statistics.as_ref().unwrap().column_statistics.len(), 2); + + let proto = protobuf::PartitionedFile::try_from_proto(&pf).unwrap(); + let decoded = PartitionedFile::try_from_proto(&proto).unwrap(); + + assert_eq!(decoded.statistics, pf.statistics); + } + #[test] fn partitioned_file_from_proto_invalid_path() { let proto = protobuf::PartitionedFile { From 2a919a87effb6a30e1d66b6b1c96a0227fc499ac Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Burak=20=C5=9Een?= Date: Thu, 30 Jul 2026 17:36:04 +0300 Subject: [PATCH 2/4] Update datafusion/proto/src/physical_plan/from_proto.rs Co-authored-by: Adrian Garcia Badaracco <1755071+adriangb@users.noreply.github.com> --- datafusion/proto/src/physical_plan/from_proto.rs | 12 +++++++----- 1 file changed, 7 insertions(+), 5 deletions(-) diff --git a/datafusion/proto/src/physical_plan/from_proto.rs b/datafusion/proto/src/physical_plan/from_proto.rs index 9dd053fd1cec3..04817239b7670 100644 --- a/datafusion/proto/src/physical_plan/from_proto.rs +++ b/datafusion/proto/src/physical_plan/from_proto.rs @@ -861,11 +861,13 @@ mod tests { fn partitioned_file_statistics_roundtrip_with_partition_values() { use datafusion_common::Statistics; - let file_schema = Schema::new(vec![Field::new("a", DataType::Int32, true)]); - let pf = PartitionedFile::new("foo/bar.parquet", 1234) - .with_partition_values(vec![ScalarValue::from("2024-01-01")]) - .with_statistics(Arc::new(Statistics::new_unknown(&file_schema))); - assert_eq!(pf.statistics.as_ref().unwrap().column_statistics.len(), 2); + // `statistics` covers the full table schema: file columns followed by one + // entry per partition column. + let expected_len = file_schema.fields().len() + pf.partition_values.len(); + assert_eq!( + pf.statistics.as_ref().unwrap().column_statistics.len(), + expected_len + ); let proto = protobuf::PartitionedFile::try_from_proto(&pf).unwrap(); let decoded = PartitionedFile::try_from_proto(&proto).unwrap(); From 6063c66d7dcd53e839046c07bbf470c000e749f7 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Burak=20=C5=9Een?= Date: Thu, 30 Jul 2026 17:36:12 +0300 Subject: [PATCH 3/4] Update datafusion/proto/src/physical_plan/from_proto.rs Co-authored-by: Adrian Garcia Badaracco <1755071+adriangb@users.noreply.github.com> --- datafusion/proto/src/physical_plan/from_proto.rs | 3 +++ 1 file changed, 3 insertions(+) diff --git a/datafusion/proto/src/physical_plan/from_proto.rs b/datafusion/proto/src/physical_plan/from_proto.rs index 04817239b7670..7085b0677d780 100644 --- a/datafusion/proto/src/physical_plan/from_proto.rs +++ b/datafusion/proto/src/physical_plan/from_proto.rs @@ -660,6 +660,9 @@ impl TryFromProto<&protobuf::PartitionedFile> for PartitionedFile { pf = pf.with_range(file_range.start, file_range.end); } if let Some(proto_stats) = val.statistics.as_ref() { + // The wire format carries statistics for the full table schema (file + partition + // columns), so assign directly — `with_statistics` would append the partition + // column stats a second time. pf.statistics = Some(Arc::new(proto_stats.try_into()?)); } Ok(pf) From cd9ce3488e49735f21a062c5bf26534a476a8b6f Mon Sep 17 00:00:00 2001 From: buraksenn Date: Thu, 30 Jul 2026 18:14:32 +0300 Subject: [PATCH 4/4] fix github suggestion deleting test setup --- datafusion/proto/src/physical_plan/from_proto.rs | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/datafusion/proto/src/physical_plan/from_proto.rs b/datafusion/proto/src/physical_plan/from_proto.rs index 7085b0677d780..7aa6376313c96 100644 --- a/datafusion/proto/src/physical_plan/from_proto.rs +++ b/datafusion/proto/src/physical_plan/from_proto.rs @@ -863,6 +863,10 @@ mod tests { #[test] fn partitioned_file_statistics_roundtrip_with_partition_values() { use datafusion_common::Statistics; + let file_schema = Schema::new(vec![Field::new("a", DataType::Int32, true)]); + let pf = PartitionedFile::new("foo/bar.parquet", 1234) + .with_partition_values(vec![ScalarValue::from("2024-01-01")]) + .with_statistics(Arc::new(Statistics::new_unknown(&file_schema))); // `statistics` covers the full table schema: file columns followed by one // entry per partition column.