Skip to content

Commit c4fdbdd

Browse files
ShayanGhoclaude
andcommitted
fix: drop partition columns already present in the file schema of a ListingTable
Files written with `keep_partition_by_columns = true` physically contain the partition column. `ListingTable::try_new` appended the configured partition columns to the inferred file schema unconditionally, so the table schema listed the column twice and any query failed with "Schema contains duplicate qualified field name". The partition column now appears once, with the declared partition type, and its values come from the path, matching what CREATE EXTERNAL TABLE with an explicit column list already did. Closes #17420 Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
1 parent 9082d6b commit c4fdbdd

4 files changed

Lines changed: 118 additions & 6 deletions

File tree

‎datafusion/catalog-listing/src/table.rs‎

Lines changed: 26 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -225,6 +225,32 @@ impl ListingTable {
225225
.options
226226
.ok_or_else(|| internal_datafusion_err!("No ListingOptions provided"))?;
227227

228+
// Files may physically contain the partition columns, for example when
229+
// they were written with `keep_partition_by_columns = true`. Partition
230+
// column values are always taken from the path, so drop such columns
231+
// from the file schema to avoid duplicated fields in the table schema.
232+
let file_schema = if options
233+
.table_partition_cols
234+
.iter()
235+
.any(|(name, _)| file_schema.field_with_name(name).is_ok())
236+
{
237+
let indices: Vec<usize> = file_schema
238+
.fields()
239+
.iter()
240+
.enumerate()
241+
.filter(|(_, field)| {
242+
!options
243+
.table_partition_cols
244+
.iter()
245+
.any(|(name, _)| name == field.name())
246+
})
247+
.map(|(idx, _)| idx)
248+
.collect();
249+
Arc::new(file_schema.project(&indices)?)
250+
} else {
251+
file_schema
252+
};
253+
228254
// Add the partition columns to the file schema
229255
let mut builder = SchemaBuilder::from(file_schema.as_ref().to_owned());
230256
for (part_col_name, part_col_type) in &options.table_partition_cols {

‎datafusion/core/src/datasource/listing/table.rs‎

Lines changed: 37 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1625,6 +1625,43 @@ mod tests {
16251625
Ok(())
16261626
}
16271627

1628+
/// Files written with `keep_partition_by_columns = true` physically contain
1629+
/// the partition column, so the (inferred) file schema and the declared
1630+
/// partition columns overlap. The table schema must list the column once.
1631+
/// See <https://github.com/apache/datafusion/issues/17420>
1632+
#[test]
1633+
fn test_partition_column_present_in_file_schema_is_not_duplicated() -> Result<()> {
1634+
let partition_type =
1635+
DataType::Dictionary(Box::new(DataType::UInt16), Box::new(DataType::Utf8));
1636+
let opt = ListingOptions::new(Arc::new(JsonFormat::default()))
1637+
.with_file_extension_opt(Some(""))
1638+
.with_table_partition_cols(vec![("pid".to_string(), partition_type.clone())]);
1639+
1640+
let table_path = ListingTableUrl::parse("test:///bucket/test/")?;
1641+
let file_schema = Schema::new(vec![
1642+
Field::new("a", DataType::Boolean, false),
1643+
Field::new("pid", DataType::Utf8, false),
1644+
]);
1645+
let config = ListingTableConfig::new(table_path)
1646+
.with_listing_options(opt)
1647+
.with_schema(Arc::new(file_schema));
1648+
1649+
let table = ListingTable::try_new(config)?;
1650+
let schema = table.schema();
1651+
1652+
let names = schema
1653+
.fields()
1654+
.iter()
1655+
.map(|f| f.name().as_str())
1656+
.collect::<Vec<_>>();
1657+
assert_eq!(names, vec!["a", "pid"]);
1658+
// The partition column keeps the declared partition type, as it does
1659+
// when the file schema is given explicitly in CREATE EXTERNAL TABLE.
1660+
assert_eq!(schema.field_with_name("pid")?.data_type(), &partition_type);
1661+
1662+
Ok(())
1663+
}
1664+
16281665
#[cfg(feature = "parquet")]
16291666
#[tokio::test]
16301667
async fn test_table_stats_behaviors() -> Result<()> {

‎datafusion/core/tests/sql/path_partition.rs‎

Lines changed: 28 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -514,7 +514,11 @@ async fn parquet_statistics() -> Result<()> {
514514
async fn parquet_overlapping_columns() -> Result<()> {
515515
let ctx = SessionContext::new();
516516

517-
// `id` is both a column of the file and a partitioning col
517+
// `id` is both a column of the file (Int32, values 0..8) and a partition
518+
// column (Int64). Files written with `keep_partition_by_columns = true`
519+
// look exactly like this. The partition column wins: it appears once in
520+
// the schema, with the partition type, and its values come from the path.
521+
// See https://github.com/apache/datafusion/issues/17420
518522
register_partitioned_alltypes_parquet(
519523
&ctx,
520524
&[
@@ -528,12 +532,30 @@ async fn parquet_overlapping_columns() -> Result<()> {
528532
)
529533
.await;
530534

531-
let result = ctx.sql("SELECT id FROM t WHERE id=1 ORDER BY id").await;
535+
let schema = ctx.table_provider("t").await?.schema();
536+
let id_fields = schema
537+
.fields()
538+
.iter()
539+
.filter(|f| f.name() == "id")
540+
.collect::<Vec<_>>();
541+
assert_eq!(id_fields.len(), 1, "partition column must appear only once");
542+
assert_eq!(id_fields[0].data_type(), &DataType::Int64);
543+
544+
// All 8 rows of the `id=2` file report the path value, not the file value
545+
let result = ctx
546+
.sql("SELECT id, count(*) AS n FROM t WHERE id = 2 GROUP BY id")
547+
.await?
548+
.collect()
549+
.await?;
550+
551+
assert_snapshot!(batches_to_sort_string(&result), @r"
552+
+----+---+
553+
| id | n |
554+
+----+---+
555+
| 2 | 8 |
556+
+----+---+
557+
");
532558

533-
assert!(
534-
result.is_err(),
535-
"Duplicate qualified name should raise error"
536-
);
537559
Ok(())
538560
}
539561

‎datafusion/sqllogictest/test_files/copy.slt‎

Lines changed: 27 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -228,6 +228,33 @@ select column1, column2 from validate_partitioned_parquet4 order by column1,colu
228228
----
229229
1 a
230230

231+
# The files written above physically contain the partition column (because
232+
# keep_partition_by_columns was enabled), so a hive-partitioned table over the
233+
# whole directory must not end up with that column twice.
234+
# See https://github.com/apache/datafusion/issues/17420
235+
statement ok
236+
CREATE EXTERNAL TABLE validate_partitioned_parquet4_hive STORED AS PARQUET
237+
LOCATION 'test_files/scratch/copy/partitioned_table4/' PARTITIONED BY (column1);
238+
239+
query TT
240+
select column1, column2 from validate_partitioned_parquet4_hive order by column1, column2;
241+
----
242+
1 a
243+
2 b
244+
3 c
245+
246+
# Same directory, but letting the factory infer the hive partitions from the path
247+
statement ok
248+
CREATE EXTERNAL TABLE validate_partitioned_parquet4_inferred STORED AS PARQUET
249+
LOCATION 'test_files/scratch/copy/partitioned_table4/';
250+
251+
query TT
252+
select column1, column2 from validate_partitioned_parquet4_inferred order by column1, column2;
253+
----
254+
1 a
255+
2 b
256+
3 c
257+
231258
# Copy more files to directory via query
232259
query I
233260
COPY (select * from source_table UNION ALL select * from source_table) to 'test_files/scratch/copy/table/' STORED AS PARQUET;

0 commit comments

Comments
 (0)