From 6dcef99d5970ef0c157625fbbc72424505e751da Mon Sep 17 00:00:00 2001 From: yangjie01 Date: Wed, 16 Sep 2026 12:17:40 +0800 Subject: [PATCH] [core] Reject map key type changes in schema merging merge() recursively merged the key type of a MAP column like any other nested type, so with merge-schema and type-widening enabled a MAP column could evolve into MAP. The read layer cannot cast map keys (createMapCastExecutor requires equal key types), so after such a merge every pre-change file crashed with IllegalStateException on scan, compaction, or stats read. Throw a descriptive UnsupportedOperationException when the key types differ, mirroring the method's other merge guards, so the schema change fails up front with an actionable reason instead of producing an unreadable table. Nullability is ignored (matching merge()'s contract and the read layer), so a key that only changes nullability still merges; map value types keep merging as before. --- .../paimon/schema/SchemaMergingUtils.java | 16 +++ .../paimon/schema/SchemaMergingUtilsTest.java | 103 ++++++++++++++++++ 2 files changed, 119 insertions(+) diff --git a/paimon-core/src/main/java/org/apache/paimon/schema/SchemaMergingUtils.java b/paimon-core/src/main/java/org/apache/paimon/schema/SchemaMergingUtils.java index d444d89ca4ad..bb9f42fc7978 100644 --- a/paimon-core/src/main/java/org/apache/paimon/schema/SchemaMergingUtils.java +++ b/paimon-core/src/main/java/org/apache/paimon/schema/SchemaMergingUtils.java @@ -158,6 +158,22 @@ public static DataType merge( return new RowType(base0.isNullable(), updatedFields); } else if (base instanceof MapType && update instanceof MapType) { + // The read layer cannot cast map keys (createMapCastExecutor requires equal key + // types), so widening a key here would make every pre-change file unreadable. + // Fail the schema change up front with a clear reason instead. Nullability is + // ignored, matching merge()'s contract and the read layer, so that a key that + // only changes nullability (Spark forces map keys to NOT NULL) still merges; only + // a genuine key type change is rejected. + if (!((MapType) base) + .getKeyType() + .equalsIgnoreNullable(((MapType) update).getKeyType())) { + throw new UnsupportedOperationException( + String.format( + "Failed to merge map types with different key types: %s and %s. " + + "Map key type cannot be changed; cast the keys manually " + + "and recreate the column if needed.", + base, update)); + } return new MapType( base0.isNullable(), merge( diff --git a/paimon-core/src/test/java/org/apache/paimon/schema/SchemaMergingUtilsTest.java b/paimon-core/src/test/java/org/apache/paimon/schema/SchemaMergingUtilsTest.java index 47fc100032f8..b5befc4158f4 100644 --- a/paimon-core/src/test/java/org/apache/paimon/schema/SchemaMergingUtilsTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/schema/SchemaMergingUtilsTest.java @@ -477,6 +477,109 @@ public void testMergeMapTypes() { assertThat(r3.getValueType() instanceof SmallIntType).isTrue(); } + @Test + public void testMergeMapTypesWithDifferentKeyTypes() { + AtomicInteger highestFieldId = new AtomicInteger(1); + + // widening a map key must be rejected at merge time: the read layer cannot cast map + // keys (SchemaEvolutionUtil.createMapCastExecutor requires equal key types), so a + // widened key would make every pre-change file unreadable + DataType source = new MapType(new IntType(), new VarCharType(VarCharType.MAX_LENGTH)); + DataType widenedKey = + new MapType(new BigIntType(), new VarCharType(VarCharType.MAX_LENGTH)); + assertThatThrownBy( + () -> + SchemaMergingUtils.merge( + source, widenedKey, highestFieldId, true, false, true)) + .isInstanceOf(UnsupportedOperationException.class) + .hasMessageContaining("different key types"); + + // same rejection when explicit casts are allowed: no cast can make old keys readable + assertThatThrownBy( + () -> + SchemaMergingUtils.merge( + source, widenedKey, highestFieldId, true, true, true)) + .isInstanceOf(UnsupportedOperationException.class) + .hasMessageContaining("different key types"); + + // without type widening the key change is rejected too: the old behavior silently + // kept the base key type, deferring the same crash to the write-alignment/read layer + assertThatThrownBy( + () -> + SchemaMergingUtils.merge( + source, widenedKey, highestFieldId, false, false, true)) + .isInstanceOf(UnsupportedOperationException.class) + .hasMessageContaining("different key types"); + + // a nullable map with equal keys still merges normally: nullability flows, the value + // widens, and the key check does not reject same-key merges + MapType nonNullableSource = new MapType(false, new IntType(), new IntType()); + MapType nullableSameKey = new MapType(true, new IntType(), new BigIntType()); + MapType merged = + (MapType) + SchemaMergingUtils.merge( + nonNullableSource, + nullableSameKey, + highestFieldId, + true, + false, + true); + assertThat(merged.isNullable()).isFalse(); + assertThat(merged.getKeyType() instanceof IntType).isTrue(); + assertThat(merged.getValueType() instanceof BigIntType).isTrue(); + } + + @Test + public void testMergeMapKeysDifferingOnlyInNullabilityStillMerges() { + AtomicInteger highestFieldId = new AtomicInteger(1); + + // a key that changes only nullability is not a key type change: merge ignores + // nullability and the read layer sees identical keys. Spark forces map keys to + // NOT NULL while core/Flink default to nullable, so this must stay a benign merge. + MapType nullableKey = new MapType(new IntType(), new IntType()); + MapType nonNullKey = new MapType(new IntType(false), new BigIntType()); + MapType merged = + (MapType) + SchemaMergingUtils.merge( + nullableKey, nonNullKey, highestFieldId, true, false, true); + // the base key's nullability flows to the result, so pre-change files stay readable + assertThat(merged.getKeyType() instanceof IntType).isTrue(); + assertThat(merged.getKeyType().isNullable()).isTrue(); + assertThat(merged.getValueType() instanceof BigIntType).isTrue(); + } + + @Test + public void testMergeMapKeyChangeNestedInRowIsRejected() { + AtomicInteger highestFieldId = new AtomicInteger(1); + + // the guard must fire on the recursive path too: a map key change nested inside a + // row (how a real column evolves) is rejected the same as a top-level map + RowType base = + new RowType( + Lists.newArrayList( + new DataField( + 0, + "m", + new MapType( + new IntType(), + new VarCharType(VarCharType.MAX_LENGTH))))); + RowType widenedKey = + new RowType( + Lists.newArrayList( + new DataField( + 0, + "m", + new MapType( + new BigIntType(), + new VarCharType(VarCharType.MAX_LENGTH))))); + assertThatThrownBy( + () -> + SchemaMergingUtils.merge( + base, widenedKey, highestFieldId, true, false, true)) + .isInstanceOf(UnsupportedOperationException.class) + .hasMessageContaining("different key types"); + } + @Test public void testMergeMultisetTypes() { AtomicInteger highestFieldId = new AtomicInteger(1);