Sitelet https://github.com/apache/paimon/pull/10367/files
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -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;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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
Expand Down
Loading