From fe2b770f3b592ad7c6932d6f85f8155ddc91c046 Mon Sep 17 00:00:00 2001 From: yangjie01 Date: Sat, 26 Sep 2026 17:07:08 +0800 Subject: [PATCH] [core] Validate the changelog format's type limits at DDL time SchemaValidation ran validateDataFields only for the base data-file format, although changelog files are written with the changelog-file format whose type limits can differ. A table with a TIMESTAMP(7-9) or TIME(4+) column and changelog-file.format=avro passed the DDL and crashed on the first changelog file write. Run the changelog format's validateDataFields over the same normal fields whenever it is set and differs from the data format. Assisted-by: GLM-5.3 --- .../paimon/schema/SchemaValidation.java | 11 ++++++ .../paimon/schema/SchemaValidationTest.java | 34 +++++++++++++++++++ 2 files changed, 45 insertions(+) diff --git a/paimon-core/src/main/java/org/apache/paimon/schema/SchemaValidation.java b/paimon-core/src/main/java/org/apache/paimon/schema/SchemaValidation.java index ed316a72238d..8d02ebc1e3c7 100644 --- a/paimon-core/src/main/java/org/apache/paimon/schema/SchemaValidation.java +++ b/paimon-core/src/main/java/org/apache/paimon/schema/SchemaValidation.java @@ -101,6 +101,7 @@ import static org.apache.paimon.CoreOptions.SNAPSHOT_NUM_RETAINED_MAX; import static org.apache.paimon.CoreOptions.SNAPSHOT_NUM_RETAINED_MIN; import static org.apache.paimon.CoreOptions.STREAMING_READ_OVERWRITE; +import static org.apache.paimon.CoreOptions.normalizeFileFormat; import static org.apache.paimon.format.FileFormat.vectorFileFormat; import static org.apache.paimon.mergetree.compact.PartialUpdateMergeFunction.isSequenceGroupOption; import static org.apache.paimon.mergetree.compact.PartialUpdateMergeFunction.isSequenceGroupOptionCandidate; @@ -277,6 +278,16 @@ public static void validateTableSchema(TableSchema schema, Set dynamicOp } fileFormat.validateDataFields(new RowType(fieldsInNormalFile)); + // changelog files are written with the changelog format, whose type limits can + // differ from the data format's + String changelogFormatIdentifier = normalizeFileFormat(options.changelogFileFormat()); + if (changelogFormatIdentifier != null + && !changelogFormatIdentifier.equalsIgnoreCase( + normalizeFileFormat(options.formatType()))) { + FileFormat.fromIdentifier(changelogFormatIdentifier, new Options(schema.options())) + .validateDataFields(new RowType(fieldsInNormalFile)); + } + for (Map.Entry entry : options.fileFormatPerLevel().entrySet()) { if (!"avro".equalsIgnoreCase(entry.getValue())) { continue; diff --git a/paimon-core/src/test/java/org/apache/paimon/schema/SchemaValidationTest.java b/paimon-core/src/test/java/org/apache/paimon/schema/SchemaValidationTest.java index 681edbf1dac0..2352820d980f 100644 --- a/paimon-core/src/test/java/org/apache/paimon/schema/SchemaValidationTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/schema/SchemaValidationTest.java @@ -74,6 +74,40 @@ void testNullablePrimaryKeyRequiresPrimaryKeyTable() { "Option 'primary-key.nullable' can only be enabled for a table with primary keys."); } + @Test + public void testChangelogFileFormatRejectsUnsupportedTypes() { + // the changelog format writes the same fields: Avro's type limits apply at DDL + // time instead of crashing the first changelog write + Map options = new HashMap<>(); + options.put(CoreOptions.CHANGELOG_FILE_FORMAT.key(), "avro"); + options.put(BUCKET.key(), String.valueOf(-1)); + TableSchema schema = + new TableSchema( + 1, + singletonList(new DataField(0, "f0", DataTypes.TIMESTAMP(9))), + 10, + emptyList(), + emptyList(), + options, + ""); + assertThatThrownBy(() -> validateTableSchema(schema)) + .hasMessageContaining("Avro does not support TIMESTAMP type with precision: 9"); + + // both formats hold the type and stay accepted, exercising the differing-format + // validation path + options.put(CoreOptions.FILE_FORMAT.key(), "parquet"); + options.put(CoreOptions.CHANGELOG_FILE_FORMAT.key(), "orc"); + validateTableSchema( + new TableSchema( + 1, + singletonList(new DataField(0, "f0", DataTypes.TIMESTAMP(9))), + 10, + emptyList(), + emptyList(), + options, + "")); + } + private void validateTableSchemaExec(Map options) { List fields = Arrays.asList(