This is an automated email from the ASF dual-hosted git repository.
JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git
The following commit(s) were added to refs/heads/master by this push:
new 77ffb4a279 [core] Compare binary elements and keys by value in collect
and merge_map (#9249)
77ffb4a279 is described below
commit 77ffb4a2794f52d15e417f19e00c52feff6ce025
Author: ZIHAN DAI <[email protected]>
AuthorDate: Mon Aug 17 14:19:49 2026 +1000
[core] Compare binary elements and keys by value in collect and merge_map
(#9249)
---
.../java/org/apache/paimon/utils/ByteArrayKey.java | 2 +-
.../mergetree/compact/aggregate/BinaryMapKeys.java | 64 ++++++++
.../compact/aggregate/FieldCollectAgg.java | 21 ++-
.../compact/aggregate/FieldMergeMapAgg.java | 22 ++-
.../aggregate/FieldMergeMapWithKeyTimeAgg.java | 15 +-
.../compact/aggregate/FieldAggregatorTest.java | 172 +++++++++++++++++++++
6 files changed, 281 insertions(+), 15 deletions(-)
diff --git
a/paimon-common/src/main/java/org/apache/paimon/utils/ByteArrayKey.java
b/paimon-common/src/main/java/org/apache/paimon/utils/ByteArrayKey.java
index 274e20abdc..cddf5bbba3 100644
--- a/paimon-common/src/main/java/org/apache/paimon/utils/ByteArrayKey.java
+++ b/paimon-common/src/main/java/org/apache/paimon/utils/ByteArrayKey.java
@@ -39,7 +39,7 @@ public final class ByteArrayKey {
this.hash = Arrays.hashCode(bytes);
}
- byte[] bytes() {
+ public byte[] bytes() {
return bytes;
}
diff --git
a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/aggregate/BinaryMapKeys.java
b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/aggregate/BinaryMapKeys.java
new file mode 100644
index 0000000000..5fc1eb68e4
--- /dev/null
+++
b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/aggregate/BinaryMapKeys.java
@@ -0,0 +1,64 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.paimon.mergetree.compact.aggregate;
+
+import org.apache.paimon.data.GenericMap;
+import org.apache.paimon.types.DataType;
+import org.apache.paimon.types.DataTypeFamily;
+import org.apache.paimon.utils.ByteArrayKey;
+
+import java.util.HashMap;
+import java.util.Map;
+
+/**
+ * Value semantics for map keys of type {@code BINARY} or {@code VARBINARY}.
+ *
+ * <p>Such a key arrives as a {@code byte[]}, which inherits identity equality
from {@link Object}.
+ * Used directly as a hash key, two keys with the same content occupy two
entries, a lookup never
+ * finds an existing entry and a removal never matches. Keys are therefore
held in a {@link
+ * ByteArrayKey} for as long as they are in a hash collection and unwrapped
when the result map is
+ * built.
+ */
+final class BinaryMapKeys {
+
+ private BinaryMapKeys() {}
+
+ static boolean isBinary(DataType keyType) {
+ return
keyType.getTypeRoot().getFamilies().contains(DataTypeFamily.BINARY_STRING);
+ }
+
+ /** Wrap a key for storage in a hash collection; a no-op for every
non-binary key type. */
+ static Object hashKey(boolean binaryKey, Object key) {
+ return binaryKey && key != null ? new ByteArrayKey((byte[]) key) : key;
+ }
+
+ /** Build the result map, restoring the original {@code byte[]} of any
wrapped key. */
+ static GenericMap toGenericMap(boolean binaryKey, Map<Object, Object> map)
{
+ if (!binaryKey) {
+ return new GenericMap(map);
+ }
+ Map<Object, Object> unwrapped = new HashMap<>(map.size());
+ map.forEach(
+ (key, value) ->
+ unwrapped.put(
+ key instanceof ByteArrayKey ? ((ByteArrayKey)
key).bytes() : key,
+ value));
+ return new GenericMap(unwrapped);
+ }
+}
diff --git
a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/aggregate/FieldCollectAgg.java
b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/aggregate/FieldCollectAgg.java
index 6fbd305085..368ae685a9 100644
---
a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/aggregate/FieldCollectAgg.java
+++
b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/aggregate/FieldCollectAgg.java
@@ -36,6 +36,7 @@ import java.util.Collections;
import java.util.HashSet;
import java.util.Iterator;
import java.util.List;
+import java.util.Set;
import java.util.function.BiFunction;
import static org.apache.paimon.codegen.CodeGenUtils.newRecordEqualiser;
@@ -54,11 +55,7 @@ public class FieldCollectAgg extends FieldAggregator {
this.distinct = distinct;
this.elementGetter =
InternalArray.createElementGetter(dataType.getElementType());
- if (distinct
- && dataType.getElementType()
- .getTypeRoot()
- .getFamilies()
- .contains(DataTypeFamily.CONSTRUCTED)) {
+ if (distinct && needsEqualiser(dataType.getElementType())) {
DataType elementType = dataType.getElementType();
List<DataType> fieldTypes =
elementType instanceof RowType
@@ -82,6 +79,20 @@ public class FieldCollectAgg extends FieldAggregator {
}
}
+ /**
+ * Whether elements of this type need the generated equaliser rather than
{@link Object#equals}.
+ *
+ * <p>Constructed types need it because two rows holding the same values
are not necessarily
+ * equal objects. Binary types need it for a blunter reason: an element of
{@code BINARY} or
+ * {@code VARBINARY} is a {@code byte[]}, which inherits identity equality
from {@link Object},
+ * so two arrays with the same content never compare equal and never share
a hash bucket.
+ */
+ private static boolean needsEqualiser(DataType elementType) {
+ Set<DataTypeFamily> families = elementType.getTypeRoot().getFamilies();
+ return families.contains(DataTypeFamily.CONSTRUCTED)
+ || families.contains(DataTypeFamily.BINARY_STRING);
+ }
+
@Override
public Object aggReversed(Object accumulator, Object inputField) {
// we don't need to actually do the reverse here for this agg
diff --git
a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/aggregate/FieldMergeMapAgg.java
b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/aggregate/FieldMergeMapAgg.java
index 487f20e3fd..d0cfaccbec 100644
---
a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/aggregate/FieldMergeMapAgg.java
+++
b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/aggregate/FieldMergeMapAgg.java
@@ -36,11 +36,23 @@ public class FieldMergeMapAgg extends FieldAggregator {
private final InternalArray.ElementGetter keyGetter;
private final InternalArray.ElementGetter valueGetter;
+ /** See {@link BinaryMapKeys}: a binary key has no value equality of its
own. */
+ private final boolean binaryKey;
+
public FieldMergeMapAgg(String name, MapType dataType) {
super(name, dataType);
this.keyGetter =
InternalArray.createElementGetter(dataType.getKeyType());
this.valueGetter =
InternalArray.createElementGetter(dataType.getValueType());
+ this.binaryKey = BinaryMapKeys.isBinary(dataType.getKeyType());
+ }
+
+ private Object hashKey(Object key) {
+ return BinaryMapKeys.hashKey(binaryKey, key);
+ }
+
+ private GenericMap toGenericMap(Map<Object, Object> map) {
+ return BinaryMapKeys.toGenericMap(binaryKey, map);
}
@Override
@@ -53,7 +65,7 @@ public class FieldMergeMapAgg extends FieldAggregator {
putToMap(resultMap, accumulator);
putToMap(resultMap, inputField);
- return new GenericMap(resultMap);
+ return toGenericMap(resultMap);
}
private void putToMap(Map<Object, Object> map, Object data) {
@@ -62,7 +74,7 @@ public class FieldMergeMapAgg extends FieldAggregator {
InternalArray valueArray = mapData.valueArray();
for (int i = 0; i < keyArray.size(); i++) {
map.put(
- keyGetter.getElementOrNull(keyArray, i),
+ hashKey(keyGetter.getElementOrNull(keyArray, i)),
valueGetter.getElementOrNull(valueArray, i));
}
}
@@ -86,7 +98,7 @@ public class FieldMergeMapAgg extends FieldAggregator {
InternalArray retractKeyArray = retract.keyArray();
Set<Object> retractKeys = new HashSet<>();
for (int i = 0; i < retractKeyArray.size(); i++) {
- retractKeys.add(keyGetter.getElementOrNull(retractKeyArray, i));
+
retractKeys.add(hashKey(keyGetter.getElementOrNull(retractKeyArray, i)));
}
InternalMap acc = (InternalMap) accumulator;
@@ -94,12 +106,12 @@ public class FieldMergeMapAgg extends FieldAggregator {
InternalArray accKeyArray = acc.keyArray();
InternalArray accValueArray = acc.valueArray();
for (int i = 0; i < accKeyArray.size(); i++) {
- Object accKey = keyGetter.getElementOrNull(accKeyArray, i);
+ Object accKey = hashKey(keyGetter.getElementOrNull(accKeyArray,
i));
if (!retractKeys.contains(accKey)) {
resultMap.put(accKey,
valueGetter.getElementOrNull(accValueArray, i));
}
}
- return new GenericMap(resultMap);
+ return toGenericMap(resultMap);
}
}
diff --git
a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/aggregate/FieldMergeMapWithKeyTimeAgg.java
b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/aggregate/FieldMergeMapWithKeyTimeAgg.java
index 7b16beec65..c059240682 100644
---
a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/aggregate/FieldMergeMapWithKeyTimeAgg.java
+++
b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/aggregate/FieldMergeMapWithKeyTimeAgg.java
@@ -18,7 +18,6 @@
package org.apache.paimon.mergetree.compact.aggregate;
-import org.apache.paimon.data.GenericMap;
import org.apache.paimon.data.InternalArray;
import org.apache.paimon.data.InternalMap;
import org.apache.paimon.data.InternalRow;
@@ -36,11 +35,19 @@ public class FieldMergeMapWithKeyTimeAgg extends
FieldAggregator {
private final InternalArray.ElementGetter valueGetter;
private final int timestampFieldIndex;
+ /** See {@link BinaryMapKeys}: a binary key has no value equality of its
own. */
+ private final boolean binaryKey;
+
public FieldMergeMapWithKeyTimeAgg(String name, MapType dataType, int
timestampFieldIndex) {
super(name, dataType);
this.keyGetter =
InternalArray.createElementGetter(dataType.getKeyType());
this.valueGetter =
InternalArray.createElementGetter(dataType.getValueType());
this.timestampFieldIndex = timestampFieldIndex;
+ this.binaryKey = BinaryMapKeys.isBinary(dataType.getKeyType());
+ }
+
+ private Object hashKey(Object key) {
+ return BinaryMapKeys.hashKey(binaryKey, key);
}
@Override
@@ -60,14 +67,14 @@ public class FieldMergeMapWithKeyTimeAgg extends
FieldAggregator {
mergeInputMap(resultMap, inputMap);
- return new GenericMap(resultMap);
+ return BinaryMapKeys.toGenericMap(binaryKey, resultMap);
}
private void putToMap(Map<Object, Object> map, InternalMap data) {
InternalArray keyArray = data.keyArray();
InternalArray valueArray = data.valueArray();
for (int i = 0; i < keyArray.size(); i++) {
- Object key = keyGetter.getElementOrNull(keyArray, i);
+ Object key = hashKey(keyGetter.getElementOrNull(keyArray, i));
Object value = valueGetter.getElementOrNull(valueArray, i);
map.put(key, value);
}
@@ -78,7 +85,7 @@ public class FieldMergeMapWithKeyTimeAgg extends
FieldAggregator {
InternalArray valueArray = inputMap.valueArray();
for (int i = 0; i < keyArray.size(); i++) {
- Object key = keyGetter.getElementOrNull(keyArray, i);
+ Object key = hashKey(keyGetter.getElementOrNull(keyArray, i));
InternalRow newRow = (InternalRow)
valueGetter.getElementOrNull(valueArray, i);
if (newRow == null) {
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 ef3603a7a0..6cc3c73a01 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
@@ -2104,6 +2104,60 @@ public class FieldAggregatorTest {
assertThat(unnest(result, elementGetter)).containsExactlyInAnyOrder(1,
2, 3);
}
+ /**
+ * Elements of a binary array are {@code byte[]}, which has identity
equality, so distinct
+ * collection has to compare them by content rather than dropping them
into a {@link
+ * java.util.HashSet}.
+ */
+ @Test
+ public void testFieldCollectAggWithDistinctBinary() {
+ FieldCollectAgg agg =
+ new FieldCollectAggFactory()
+ .create(
+ DataTypes.ARRAY(DataTypes.VARBINARY(10)),
+ CoreOptions.fromMap(
+
ImmutableMap.of("fields.fieldName.distinct", "true")),
+ "fieldName");
+ InternalArray.ElementGetter elementGetter =
+ InternalArray.createElementGetter(DataTypes.VARBINARY(10));
+
+ InternalArray result =
+ (InternalArray)
+ agg.agg(
+ new GenericArray(new Object[] {new byte[] {1,
2}}),
+ new GenericArray(
+ new Object[] {new byte[] {1, 2}, new
byte[] {3, 4}}));
+
+ assertThat(unnest(result, elementGetter))
+ .usingRecursiveFieldByFieldElementComparator()
+ .containsExactlyInAnyOrder(new byte[] {1, 2}, new byte[] {3,
4});
+ }
+
+ /** Retraction of a binary element must match by content too. */
+ @Test
+ public void testFieldCollectAggRetractWithDistinctBinary() {
+ FieldCollectAgg agg =
+ new FieldCollectAggFactory()
+ .create(
+ DataTypes.ARRAY(DataTypes.VARBINARY(10)),
+ CoreOptions.fromMap(
+
ImmutableMap.of("fields.fieldName.distinct", "true")),
+ "fieldName");
+ InternalArray.ElementGetter elementGetter =
+ InternalArray.createElementGetter(DataTypes.VARBINARY(10));
+
+ InternalArray result =
+ (InternalArray)
+ agg.retract(
+ new GenericArray(
+ new Object[] {new byte[] {1, 2}, new
byte[] {3, 4}}),
+ new GenericArray(new Object[] {new byte[] {1,
2}}));
+
+ assertThat(unnest(result, elementGetter))
+ .usingRecursiveFieldByFieldElementComparator()
+ .containsExactly(new byte[] {3, 4});
+ }
+
@Test
public void testFiledCollectAggWithRowType() {
RowType rowType = RowType.of(DataTypes.INT(), DataTypes.STRING());
@@ -2492,6 +2546,72 @@ public class FieldAggregatorTest {
assertThat(toJavaMap(result)).containsExactlyInAnyOrderEntriesOf(toMap(3, "C"));
}
+ /**
+ * A binary key is a {@code byte[]}, which has identity equality, so
without wrapping it the
+ * merged map keeps one entry per occurrence instead of one per distinct
key.
+ */
+ @Test
+ public void testFieldMergeMapAggWithBinaryKey() {
+ FieldMergeMapAgg agg =
+ new FieldMergeMapAggFactory()
+ .create(
+ DataTypes.MAP(DataTypes.VARBINARY(10),
DataTypes.INT()),
+ null,
+ null);
+
+ Map<Object, Object> first = new HashMap<>();
+ first.put(new byte[] {1, 2}, 1);
+ Map<Object, Object> second = new HashMap<>();
+ second.put(new byte[] {1, 2}, 2);
+ second.put(new byte[] {3, 4}, 3);
+
+ InternalMap merged = (InternalMap) agg.agg(new GenericMap(first), new
GenericMap(second));
+
+ assertThat(merged.size()).isEqualTo(2);
+ assertThat(binaryKeyed(merged)).containsOnlyKeys("0102",
"0304").containsValues(2, 3);
+ }
+
+ /** The same for retraction: a retracted binary key must match the
accumulated one. */
+ @Test
+ public void testFieldMergeMapAggRetractWithBinaryKey() {
+ FieldMergeMapAgg agg =
+ new FieldMergeMapAggFactory()
+ .create(
+ DataTypes.MAP(DataTypes.VARBINARY(10),
DataTypes.INT()),
+ null,
+ null);
+
+ Map<Object, Object> acc = new HashMap<>();
+ acc.put(new byte[] {1, 2}, 1);
+ acc.put(new byte[] {3, 4}, 2);
+ Map<Object, Object> retract = new HashMap<>();
+ retract.put(new byte[] {1, 2}, 1);
+
+ InternalMap result =
+ (InternalMap) agg.retract(new GenericMap(acc), new
GenericMap(retract));
+
+ assertThat(result.size()).isEqualTo(1);
+ assertThat(binaryKeyed(result)).containsOnlyKeys("0304");
+ }
+
+ /** Render an {@code InternalMap} with binary keys as hex so it can be
asserted by value. */
+ private Map<String, Object> binaryKeyed(InternalMap map) {
+ InternalArray.ElementGetter keyGetter =
+ InternalArray.createElementGetter(DataTypes.VARBINARY(10));
+ InternalArray.ElementGetter valueGetter =
+ InternalArray.createElementGetter(DataTypes.INT());
+ Map<String, Object> out = new HashMap<>();
+ for (int i = 0; i < map.size(); i++) {
+ byte[] key = (byte[]) keyGetter.getElementOrNull(map.keyArray(),
i);
+ StringBuilder hex = new StringBuilder();
+ for (byte b : key) {
+ hex.append(String.format("%02x", b));
+ }
+ out.put(hex.toString(),
valueGetter.getElementOrNull(map.valueArray(), i));
+ }
+ return out;
+ }
+
@Test
public void testFieldThetaSketchAgg() {
FieldThetaSketchAgg agg =
@@ -2787,6 +2907,58 @@ public class FieldAggregatorTest {
createExpectedEntry("key3", "C"));
}
+ /**
+ * 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
+ * null row fails to remove anything.
+ */
+ @Test
+ public void testFieldMergeMapWithKeyTimeAggWithBinaryKey() {
+ MapType mapType =
+ DataTypes.MAP(
+ DataTypes.VARBINARY(10),
+ DataTypes.ROW(
+ DataTypes.FIELD(0, "actual_value",
DataTypes.STRING()),
+ DataTypes.FIELD(1, "dbsync_ts",
DataTypes.STRING())));
+ FieldMergeMapWithKeyTimeAgg agg = new
FieldMergeMapWithKeyTimeAgg("test", mapType, 1);
+
+ Object acc = agg.agg(null, binaryKeyedMap(new byte[] {1, 2}, "A",
"100"));
+
+ // Newer timestamp for the same key wins, and does not become a second
entry.
+ acc = agg.agg(acc, binaryKeyedMap(new byte[] {1, 2}, "A1", "200"));
+ InternalMap merged = (InternalMap) acc;
+ assertThat(merged.size()).isEqualTo(1);
+ assertThat(firstRowValue(merged)).isEqualTo("A1");
+
+ // Older timestamp is ignored rather than appended.
+ acc = agg.agg(acc, binaryKeyedMap(new byte[] {1, 2}, "A0", "050"));
+ merged = (InternalMap) acc;
+ assertThat(merged.size()).isEqualTo(1);
+ assertThat(firstRowValue(merged)).isEqualTo("A1");
+
+ // A null row is a tombstone and must remove the entry.
+ Map<Object, Object> tombstone = new HashMap<>();
+ tombstone.put(new byte[] {1, 2}, null);
+ acc = agg.agg(acc, new GenericMap(tombstone));
+ assertThat(((InternalMap) acc).size()).isEqualTo(0);
+ }
+
+ private GenericMap binaryKeyedMap(byte[] key, String value, String ts) {
+ Map<Object, Object> map = new HashMap<>();
+ map.put(key, GenericRow.of(BinaryString.fromString(value),
BinaryString.fromString(ts)));
+ return new GenericMap(map);
+ }
+
+ private String firstRowValue(InternalMap map) {
+ InternalArray.ElementGetter valueGetter =
+ InternalArray.createElementGetter(
+ DataTypes.ROW(
+ DataTypes.FIELD(0, "actual_value",
DataTypes.STRING()),
+ DataTypes.FIELD(1, "dbsync_ts",
DataTypes.STRING())));
+ InternalRow row = (InternalRow)
valueGetter.getElementOrNull(map.valueArray(), 0);
+ return row.getString(0).toString();
+ }
+
private Map.Entry<BinaryString, InternalRow> createEntry(String key,
String value, String ts) {
return new AbstractMap.SimpleEntry<>(
BinaryString.fromString(key),