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),

Reply via email to