diff --git a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/aggregate/factory/FieldMergeMapWithKeyTimeAggFactory.java b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/aggregate/factory/FieldMergeMapWithKeyTimeAggFactory.java index 452b3e473f0b..7b248909249e 100644 --- a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/aggregate/factory/FieldMergeMapWithKeyTimeAggFactory.java +++ b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/aggregate/factory/FieldMergeMapWithKeyTimeAggFactory.java @@ -21,6 +21,7 @@ import org.apache.paimon.CoreOptions; import org.apache.paimon.mergetree.compact.aggregate.FieldMergeMapWithKeyTimeAgg; import org.apache.paimon.types.DataType; +import org.apache.paimon.types.DataTypeFamily; import org.apache.paimon.types.MapType; import org.apache.paimon.types.RowType; @@ -77,6 +78,14 @@ private int resolveTsFieldIndex(RowType rowType, CoreOptions options, @Nullable field, rowType.getFieldNames()); } + + DataType tsFieldType = rowType.getTypeAt(tsFieldIndex); + checkArgument( + tsFieldType.getTypeRoot().getFamilies().contains(DataTypeFamily.CHARACTER_STRING), + "Timestamp field '%s' for field '%s' must be a string type (CHAR/VARCHAR) but was '%s'.", + rowType.getFieldNames().get(tsFieldIndex), + field, + tsFieldType); return tsFieldIndex; } diff --git a/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/aggregate/FieldAggregatorTest.java b/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/aggregate/FieldAggregatorTest.java index 820888d7a152..b915866e7f01 100644 --- a/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/aggregate/FieldAggregatorTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/aggregate/FieldAggregatorTest.java @@ -40,6 +40,7 @@ import org.apache.paimon.mergetree.compact.aggregate.factory.FieldListaggAggFactory; import org.apache.paimon.mergetree.compact.aggregate.factory.FieldMaxAggFactory; import org.apache.paimon.mergetree.compact.aggregate.factory.FieldMergeMapAggFactory; +import org.apache.paimon.mergetree.compact.aggregate.factory.FieldMergeMapWithKeyTimeAggFactory; import org.apache.paimon.mergetree.compact.aggregate.factory.FieldMinAggFactory; import org.apache.paimon.mergetree.compact.aggregate.factory.FieldNestedPartialUpdateAggFactory; import org.apache.paimon.mergetree.compact.aggregate.factory.FieldNestedUpdateAggFactory; @@ -3071,6 +3072,40 @@ public void testFieldMergeMapWithKeyTimeAgg() { createExpectedEntry("key3", "C")); } + @Test + public void testFieldMergeMapWithKeyTimeAggFactoryRejectsNonStringTsField() { + // The merge path reads the ts field via InternalRow#getString, so a non-string ts field + // (declared BIGINT/TIMESTAMP/INT, which is natural for a "ts" field) is accepted at DDL + // today and only throws ClassCastException at the first merge/compaction, naming neither + // the field nor the function. Validate it is a string type at factory creation. + MapType mapType = + DataTypes.MAP( + DataTypes.STRING(), + DataTypes.ROW( + DataTypes.FIELD(0, "actual_value", DataTypes.STRING()), + DataTypes.FIELD(1, "dbsync_ts", DataTypes.BIGINT()))); + assertThatThrownBy( + () -> + new FieldMergeMapWithKeyTimeAggFactory() + .create(mapType, new CoreOptions(new HashMap<>()), "f")) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("must be a string type"); + } + + @Test + public void testFieldMergeMapWithKeyTimeAggFactoryAcceptsStringTsField() { + MapType mapType = + DataTypes.MAP( + DataTypes.STRING(), + DataTypes.ROW( + DataTypes.FIELD(0, "actual_value", DataTypes.STRING()), + DataTypes.FIELD(1, "dbsync_ts", DataTypes.STRING()))); + assertThat( + new FieldMergeMapWithKeyTimeAggFactory() + .create(mapType, new CoreOptions(new HashMap<>()), "f")) + .isNotNull(); + } + /** * With a binary key the timestamp comparison never runs, because the lookup of the existing * entry misses: the newer row is appended as a second entry under the same logical key, and a