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(