Skip to content

Commit 745aed7

Browse files
committed
Improve error message for metadata conflict in schema
1 parent 40b5410 commit 745aed7

3 files changed

Lines changed: 62 additions & 16 deletions

File tree

datafusion/core/src/physical_planner.rs

Lines changed: 32 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -3258,15 +3258,14 @@ impl<'a> OptimizationInvariantChecker<'a> {
32583258
previous_schema: &Arc<Schema>,
32593259
) -> Result<()> {
32603260
// if the rule is not permitted to change the schema, confirm that it did not change.
3261-
if self.rule.schema_check()
3262-
&& !is_allowed_schema_change(previous_schema.as_ref(), plan.schema().as_ref())
3263-
{
3264-
internal_err!(
3265-
"PhysicalOptimizer rule '{}' failed. Schema mismatch. Expected original schema: {}, got new schema: {}",
3266-
self.rule.name(),
3267-
previous_schema,
3268-
plan.schema()
3269-
)?
3261+
if self.rule.schema_check() {
3262+
is_allowed_schema_change(previous_schema.as_ref(), plan.schema().as_ref())
3263+
.map_err(|e| {
3264+
e.context(format!(
3265+
"PhysicalOptimizer rule '{}' failed. Schema mismatch.",
3266+
self.rule.name(),
3267+
))
3268+
})?
32703269
}
32713270

32723271
// check invariants per each ExecutionPlan node
@@ -3285,28 +3284,45 @@ impl<'a> OptimizationInvariantChecker<'a> {
32853284
/// This change is allowed because for any field the non-nullable domain `F` is a strict subset
32863285
/// of the nullable domain `F ∪ { NULL }`. A physical schema that guarantees a stricter subset
32873286
/// of values will not violate any assumptions made based on the less strict schema.
3288-
fn is_allowed_schema_change(old: &Schema, new: &Schema) -> bool {
3287+
fn is_allowed_schema_change(old: &Schema, new: &Schema) -> Result<()> {
32893288
if new.metadata != old.metadata {
3290-
return false;
3289+
return internal_err!(
3290+
"Schema metadata mismatch: Expected original metadata: {:?}, got metadata: {:?}",
3291+
old.metadata,
3292+
new.metadata
3293+
);
32913294
}
32923295

32933296
if new.fields.len() != old.fields.len() {
3294-
return false;
3297+
return internal_err!(
3298+
"Schema field mismatch: Expected original field count: {}, got field count: {}",
3299+
old.fields.len(),
3300+
new.fields.len()
3301+
);
32953302
}
32963303

32973304
let new_fields = new.fields.iter().map(|f| f.as_ref());
32983305
let old_fields = old.fields.iter().map(|f| f.as_ref());
32993306
old_fields
33003307
.zip(new_fields)
3301-
.all(|(old, new)| is_allowed_field_change(old, new))
3308+
.try_for_each(|(old, new)| is_allowed_field_change(old, new))
33023309
}
33033310

3304-
fn is_allowed_field_change(old_field: &Field, new_field: &Field) -> bool {
3305-
new_field.name() == old_field.name()
3311+
fn is_allowed_field_change(old_field: &Field, new_field: &Field) -> Result<()> {
3312+
if new_field.name() == old_field.name()
33063313
&& new_field.data_type() == old_field.data_type()
33073314
&& new_field.metadata() == old_field.metadata()
33083315
&& (new_field.is_nullable() == old_field.is_nullable()
33093316
|| !new_field.is_nullable())
3317+
{
3318+
Ok(())
3319+
} else {
3320+
internal_err!(
3321+
"Schema field unallowed change: old field: {:?}, new field: {:?}",
3322+
old_field,
3323+
new_field
3324+
)
3325+
}
33103326
}
33113327

33123328
impl<'n> TreeNodeVisitor<'n> for OptimizationInvariantChecker<'_> {
@@ -4973,7 +4989,7 @@ digraph {
49734989
let expected_err = OptimizationInvariantChecker::new(&rule)
49744990
.check(&ok_plan, &different_schema)
49754991
.unwrap_err();
4976-
assert!(expected_err.to_string().contains("PhysicalOptimizer rule 'OptimizerRuleWithSchemaCheck' failed. Schema mismatch. Expected original schema"));
4992+
assert!(expected_err.to_string().contains("PhysicalOptimizer rule 'OptimizerRuleWithSchemaCheck' failed. Schema mismatch."));
49774993

49784994
// The recursive `check_invariants` walk only runs under `debug_assertions`
49794995
// (see `OptimizationInvariantChecker::check`). In release builds the walk is

datafusion/sqllogictest/src/test_context.rs

Lines changed: 24 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -182,6 +182,7 @@ impl TestContext {
182182
"metadata.slt" | "arrow_field.slt" => {
183183
info!("Registering metadata table tables");
184184
register_metadata_tables(test_ctx.session_ctx());
185+
register_conflicting_metadata_tables(test_ctx.session_ctx())
185186
}
186187
"union_function.slt" => {
187188
info!("Registering table with union column");
@@ -765,3 +766,26 @@ fn register_async_abs_udf(ctx: &SessionContext) {
765766
let udf = AsyncScalarUDF::new(Arc::new(async_abs));
766767
ctx.register_udf(udf.into_scalar_udf());
767768
}
769+
770+
fn register_conflicting_metadata_tables(ctx: &SessionContext) {
771+
let schema_left =
772+
Schema::new(vec![Field::new("a", DataType::Int32, false)]).with_metadata(
773+
HashMap::from([(String::from("metadata_key"), String::from("left"))]),
774+
);
775+
let data_left =
776+
Arc::new(Int32Array::from(vec![1, 2, 3, 4, 5, 6, 7, 8, 9, 10])) as ArrayRef;
777+
778+
let batch_left =
779+
RecordBatch::try_new(Arc::new(schema_left), vec![Arc::new(data_left)]).unwrap();
780+
ctx.register_batch("larger_table", batch_left).unwrap();
781+
782+
let schema_right =
783+
Schema::new(vec![Field::new("b", DataType::Int32, false)]).with_metadata(
784+
HashMap::from([(String::from("metadata_key"), String::from("right"))]),
785+
);
786+
let data_right = Arc::new(Int32Array::from(vec![1])) as ArrayRef;
787+
788+
let batch_right =
789+
RecordBatch::try_new(Arc::new(schema_right), vec![Arc::new(data_right)]).unwrap();
790+
ctx.register_batch("smaller_table", batch_right).unwrap();
791+
}

datafusion/sqllogictest/test_files/metadata.slt

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -520,3 +520,9 @@ NULL the id field
520520

521521
statement ok
522522
drop table table_with_metadata;
523+
524+
# Test that metadata on conflicting values raises an error.
525+
# The larger_table has 10 values, smaller_tables 1 value and the fields of each table
526+
# have conflicting metadata, same key different values See test:context.rs register_conflicting_metadata_tables
527+
statement error DataFusion error: PhysicalOptimizer rule 'join_selection' failed\. Schema mismatch\.\ncaused by\nInternal error: Schema metadata mismatch: Expected original metadata: \{"metadata_key": "right"\}, got metadata: \{"metadata_key": "left"\}\.
528+
select * from larger_table cross join smaller_table;

0 commit comments

Comments
 (0)