From fda51119aae39409124309a1501244e78cf5d975 Mon Sep 17 00:00:00 2001 From: batu <109382994+sovsparrow@users.noreply.github.com> Date: Sat, 15 Aug 2026 00:46:36 +0300 Subject: [PATCH] fix: reject duplicate field names during CSV and Parquet schema inference --- .../core/src/datasource/file_format/csv.rs | 27 +++++++++++++++ .../src/datasource/file_format/parquet.rs | 34 +++++++++++++++++++ datafusion/datasource-csv/src/file_format.rs | 7 ++++ .../datasource-parquet/src/file_format.rs | 13 ++++++- datafusion/datasource/src/file_format.rs | 24 +++++++++++-- 5 files changed, 101 insertions(+), 4 deletions(-) diff --git a/datafusion/core/src/datasource/file_format/csv.rs b/datafusion/core/src/datasource/file_format/csv.rs index 90d7eb3b41388..447d5244f9e94 100644 --- a/datafusion/core/src/datasource/file_format/csv.rs +++ b/datafusion/core/src/datasource/file_format/csv.rs @@ -1628,4 +1628,31 @@ mod tests { Ok(()) } + + #[tokio::test] + async fn infer_schema_rejects_duplicate_header_names() -> Result<()> { + let directory = tempfile::tempdir()?; + let path = directory.path().join("duplicate_header.csv"); + std::fs::write(&path, "id,value,value\n1,10,100\n")?; + + let store = Arc::new(LocalFileSystem::new()) as _; + let meta = crate::test::object_store::local_unpartitioned_file(&path); + + let ctx = SessionContext::new().state(); + let error = CsvFormat::default() + .with_has_header(true) + .infer_schema(&ctx, &store, std::slice::from_ref(&meta)) + .await + .expect_err("duplicate header names must not infer a schema") + .to_string(); + + assert!( + error.contains("duplicate unqualified field name") + && error.contains("value") + && error.contains("duplicate_header.csv"), + "unexpected error: {error}" + ); + + Ok(()) + } } diff --git a/datafusion/core/src/datasource/file_format/parquet.rs b/datafusion/core/src/datasource/file_format/parquet.rs index 0f5db4a057d76..bbd9d0937ad83 100644 --- a/datafusion/core/src/datasource/file_format/parquet.rs +++ b/datafusion/core/src/datasource/file_format/parquet.rs @@ -1822,4 +1822,38 @@ mod tests { Ok(()) } + + #[tokio::test] + async fn infer_schema_rejects_duplicate_field_names() -> Result<()> { + let schema = Arc::new(Schema::new(vec![ + Field::new("id", DataType::Int64, false), + Field::new("value", DataType::Int64, false), + Field::new("value", DataType::Int64, false), + ])); + let batch = RecordBatch::try_new( + schema, + vec![ + Arc::new(Int64Array::from(vec![1, 2, 3])) as ArrayRef, + Arc::new(Int64Array::from(vec![10, 20, 30])) as ArrayRef, + Arc::new(Int64Array::from(vec![100, 200, 300])) as ArrayRef, + ], + )?; + + let store = Arc::new(LocalFileSystem::new()) as _; + let (meta, _files) = store_parquet(vec![batch], false).await?; + + let ctx = SessionContext::new().state(); + let error = ParquetFormat::default() + .infer_schema(&ctx, &store, &meta) + .await + .expect_err("duplicate field names must not infer a schema") + .to_string(); + + assert!( + error.contains("duplicate unqualified field name") && error.contains("value"), + "unexpected error: {error}" + ); + + Ok(()) + } } diff --git a/datafusion/datasource-csv/src/file_format.rs b/datafusion/datasource-csv/src/file_format.rs index c0d22b80f08d0..aff13f7c40a30 100644 --- a/datafusion/datasource-csv/src/file_format.rs +++ b/datafusion/datasource-csv/src/file_format.rs @@ -41,6 +41,7 @@ use datafusion_datasource::file::FileSource; use datafusion_datasource::file_compression_type::FileCompressionType; use datafusion_datasource::file_format::{ DEFAULT_SCHEMA_INFER_MAX_RECORD, FileFormat, FileFormatFactory, + ensure_unique_field_names, }; use datafusion_datasource::file_scan_config::{FileScanConfig, FileScanConfigBuilder}; use datafusion_datasource::file_sink_config::{FileSink, FileSinkConfig}; @@ -396,6 +397,12 @@ impl FileFormat for CsvFormat { Box::new(err), ) })?; + ensure_unique_field_names(&schema).map_err(|err| { + DataFusionError::Context( + format!("Error when processing CSV file {}", object.location), + Box::new(err), + ) + })?; records_to_read -= records_read; schemas.push(schema); if records_to_read == 0 { diff --git a/datafusion/datasource-parquet/src/file_format.rs b/datafusion/datasource-parquet/src/file_format.rs index 6358201c06fa5..2d39a16c776fd 100644 --- a/datafusion/datasource-parquet/src/file_format.rs +++ b/datafusion/datasource-parquet/src/file_format.rs @@ -37,7 +37,9 @@ use datafusion_datasource::TableSchema; use datafusion_datasource::file_compression_type::FileCompressionType; use datafusion_datasource::file_sink_config::FileSinkConfig; -use datafusion_datasource::file_format::{FileFormat, FileFormatFactory}; +use datafusion_datasource::file_format::{ + FileFormat, FileFormatFactory, ensure_unique_field_names, +}; use datafusion_common::Statistics; use datafusion_common::config::{ConfigField, ConfigFileType, TableParquetOptions}; @@ -388,6 +390,15 @@ impl FileFormat for ParquetFormat { schemas .sort_unstable_by(|(location1, _), (location2, _)| location1.cmp(location2)); + for (location, schema) in &schemas { + ensure_unique_field_names(schema).map_err(|err| { + DataFusionError::Context( + format!("Error when processing Parquet file {location}"), + Box::new(err), + ) + })?; + } + let schemas = schemas.into_iter().map(|(_, schema)| schema); let schema = if self.skip_metadata() { diff --git a/datafusion/datasource/src/file_format.rs b/datafusion/datasource/src/file_format.rs index dd30881610f36..0d2541ad39a3c 100644 --- a/datafusion/datasource/src/file_format.rs +++ b/datafusion/datasource/src/file_format.rs @@ -19,7 +19,7 @@ //! See write.rs for write related helper methods use std::any::Any; -use std::collections::HashMap; +use std::collections::{HashMap, HashSet}; use std::fmt; use std::sync::Arc; @@ -28,9 +28,11 @@ use crate::file_compression_type::FileCompressionType; use crate::file_scan_config::FileScanConfig; use crate::file_sink_config::FileSinkConfig; -use arrow::datatypes::SchemaRef; +use arrow::datatypes::{Schema, SchemaRef}; use datafusion_common::file_options::file_type::FileType; -use datafusion_common::{GetExt, Result, Statistics, internal_err, not_impl_err}; +use datafusion_common::{ + GetExt, Result, SchemaError, Statistics, internal_err, not_impl_err, schema_err, +}; use datafusion_physical_expr::LexRequirement; use datafusion_physical_expr_common::sort_expr::LexOrdering; use datafusion_physical_plan::ExecutionPlan; @@ -42,6 +44,22 @@ use object_store::{ObjectMeta, ObjectStore}; /// Default max records to scan to infer the schema pub const DEFAULT_SCHEMA_INFER_MAX_RECORD: usize = 1000; +/// Rejects an inferred schema that names the same field more than once. +/// +/// [`Schema::try_merge`] coalesces fields by name, so callers validate each +/// inferred file schema before merging. +pub fn ensure_unique_field_names(schema: &Schema) -> Result<()> { + let mut seen = HashSet::with_capacity(schema.fields().len()); + for field in schema.fields() { + if !seen.insert(field.name()) { + return schema_err!(SchemaError::DuplicateUnqualifiedField { + name: field.name().clone(), + }); + } + } + Ok(()) +} + /// Metadata fetched from a file, including statistics and ordering. /// /// This struct is returned by [`FileFormat::infer_stats_and_ordering`] to