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 3de50ecc20 [format][python] Support binary keys for MAP BLOB (#9057)
3de50ecc20 is described below
commit 3de50ecc20dff1c336fcec4cf4d6f1ef6e0bd911
Author: QuakeWang <[email protected]>
AuthorDate: Fri Aug 7 13:15:31 2026 +0800
[format][python] Support binary keys for MAP BLOB (#9057)
### Purpose
`MAP<K, BLOB>` rejects `BINARY` and `VARBINARY` keys. In Java, `byte[]`
map keys also use identity equality, so equal key bytes can survive as
duplicate entries and overridden BLOB payloads can still be
externalized.
Add Java and Python binary key codecs, normalize Java binary keys by
content across map reads, copies, and primary-key externalization, and
document the byte encoding.
### Tests
- Java MAP BLOB format and map semantics tests
- Python MAP BLOB key tests
- Java/Python interoperability tests
- Checkstyle, Spotless, Flake8, license check, and docs build
Signed-off-by: QuakeWang <[email protected]>
---
docs/docs/concepts/spec/fileformat.md | 9 +-
docs/docs/multimodal-table/blob.mdx | 5 +-
docs/docs/primary-key-table/blob-storage.md | 4 +-
.../java/org/apache/paimon/data/GenericMap.java | 156 ++++++++++++++++-
.../data/serializer/InternalMapSerializer.java | 4 +-
.../org/apache/paimon/data/GenericMapTest.java | 188 +++++++++++++++++++++
.../data/serializer/InternalMapSerializerTest.java | 20 +++
.../paimon/blob/PrimaryKeyBlobExternalizer.java | 33 ++--
.../test/java/org/apache/paimon/JavaPyE2ETest.java | 42 ++++-
.../blob/PrimaryKeyBlobExternalizerTest.java | 136 ++++++++++-----
.../operation/PrimaryKeyManagedBlobStoreTest.java | 50 ++++++
.../format/blob/MapBlobElementSerializer.java | 25 ++-
.../paimon/format/blob/BlobFileFormatTest.java | 157 ++++++++++++++++-
.../pypaimon/common/map_blob_key_serializer.py | 23 +++
paimon-python/pypaimon/tests/blob_test.py | 19 +++
.../pypaimon/tests/e2e/java_py_read_write_test.py | 37 +++-
16 files changed, 828 insertions(+), 80 deletions(-)
diff --git a/docs/docs/concepts/spec/fileformat.md
b/docs/docs/concepts/spec/fileformat.md
index 2d199e6f15..252be7f8de 100644
--- a/docs/docs/concepts/spec/fileformat.md
+++ b/docs/docs/concepts/spec/fileformat.md
@@ -896,11 +896,14 @@ their encodings are:
| `DECIMAL(p, s)`, `p > 18` | Minimal-length signed big-endian
two's-complement unscaled integer |
| `DATE` | Four-byte little-endian signed count of days since 1970-01-01 |
| `TIME(p)` | Four-byte little-endian signed count of milliseconds since
midnight |
+| `BINARY`, `VARBINARY` (`BYTES`) | Raw bytes |
| `CHAR`, `VARCHAR` | UTF-8 bytes |
-The DECIMAL scale is defined by the field type and is not stored in each key.
An empty
-map has an entry count of zero and is distinct from a null map. The `TIME(p)`
encoding
-uses Paimon's millisecond internal representation and does not add nanosecond
precision.
+The DECIMAL scale is defined by the field type and is not stored in each key.
`BINARY`
+and `VARBINARY` keys are not padded, truncated, or validated against the
declared length.
+An empty map has an entry count of zero and is distinct from a null map. The
`TIME(p)`
+encoding uses Paimon's millisecond internal representation and does not add
nanosecond
+precision.
At the outer file index level, `-1` represents a null field and `-2`
represents a
field placeholder used by data evolution.
diff --git a/docs/docs/multimodal-table/blob.mdx
b/docs/docs/multimodal-table/blob.mdx
index 064ea567a5..ed92e7a832 100644
--- a/docs/docs/multimodal-table/blob.mdx
+++ b/docs/docs/multimodal-table/blob.mdx
@@ -87,8 +87,9 @@ Paimon supports three storage modes for BLOB fields, selected
via **comment dire
This allows one table to mix different storage modes for different BLOB
columns.
`ARRAY<BLOB>` and `MAP<K, BLOB>` are supported only by `__BLOB_FIELD`;
descriptor-only and blob-view comment directives accept scalar BLOB fields
only.
-Map keys support the integer family, `BOOLEAN`, `DECIMAL`, `DATE`, `TIME`,
`CHAR`, and
-`VARCHAR`. Use non-null keys for compatibility across Flink, Spark, and Python.
+Map keys support the integer family, `BOOLEAN`, `DECIMAL`, `DATE`, `TIME`,
`BINARY`,
+`VARBINARY` (`BYTES`), `CHAR`, and `VARCHAR`. Use non-null keys for
compatibility across
+Flink, Spark, and Python.
## Table Options
diff --git a/docs/docs/primary-key-table/blob-storage.md
b/docs/docs/primary-key-table/blob-storage.md
index 05acbaf3f6..871822087e 100644
--- a/docs/docs/primary-key-table/blob-storage.md
+++ b/docs/docs/primary-key-table/blob-storage.md
@@ -94,8 +94,8 @@ array order, a null array, and null elements are preserved.
An empty array write
`MAP<K, BLOB>` is externalized value by value. Keys remain in the normal data
file and every non-null value is replaced
with a descriptor to managed storage. A null map, an empty map, and null
values are preserved. Supported key types are
-the integer family, `BOOLEAN`, `DECIMAL`, `DATE`, `TIME`, `CHAR`, and
`VARCHAR`; `blob-descriptor-field` and
-`blob-view-field` remain scalar-only declarations.
+the integer family, `BOOLEAN`, `DECIMAL`, `DATE`, `TIME`, `BINARY`,
`VARBINARY` (`BYTES`), `CHAR`, and `VARCHAR`;
+`blob-descriptor-field` and `blob-view-field` remain scalar-only declarations.
`blob.target-file-size` controls when a writer rolls to a new managed payload
pack. A pack can contain payloads from
multiple rows, and a row descriptor records its URI, offset, and length.
diff --git a/paimon-common/src/main/java/org/apache/paimon/data/GenericMap.java
b/paimon-common/src/main/java/org/apache/paimon/data/GenericMap.java
index 0e07e80a5f..f70fc65300 100644
--- a/paimon-common/src/main/java/org/apache/paimon/data/GenericMap.java
+++ b/paimon-common/src/main/java/org/apache/paimon/data/GenericMap.java
@@ -23,6 +23,10 @@ import org.apache.paimon.types.MapType;
import org.apache.paimon.types.MultisetType;
import java.io.Serializable;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.LinkedHashMap;
+import java.util.List;
import java.util.Map;
import java.util.Objects;
@@ -46,6 +50,7 @@ public final class GenericMap implements InternalMap,
Serializable {
private static final long serialVersionUID = 1L;
private final Map<?, ?> map;
+ private final boolean binaryKeys;
/**
* Creates an instance of {@link GenericMap} using the given Java map.
@@ -53,7 +58,33 @@ public final class GenericMap implements InternalMap,
Serializable {
* <p>Note: All keys and values of the map must be internal data
structures.
*/
public GenericMap(Map<?, ?> map) {
- this.map = map;
+ this(map, false);
+ }
+
+ private GenericMap(Map<?, ?> map, boolean binaryKeys) {
+ this.binaryKeys = binaryKeys;
+ this.map = binaryKeys ? normalizeBinaryKeys(map) : map;
+ }
+
+ /**
+ * Creates a map whose binary keys use content equality.
+ *
+ * @since 2.1
+ */
+ public static GenericMap fromBinaryKeyMap(Map<?, ?> map) {
+ return new GenericMap(map, true);
+ }
+
+ private static Map<BinaryKey, Object> normalizeBinaryKeys(Map<?, ?> map) {
+ Map<BinaryKey, Object> binaryMap = new LinkedHashMap<>();
+ for (Map.Entry<?, ?> entry : map.entrySet()) {
+ Object key = entry.getKey();
+ if (key != null && !(key instanceof byte[])) {
+ throw new IllegalArgumentException("Binary key must be
byte[].");
+ }
+ binaryMap.put(copyBinaryKey(key), entry.getValue());
+ }
+ return binaryMap;
}
/**
@@ -61,10 +92,16 @@ public final class GenericMap implements InternalMap,
Serializable {
* no mapping for the key. The returned value is in internal data
structure.
*/
public Object get(Object key) {
+ if (binaryKeys) {
+ return isBinaryKey(key) ? map.get(lookupBinaryKey(key)) : null;
+ }
return map.get(key);
}
public boolean contains(Object key) {
+ if (binaryKeys) {
+ return isBinaryKey(key) && map.containsKey(lookupBinaryKey(key));
+ }
return map.containsKey(key);
}
@@ -75,7 +112,11 @@ public final class GenericMap implements InternalMap,
Serializable {
@Override
public InternalArray keyArray() {
- Object[] keys = map.keySet().toArray();
+ Object[] keys = new Object[map.size()];
+ int index = 0;
+ for (Object key : map.keySet()) {
+ keys[index++] = copyUnwrappedBinaryKey(key);
+ }
return new GenericArray(keys);
}
@@ -94,18 +135,50 @@ public final class GenericMap implements InternalMap,
Serializable {
return false;
}
// deepEquals for values of byte[]
- return deepEquals(map, ((GenericMap) o).map);
+ return deepEquals(this, (GenericMap) o);
+ }
+
+ private static boolean deepEquals(GenericMap m1, GenericMap m2) {
+ if (m1.map.size() != m2.map.size()) {
+ return false;
+ }
+ if ((m1.binaryKeys && m2.binaryKeys) || (!m1.hasBinaryKeys() &&
!m2.hasBinaryKeys())) {
+ return deepEquals(m1.map, m2.map);
+ }
+
+ List<Map.Entry<?, ?>> entries2 = new ArrayList<>(m2.map.entrySet());
+ boolean[] matched = new boolean[entries2.size()];
+ for (Map.Entry<?, ?> entry1 : m1.map.entrySet()) {
+ boolean found = false;
+ for (int i = 0; i < entries2.size(); i++) {
+ if (matched[i]) {
+ continue;
+ }
+ Map.Entry<?, ?> entry2 = entries2.get(i);
+ if (Objects.deepEquals(
+ unwrapBinaryKey(entry1.getKey()),
unwrapBinaryKey(entry2.getKey()))
+ && Objects.deepEquals(entry1.getValue(),
entry2.getValue())) {
+ matched[i] = true;
+ found = true;
+ break;
+ }
+ }
+ if (!found) {
+ return false;
+ }
+ }
+ return true;
}
- private static <K, V> boolean deepEquals(Map<K, V> m1, Map<?, ?> m2) {
+ private static boolean deepEquals(Map<?, ?> m1, Map<?, ?> m2) {
// copied from HashMap.equals but with deepEquals comparison
if (m1.size() != m2.size()) {
return false;
}
try {
- for (Map.Entry<K, V> e : m1.entrySet()) {
- K key = e.getKey();
- V value = e.getValue();
+ for (Map.Entry<?, ?> entry : m1.entrySet()) {
+ Object key = entry.getKey();
+ Object value = entry.getValue();
if (value == null) {
if (!(m2.get(key) == null && m2.containsKey(key))) {
return false;
@@ -126,9 +199,76 @@ public final class GenericMap implements InternalMap,
Serializable {
public int hashCode() {
int result = 0;
for (Object key : map.keySet()) {
+ key = unwrapBinaryKey(key);
// only include key because values can contain byte[]
- result += 31 * Objects.hashCode(key);
+ result +=
+ 31
+ * (key instanceof byte[]
+ ? Arrays.hashCode((byte[]) key)
+ : Objects.hashCode(key));
}
return result;
}
+
+ private boolean hasBinaryKeys() {
+ return binaryKeys || hasBinaryKey(map);
+ }
+
+ private static Object unwrapBinaryKey(Object key) {
+ return key instanceof BinaryKey ? ((BinaryKey) key).bytes : key;
+ }
+
+ private static Object copyUnwrappedBinaryKey(Object key) {
+ return key instanceof BinaryKey ? ((BinaryKey) key).copyBytes() : key;
+ }
+
+ private static boolean isBinaryKey(Object key) {
+ return key == null || key instanceof byte[];
+ }
+
+ private static BinaryKey copyBinaryKey(Object key) {
+ return key == null ? null : new BinaryKey((byte[]) key, true);
+ }
+
+ private static BinaryKey lookupBinaryKey(Object key) {
+ return key == null ? null : new BinaryKey((byte[]) key, false);
+ }
+
+ private static boolean hasBinaryKey(Map<?, ?> map) {
+ for (Object key : map.keySet()) {
+ if (key instanceof byte[]) {
+ return true;
+ }
+ }
+ return false;
+ }
+
+ private static final class BinaryKey implements Serializable {
+
+ private static final long serialVersionUID = 1L;
+
+ private final byte[] bytes;
+ private final int hash;
+
+ private BinaryKey(byte[] bytes, boolean copy) {
+ this.bytes = copy ? Arrays.copyOf(bytes, bytes.length) : bytes;
+ this.hash = Arrays.hashCode(this.bytes);
+ }
+
+ private byte[] copyBytes() {
+ return Arrays.copyOf(bytes, bytes.length);
+ }
+
+ @Override
+ public boolean equals(Object object) {
+ return object == this
+ || (object instanceof BinaryKey
+ && Arrays.equals(bytes, ((BinaryKey)
object).bytes));
+ }
+
+ @Override
+ public int hashCode() {
+ return hash;
+ }
+ }
}
diff --git
a/paimon-common/src/main/java/org/apache/paimon/data/serializer/InternalMapSerializer.java
b/paimon-common/src/main/java/org/apache/paimon/data/serializer/InternalMapSerializer.java
index 75533dd4fc..ea09ac659d 100644
---
a/paimon-common/src/main/java/org/apache/paimon/data/serializer/InternalMapSerializer.java
+++
b/paimon-common/src/main/java/org/apache/paimon/data/serializer/InternalMapSerializer.java
@@ -100,6 +100,8 @@ public class InternalMapSerializer implements
Serializer<InternalMap> {
}
private GenericMap copyBlobMap(InternalMap map) {
+ DataTypeRoot keyRoot = keyType.getTypeRoot();
+ boolean binaryKey = keyRoot == DataTypeRoot.BINARY || keyRoot ==
DataTypeRoot.VARBINARY;
Map<Object, Object> copied = new LinkedHashMap<>();
InternalArray keys = map.keyArray();
InternalArray values = map.valueArray();
@@ -110,7 +112,7 @@ public class InternalMapSerializer implements
Serializer<InternalMap> {
key == null ? null : keySerializer.copy(key),
value == null ? null : valueSerializer.copy(value));
}
- return new GenericMap(copied);
+ return binaryKey ? GenericMap.fromBinaryKeyMap(copied) : new
GenericMap(copied);
}
@Override
diff --git
a/paimon-common/src/test/java/org/apache/paimon/data/GenericMapTest.java
b/paimon-common/src/test/java/org/apache/paimon/data/GenericMapTest.java
new file mode 100644
index 0000000000..aa4f4db035
--- /dev/null
+++ b/paimon-common/src/test/java/org/apache/paimon/data/GenericMapTest.java
@@ -0,0 +1,188 @@
+/*
+ * 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.data;
+
+import org.apache.paimon.types.DataTypes;
+import org.apache.paimon.utils.InstantiationUtil;
+import org.apache.paimon.utils.InternalRowUtils;
+
+import org.junit.jupiter.api.Test;
+
+import java.util.Base64;
+import java.util.LinkedHashMap;
+import java.util.Map;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+/** Tests for {@link GenericMap}. */
+class GenericMapTest {
+
+ // Serialized by GenericMap with serialVersionUID 1L and only the map
field.
+ private static final String LEGACY_DUPLICATE_BINARY_KEY_MAP =
+
"rO0ABXNyACFvcmcuYXBhY2hlLnBhaW1vbi5kYXRhLkdlbmVyaWNNYXAAAAAAAAAAAQIAAUwAA21hcHQA"
+ +
"D0xqYXZhL3V0aWwvTWFwO3hwc3IAF2phdmEudXRpbC5MaW5rZWRIYXNoTWFwNMBOXBBswPsCAAFaAAth"
+ +
"Y2Nlc3NPcmRlcnhyABFqYXZhLnV0aWwuSGFzaE1hcAUH2sHDFmDRAwACRgAKbG9hZEZhY3RvckkACXRo"
+ +
"cmVzaG9sZHhwP0AAAAAAAAx3CAAAABAAAAACdXIAAltCrPMX+AYIVOACAAB4cAAAAAEBc3IAEWphdmEu"
+ +
"bGFuZy5JbnRlZ2VyEuKgpPeBhzgCAAFJAAV2YWx1ZXhyABBqYXZhLmxhbmcuTnVtYmVyhqyVHQuU4IsC"
+ + "AAB4cAAAAAF1cQB+AAYAAAABAXNxAH4ACAAAAAJ4AA==";
+
+ @Test
+ void testBinarySemanticsDoNotDependOnInitialContents() throws Exception {
+ Map<Object, Object> initiallyEmptyEntries = new LinkedHashMap<>();
+ GenericMap initiallyEmpty = new GenericMap(initiallyEmptyEntries);
+ byte[] initiallyEmptyKey = new byte[] {1};
+ initiallyEmptyEntries.put(initiallyEmptyKey, 2);
+
+ Map<Object, Object> populatedEntries = new LinkedHashMap<>();
+ byte[] populatedKey = new byte[] {1};
+ populatedEntries.put(populatedKey, 2);
+ GenericMap initiallyPopulated = new GenericMap(populatedEntries);
+ GenericMap binary = GenericMap.fromBinaryKeyMap(populatedEntries);
+
+ assertThat(initiallyEmpty.contains(initiallyEmptyKey)).isTrue();
+ assertThat(initiallyEmpty.contains(new byte[] {1})).isFalse();
+ assertThat(initiallyPopulated.contains(populatedKey)).isTrue();
+ assertThat(initiallyPopulated.contains(new byte[] {1})).isFalse();
+
+
assertThat(initiallyEmpty).isEqualTo(initiallyPopulated).isEqualTo(binary);
+
assertThat(initiallyPopulated).isEqualTo(initiallyEmpty).isEqualTo(binary);
+
assertThat(binary).isEqualTo(initiallyEmpty).isEqualTo(initiallyPopulated);
+ assertThat(initiallyEmpty.hashCode())
+ .isEqualTo(initiallyPopulated.hashCode())
+ .isEqualTo(binary.hashCode());
+
+ GenericMap restored = InstantiationUtil.clone(initiallyPopulated);
+ assertThat(restored.contains(new byte[] {1})).isFalse();
+ assertThat(restored).isEqualTo(binary);
+ assertThat(binary).isEqualTo(restored);
+ assertThat(restored.hashCode()).isEqualTo(binary.hashCode());
+ }
+
+ @Test
+ void testOrdinaryDuplicateBinaryKeysPreservePhysicalEntries() {
+ Map<Object, Object> entries = new LinkedHashMap<>();
+ entries.put(new byte[] {1}, 1);
+ entries.put(new byte[] {1}, 2);
+ GenericMap ordinary = new GenericMap(entries);
+
+ Map<Object, Object> sameEntries = new LinkedHashMap<>();
+ sameEntries.put(new byte[] {1}, 2);
+ sameEntries.put(new byte[] {1}, 1);
+ GenericMap same = new GenericMap(sameEntries);
+
+ GenericMap normalized = GenericMap.fromBinaryKeyMap(entries);
+
+ assertThat(ordinary.size()).isEqualTo(2);
+ assertThat(ordinary.keyArray().size()).isEqualTo(2);
+ assertThat(ordinary.contains(new byte[] {1})).isFalse();
+ assertThat(normalized.size()).isOne();
+ assertThat(normalized.keyArray().size()).isOne();
+ assertThat(normalized.contains(new byte[] {1})).isTrue();
+
+ assertThat(ordinary).isEqualTo(same);
+ assertThat(same).isEqualTo(ordinary);
+ assertThat(ordinary.hashCode()).isEqualTo(same.hashCode());
+ assertThat(ordinary).isNotEqualTo(normalized);
+ assertThat(normalized).isNotEqualTo(ordinary);
+ assertThat(ordinary.hashCode()).isNotEqualTo(normalized.hashCode());
+
assertThat(GenericRow.of(ordinary)).isNotEqualTo(GenericRow.of(normalized));
+
assertThat(GenericRow.of(normalized)).isNotEqualTo(GenericRow.of(ordinary));
+ assertThat(
+ InternalRowUtils.equals(
+ ordinary,
+ normalized,
+ DataTypes.MAP(DataTypes.BYTES(),
DataTypes.INT())))
+ .isFalse();
+ }
+
+ @Test
+ void testDeserializeLegacyMapPreservesDuplicateBinaryKeys() throws
Exception {
+ GenericMap legacy =
+ InstantiationUtil.deserializeObject(
+
Base64.getDecoder().decode(LEGACY_DUPLICATE_BINARY_KEY_MAP),
+ GenericMap.class.getClassLoader());
+
+ assertThat(legacy.size()).isEqualTo(2);
+ assertThat(legacy.keyArray().size()).isEqualTo(2);
+ assertThat(legacy.keyArray().getBinary(0)).isEqualTo(new byte[] {1});
+ assertThat(legacy.keyArray().getBinary(1)).isEqualTo(new byte[] {1});
+ assertThat(legacy.valueArray().getInt(0)).isEqualTo(1);
+ assertThat(legacy.valueArray().getInt(1)).isEqualTo(2);
+ assertThat(legacy.contains(new byte[] {1})).isFalse();
+ assertThat(legacy.get(new byte[] {1})).isNull();
+ }
+
+ @Test
+ void testDuplicateBinaryKeysPreserveEqualsContract() {
+ GenericMap left =
+ binaryMap(new byte[][] {new byte[] {1}, new byte[] {1}}, new
Object[] {1, 2});
+ GenericMap right =
+ binaryMap(new byte[][] {new byte[] {1}, new byte[] {1}}, new
Object[] {2, 2});
+ GenericMap canonical = binaryMap(new byte[][] {new byte[] {1}}, new
Object[] {2});
+
+ assertThat(left.size()).isOne();
+ assertThat(left.get(new byte[] {1})).isEqualTo(2);
+ assertThat(left).isEqualTo(right);
+ assertThat(right).isEqualTo(left);
+ assertThat(right).isEqualTo(canonical);
+ assertThat(canonical).isEqualTo(right);
+ assertThat(left).isEqualTo(canonical);
+
assertThat(left.hashCode()).isEqualTo(right.hashCode()).isEqualTo(canonical.hashCode());
+
+ GenericRow leftRow = GenericRow.of(left);
+ GenericRow rightRow = GenericRow.of(right);
+ GenericRow canonicalRow = GenericRow.of(canonical);
+ assertThat(leftRow).isEqualTo(rightRow);
+ assertThat(rightRow).isEqualTo(leftRow);
+ assertThat(rightRow).isEqualTo(canonicalRow);
+ assertThat(leftRow).isEqualTo(canonicalRow);
+ assertThat(leftRow.hashCode())
+ .isEqualTo(rightRow.hashCode())
+ .isEqualTo(canonicalRow.hashCode());
+ }
+
+ @Test
+ void testBinaryKeyOwnershipIsIsolated() throws Exception {
+ byte[] key = new byte[] {1};
+ GenericMap map = binaryMap(new byte[][] {key}, new Object[] {2});
+
+ key[0] = 2;
+ assertThat(map.contains(new byte[] {1})).isTrue();
+ assertThat(map.contains(new byte[] {2})).isFalse();
+
+ byte[] exposed = map.keyArray().getBinary(0);
+ exposed[0] = 3;
+ assertThat(map.contains(new byte[] {1})).isTrue();
+ assertThat(map.contains(new byte[] {3})).isFalse();
+ assertThat(map.keyArray().getBinary(0)).isEqualTo(new byte[] {1});
+
+ GenericMap restored = InstantiationUtil.clone(map);
+ assertThat(restored.contains(new byte[] {1})).isTrue();
+ assertThat(restored).isEqualTo(map);
+ assertThat(restored.hashCode()).isEqualTo(map.hashCode());
+ }
+
+ private static GenericMap binaryMap(byte[][] keys, Object[] values) {
+ Map<Object, Object> entries = new LinkedHashMap<>();
+ for (int i = 0; i < keys.length; i++) {
+ entries.put(keys[i], values[i]);
+ }
+ return GenericMap.fromBinaryKeyMap(entries);
+ }
+}
diff --git
a/paimon-common/src/test/java/org/apache/paimon/data/serializer/InternalMapSerializerTest.java
b/paimon-common/src/test/java/org/apache/paimon/data/serializer/InternalMapSerializerTest.java
index d4892777f6..aff3a75e8d 100644
---
a/paimon-common/src/test/java/org/apache/paimon/data/serializer/InternalMapSerializerTest.java
+++
b/paimon-common/src/test/java/org/apache/paimon/data/serializer/InternalMapSerializerTest.java
@@ -22,6 +22,7 @@ import org.apache.paimon.data.BinaryArray;
import org.apache.paimon.data.BinaryArrayWriter;
import org.apache.paimon.data.BinaryMap;
import org.apache.paimon.data.BinaryString;
+import org.apache.paimon.data.BlobData;
import org.apache.paimon.data.BlobDescriptor;
import org.apache.paimon.data.GenericArray;
import org.apache.paimon.data.GenericMap;
@@ -40,11 +41,13 @@ import org.junit.jupiter.api.io.TempDir;
import java.nio.charset.StandardCharsets;
import java.util.Arrays;
import java.util.HashMap;
+import java.util.LinkedHashMap;
import java.util.Map;
import java.util.Objects;
import static
org.apache.paimon.data.serializer.InternalMapSerializer.convertToJavaMap;
import static org.apache.paimon.types.DataTypes.BLOB;
+import static org.apache.paimon.types.DataTypes.BYTES;
import static org.apache.paimon.types.DataTypes.INT;
import static org.apache.paimon.types.DataTypes.STRING;
import static org.assertj.core.api.Assertions.assertThat;
@@ -141,6 +144,23 @@ public class InternalMapSerializerTest extends
SerializerTestBase<InternalMap> {
assertThat(copied.valueArray().isNullAt(1)).isTrue();
}
+ @Test
+ void testCopyBinaryBlobMapPreservesContentEquality() {
+ Map<Object, Object> entries = new LinkedHashMap<>();
+ entries.put(new byte[] {1}, new
BlobData("first".getBytes(StandardCharsets.UTF_8)));
+ entries.put(new byte[] {1}, new
BlobData("second".getBytes(StandardCharsets.UTF_8)));
+
+ GenericMap copied =
+ (GenericMap)
+ new InternalMapSerializer(BYTES(), BLOB()).copy(new
GenericMap(entries));
+
+ assertThat(copied.size()).isOne();
+ assertThat(copied.contains(new byte[] {1})).isTrue();
+ assertThat(((BlobData) copied.get(new byte[] {1})).toData())
+ .isEqualTo("second".getBytes(StandardCharsets.UTF_8));
+ assertThat(copied.keyArray().getBinary(0)).isEqualTo(new byte[] {1});
+ }
+
private static BinaryArray createArray(int... vs) {
BinaryArray array = new BinaryArray();
BinaryArrayWriter writer = new BinaryArrayWriter(array, vs.length, 4);
diff --git
a/paimon-core/src/main/java/org/apache/paimon/blob/PrimaryKeyBlobExternalizer.java
b/paimon-core/src/main/java/org/apache/paimon/blob/PrimaryKeyBlobExternalizer.java
index c48c971e64..aa46e83972 100644
---
a/paimon-core/src/main/java/org/apache/paimon/blob/PrimaryKeyBlobExternalizer.java
+++
b/paimon-core/src/main/java/org/apache/paimon/blob/PrimaryKeyBlobExternalizer.java
@@ -223,37 +223,48 @@ public class PrimaryKeyBlobExternalizer {
keys.size() == map.size() && values.size() == map.size(),
"MAP<X, BLOB> key/value array size does not match map size.");
- Map<Object, Blob> blobs = new LinkedHashMap<>();
+ Map<Object, Object> copied = new LinkedHashMap<>();
for (int i = 0; i < map.size(); i++) {
Object key =
InternalRowUtils.copy(keyGetter.getElementOrNull(keys, i), keyType);
- blobs.put(key, values.isNullAt(i) ? null : values.getBlob(i));
+ copied.put(key, values.isNullAt(i) ? null : values.getBlob(i));
}
+ GenericMap blobs = createBlobMap(copied, keyType);
+ InternalArray normalizedKeys = blobs.keyArray();
+ InternalArray normalizedValues = blobs.valueArray();
boolean hasBlob = false;
- for (Blob blob : blobs.values()) {
- if (blob != null) {
+ for (int i = 0; i < blobs.size(); i++) {
+ if (!normalizedValues.isNullAt(i)) {
hasBlob = true;
break;
}
}
if (!hasBlob) {
- return blobs.size() == map.size() ? null : new GenericMap(blobs);
+ return blobs.size() == map.size() ? null : blobs;
}
Map<Object, Object> externalized = new LinkedHashMap<>();
- for (Map.Entry<Object, Blob> entry : blobs.entrySet()) {
- Blob blob = entry.getValue();
- if (blob == null) {
- externalized.put(entry.getKey(), null);
+ for (int i = 0; i < blobs.size(); i++) {
+ Object key = keyGetter.getElementOrNull(normalizedKeys, i);
+ if (normalizedValues.isNullAt(i)) {
+ externalized.put(key, null);
continue;
}
+ Blob blob = normalizedValues.getBlob(i);
BlobDescriptor descriptor = packWriter.write(blob);
externalized.put(
- entry.getKey(),
+ key,
Blob.fromFile(
fileIO, descriptor.uri(), descriptor.offset(),
descriptor.length()));
}
- return new GenericMap(externalized);
+ return createBlobMap(externalized, keyType);
+ }
+
+ private static GenericMap createBlobMap(Map<?, ?> map, DataType keyType) {
+ DataTypeRoot keyRoot = keyType.getTypeRoot();
+ return keyRoot == DataTypeRoot.BINARY || keyRoot ==
DataTypeRoot.VARBINARY
+ ? GenericMap.fromBinaryKeyMap(map)
+ : new GenericMap(map);
}
public void prepareCommit() throws IOException {
diff --git a/paimon-core/src/test/java/org/apache/paimon/JavaPyE2ETest.java
b/paimon-core/src/test/java/org/apache/paimon/JavaPyE2ETest.java
index 81188e64a9..69de8c06f9 100644
--- a/paimon-core/src/test/java/org/apache/paimon/JavaPyE2ETest.java
+++ b/paimon-core/src/test/java/org/apache/paimon/JavaPyE2ETest.java
@@ -27,6 +27,7 @@ import org.apache.paimon.catalog.Identifier;
import org.apache.paimon.data.BinaryRow;
import org.apache.paimon.data.BinaryString;
import org.apache.paimon.data.BinaryVector;
+import org.apache.paimon.data.Blob;
import org.apache.paimon.data.BlobData;
import org.apache.paimon.data.DataFormatTestUtil;
import org.apache.paimon.data.Decimal;
@@ -1460,6 +1461,12 @@ public class JavaPyE2ETest {
DataTypes.MAP(DataTypes.DECIMAL(20, 2),
DataTypes.BLOB()))
.column("date_payloads",
DataTypes.MAP(DataTypes.DATE(), DataTypes.BLOB()))
.column("time_payloads",
DataTypes.MAP(DataTypes.TIME(3), DataTypes.BLOB()))
+ .column(
+ "binary_payloads",
+ DataTypes.MAP(DataTypes.BINARY(4),
DataTypes.BLOB()))
+ .column(
+ "varbinary_payloads",
+ DataTypes.MAP(DataTypes.VARBINARY(8),
DataTypes.BLOB()))
.option(ROW_TRACKING_ENABLED.key(), "true")
.option(DATA_EVOLUTION_ENABLED.key(), "true")
.option(BUCKET.key(), "-1")
@@ -1486,6 +1493,16 @@ public class JavaPyE2ETest {
datePayloads.put(-1, new
BlobData("java-date".getBytes(StandardCharsets.UTF_8)));
Map<Object, Object> timePayloads = new LinkedHashMap<>();
timePayloads.put(45_296_789, new
BlobData("java-time".getBytes(StandardCharsets.UTF_8)));
+ Map<Object, Object> binaryPayloads = new LinkedHashMap<>();
+ binaryPayloads.put(
+ new byte[] {0, (byte) 0xff, 1, 2},
+ new
BlobData("java-binary-first".getBytes(StandardCharsets.UTF_8)));
+ binaryPayloads.put(
+ new byte[] {0, (byte) 0xff, 1, 2},
+ new BlobData("java-binary".getBytes(StandardCharsets.UTF_8)));
+ Map<Object, Object> varbinaryPayloads = new LinkedHashMap<>();
+ varbinaryPayloads.put(
+ new byte[0], new
BlobData("java-varbinary".getBytes(StandardCharsets.UTF_8)));
FileStoreTable table = (FileStoreTable) catalog.getTable(identifier);
BatchWriteBuilder writeBuilder = table.newBatchWriteBuilder();
@@ -1499,7 +1516,9 @@ public class JavaPyE2ETest {
new GenericMap(compactDecimalPayloads),
new GenericMap(highDecimalPayloads),
new GenericMap(datePayloads),
- new GenericMap(timePayloads)));
+ new GenericMap(timePayloads),
+ new GenericMap(binaryPayloads),
+ new GenericMap(varbinaryPayloads)));
write.write(
GenericRow.of(
2,
@@ -1508,9 +1527,13 @@ public class JavaPyE2ETest {
null,
null,
null,
+ null,
+ null,
null));
- write.write(GenericRow.of(3, null, null, null, null, null, null));
- write.write(GenericRow.of(4, new GenericMap(last), null, null,
null, null, null));
+ write.write(GenericRow.of(3, null, null, null, null, null, null,
null, null));
+ write.write(
+ GenericRow.of(
+ 4, new GenericMap(last), null, null, null, null,
null, null, null));
commit.commit(write.prepareCommit());
}
@@ -1604,6 +1627,19 @@ public class JavaPyE2ETest {
InternalMap timeMap = row.getMap(6);
assertThat(timeMap.keyArray().getInt(0)).isEqualTo(45_296_789);
assertSingleBlobValue(timeMap, valuePrefix + "-time");
+
+ GenericMap binaryMap = (GenericMap) row.getMap(7);
+ byte[] binaryKey = new byte[] {0, (byte) 0xff, 1, 2};
+
assertThat(binaryMap.keyArray().getBinary(0)).isEqualTo(binaryKey);
+ assertThat(binaryMap.contains(binaryKey)).isTrue();
+ assertThat(((Blob) binaryMap.get(binaryKey)).toData())
+ .isEqualTo(
+ (valuePrefix +
"-binary").getBytes(StandardCharsets.UTF_8));
+ assertThat(binaryMap.size()).isOne();
+
+ InternalMap varbinaryMap = row.getMap(8);
+
assertThat(varbinaryMap.keyArray().getBinary(0)).isEmpty();
+ assertSingleBlobValue(varbinaryMap, valuePrefix +
"-varbinary");
});
}
assertThat(found[0]).isTrue();
diff --git
a/paimon-core/src/test/java/org/apache/paimon/blob/PrimaryKeyBlobExternalizerTest.java
b/paimon-core/src/test/java/org/apache/paimon/blob/PrimaryKeyBlobExternalizerTest.java
index 04c09b0a94..caf4b61a16 100644
---
a/paimon-core/src/test/java/org/apache/paimon/blob/PrimaryKeyBlobExternalizerTest.java
+++
b/paimon-core/src/test/java/org/apache/paimon/blob/PrimaryKeyBlobExternalizerTest.java
@@ -459,26 +459,12 @@ class PrimaryKeyBlobExternalizerTest {
1L);
byte[] expected = "last".getBytes(StandardCharsets.UTF_8);
InternalMap duplicateKeys =
- new InternalMap() {
- @Override
- public int size() {
- return 2;
- }
-
- @Override
- public InternalArray keyArray() {
- return new GenericArray(new Object[] {1, 1});
- }
-
- @Override
- public InternalArray valueArray() {
- return new GenericArray(
- new Object[] {
-
Blob.fromData("overridden".getBytes(StandardCharsets.UTF_8)),
- Blob.fromData(expected)
- });
- }
- };
+ duplicateMap(
+ new Object[] {1, 1},
+ new Object[] {
+
Blob.fromData("overridden".getBytes(StandardCharsets.UTF_8)),
+ Blob.fromData(expected)
+ });
InternalMap result =
externalizer.externalize(RowKind.INSERT,
GenericRow.of(duplicateKeys)).getMap(0);
@@ -509,32 +495,85 @@ class PrimaryKeyBlobExternalizerTest {
pathFactory,
1024L);
InternalMap duplicateKeys =
- new InternalMap() {
- @Override
- public int size() {
- return 2;
- }
+ duplicateMap(
+ new Object[] {1, 1},
+ new Object[] {
+
Blob.fromData("overridden".getBytes(StandardCharsets.UTF_8)), null
+ });
- @Override
- public InternalArray keyArray() {
- return new GenericArray(new Object[] {1, 1});
- }
+ InternalMap result =
+ externalizer.externalize(RowKind.INSERT,
GenericRow.of(duplicateKeys)).getMap(0);
- @Override
- public InternalArray valueArray() {
- return new GenericArray(
- new Object[] {
-
Blob.fromData("overridden".getBytes(StandardCharsets.UTF_8)),
- null
- });
- }
- };
+ assertThat(result.size()).isEqualTo(1);
+ assertThat(result.keyArray().getInt(0)).isEqualTo(1);
+ assertThat(result.valueArray().isNullAt(0)).isTrue();
+ assertThat(fileIO.listStatus(bucketPath)).isEmpty();
+ }
+
+ @Test
+ void testDuplicateBinaryMapKeyUsesLastValueWithoutWritingOverriddenBlob()
throws Exception {
+ LocalFileIO fileIO = LocalFileIO.create();
+ Path bucketPath = new Path(tempDir.resolve("bucket-0").toUri());
+ fileIO.mkdirs(bucketPath);
+ DataFilePathFactory pathFactory =
+ new DataFilePathFactory(
+ bucketPath, "avro", "data-", "changelog-", false,
null, null);
+ PrimaryKeyBlobExternalizer externalizer =
+ newExternalizer(
+ fileIO,
+ RowType.of(DataTypes.MAP(DataTypes.BYTES(),
DataTypes.BLOB())),
+ Collections.singleton("f0"),
+ pathFactory,
+ 1L);
+ byte[] expected = "last".getBytes(StandardCharsets.UTF_8);
+ InternalMap duplicateKeys =
+ duplicateMap(
+ new Object[] {new byte[] {1}, new byte[] {1}},
+ new Object[] {
+
Blob.fromData("overridden".getBytes(StandardCharsets.UTF_8)),
+ Blob.fromData(expected)
+ });
InternalMap result =
externalizer.externalize(RowKind.INSERT,
GenericRow.of(duplicateKeys)).getMap(0);
assertThat(result.size()).isEqualTo(1);
- assertThat(result.keyArray().getInt(0)).isEqualTo(1);
+ assertThat(result.keyArray().getBinary(0)).isEqualTo(new byte[] {1});
+
assertThat(result.valueArray().getBlob(0).toData()).isEqualTo(expected);
+ assertThat(fileIO.listStatus(bucketPath))
+ .singleElement()
+ .extracting(status -> status.getPath().getName())
+ .asString()
+ .endsWith(ManagedBlobReferenceFile.MANAGED_BLOB_SUFFIX);
+ }
+
+ @Test
+ void testDuplicateBinaryMapKeyUsesLastNullWithoutWritingBlob() throws
Exception {
+ LocalFileIO fileIO = LocalFileIO.create();
+ Path bucketPath = new Path(tempDir.resolve("bucket-0").toUri());
+ fileIO.mkdirs(bucketPath);
+ DataFilePathFactory pathFactory =
+ new DataFilePathFactory(
+ bucketPath, "avro", "data-", "changelog-", false,
null, null);
+ PrimaryKeyBlobExternalizer externalizer =
+ newExternalizer(
+ fileIO,
+ RowType.of(DataTypes.MAP(DataTypes.VARBINARY(8),
DataTypes.BLOB())),
+ Collections.singleton("f0"),
+ pathFactory,
+ 1024L);
+ InternalMap duplicateKeys =
+ duplicateMap(
+ new Object[] {new byte[] {1}, new byte[] {1}},
+ new Object[] {
+
Blob.fromData("overridden".getBytes(StandardCharsets.UTF_8)), null
+ });
+
+ InternalMap result =
+ externalizer.externalize(RowKind.INSERT,
GenericRow.of(duplicateKeys)).getMap(0);
+
+ assertThat(result.size()).isEqualTo(1);
+ assertThat(result.keyArray().getBinary(0)).isEqualTo(new byte[] {1});
assertThat(result.valueArray().isNullAt(0)).isTrue();
assertThat(fileIO.listStatus(bucketPath)).isEmpty();
}
@@ -632,4 +671,23 @@ class PrimaryKeyBlobExternalizerTest {
targetFileSize,
BlobFormatWriter.DEFAULT_COPY_BUFFER_SIZE);
}
+
+ private static InternalMap duplicateMap(Object[] keys, Object[] values) {
+ return new InternalMap() {
+ @Override
+ public int size() {
+ return keys.length;
+ }
+
+ @Override
+ public InternalArray keyArray() {
+ return new GenericArray(keys);
+ }
+
+ @Override
+ public InternalArray valueArray() {
+ return new GenericArray(values);
+ }
+ };
+ }
}
diff --git
a/paimon-core/src/test/java/org/apache/paimon/operation/PrimaryKeyManagedBlobStoreTest.java
b/paimon-core/src/test/java/org/apache/paimon/operation/PrimaryKeyManagedBlobStoreTest.java
index 6001eb6f5c..01a0455e0f 100644
---
a/paimon-core/src/test/java/org/apache/paimon/operation/PrimaryKeyManagedBlobStoreTest.java
+++
b/paimon-core/src/test/java/org/apache/paimon/operation/PrimaryKeyManagedBlobStoreTest.java
@@ -303,6 +303,56 @@ class PrimaryKeyManagedBlobStoreTest {
assertThat(actual.get("second").toData()).isEqualTo(second);
}
+ @Test
+ void testExternalizeAndReadDuplicateBinaryKeyBlobMap() throws Exception {
+ FileIO fileIO = LocalFileIO.create();
+ TestFileStore store =
+ createStore(fileIO, "payloads",
DataTypes.MAP(DataTypes.BYTES(), DataTypes.BLOB()));
+ byte[] expected = "last-map-payload".getBytes(StandardCharsets.UTF_8);
+ InternalMap input =
+ new InternalMap() {
+ @Override
+ public int size() {
+ return 2;
+ }
+
+ @Override
+ public InternalArray keyArray() {
+ return new GenericArray(new Object[] {new byte[] {1},
new byte[] {1}});
+ }
+
+ @Override
+ public InternalArray valueArray() {
+ return new GenericArray(
+ new Object[] {
+ Blob.fromData(
+ "overridden-map-payload"
+
.getBytes(StandardCharsets.UTF_8)),
+ Blob.fromData(expected)
+ });
+ }
+ };
+
+ store.commitData(
+ Collections.singletonList(
+ new KeyValue()
+ .replace(
+ GenericRow.of(1), RowKind.INSERT,
GenericRow.of(1, input))),
+ ignored -> BinaryRow.EMPTY_ROW,
+ ignored -> 0);
+
+ ManifestEntry entry = store.newScan().plan().files().get(0);
+ assertThat(references(fileIO, store, entry)).hasSize(1);
+ InternalMap result =
+
store.readKvsFromSnapshot(store.snapshotManager().latestSnapshotId())
+ .get(0)
+ .value()
+ .getMap(1);
+ assertThat(result.size()).isOne();
+ assertThat(result.keyArray().getBinary(0)).isEqualTo(new byte[] {1});
+
assertThat(result.valueArray().getBlob(0).toData()).isEqualTo(expected);
+ }
+
@Test
void testCompactionRebuildsExactBlobReferences() throws Exception {
FileIO fileIO = LocalFileIO.create();
diff --git
a/paimon-format/src/main/java/org/apache/paimon/format/blob/MapBlobElementSerializer.java
b/paimon-format/src/main/java/org/apache/paimon/format/blob/MapBlobElementSerializer.java
index 6d73bfc5d6..b9fb6a74b9 100644
---
a/paimon-format/src/main/java/org/apache/paimon/format/blob/MapBlobElementSerializer.java
+++
b/paimon-format/src/main/java/org/apache/paimon/format/blob/MapBlobElementSerializer.java
@@ -350,6 +350,7 @@ final class MapBlobElementSerializer implements
BlobElementSerializer {
}
// 3. deserialize values and construct map
+ boolean binaryKey = keySerializer instanceof
BinaryKeySerializer;
Map<Object, Object> map = new LinkedHashMap<>();
long valueOffset = dataStart + keyDataLength;
for (int i = 0; i < entryCount; i++) {
@@ -363,7 +364,7 @@ final class MapBlobElementSerializer implements
BlobElementSerializer {
}
map.put(keys[i], value);
}
- return new GenericMap(map);
+ return binaryKey ? GenericMap.fromBinaryKeyMap(map) : new
GenericMap(map);
} catch (IOException e) {
throw new RuntimeException(e);
}
@@ -470,6 +471,9 @@ final class MapBlobElementSerializer implements
BlobElementSerializer {
case DATE:
case TIME_WITHOUT_TIME_ZONE:
return new IntKeySerializer();
+ case BINARY:
+ case VARBINARY:
+ return new BinaryKeySerializer();
case CHAR:
case VARCHAR:
return new StringKeySerializer();
@@ -645,6 +649,25 @@ final class MapBlobElementSerializer implements
BlobElementSerializer {
}
}
+ /** {@link KeySerializer} for Binary and VarBinary Types. */
+ private static final class BinaryKeySerializer implements KeySerializer {
+
+ @Override
+ public byte[] serialize(Object key) {
+ return (byte[]) key;
+ }
+
+ @Override
+ public Object deserialize(byte[] bytes) {
+ return bytes;
+ }
+
+ @Override
+ public int fixedLength() {
+ return -1;
+ }
+ }
+
/** {@link KeySerializer} for String Type. */
private static final class StringKeySerializer implements KeySerializer {
diff --git
a/paimon-format/src/test/java/org/apache/paimon/format/blob/BlobFileFormatTest.java
b/paimon-format/src/test/java/org/apache/paimon/format/blob/BlobFileFormatTest.java
index a155100b24..bd11580bb0 100644
---
a/paimon-format/src/test/java/org/apache/paimon/format/blob/BlobFileFormatTest.java
+++
b/paimon-format/src/test/java/org/apache/paimon/format/blob/BlobFileFormatTest.java
@@ -30,7 +30,10 @@ import org.apache.paimon.data.GenericArray;
import org.apache.paimon.data.GenericMap;
import org.apache.paimon.data.GenericRow;
import org.apache.paimon.data.InternalArray;
+import org.apache.paimon.data.InternalMap;
import org.apache.paimon.data.InternalRow;
+import org.apache.paimon.data.serializer.InternalMapSerializer;
+import org.apache.paimon.data.serializer.InternalRowSerializer;
import org.apache.paimon.format.FormatReaderContext;
import org.apache.paimon.format.FormatReaderFactory;
import org.apache.paimon.format.FormatWriter;
@@ -44,6 +47,7 @@ import org.apache.paimon.reader.FileRecordReader;
import org.apache.paimon.types.DataType;
import org.apache.paimon.types.DataTypes;
import org.apache.paimon.types.RowType;
+import org.apache.paimon.utils.InstantiationUtil;
import org.apache.paimon.utils.ProjectedRow;
import org.apache.paimon.utils.RoaringBitmap32;
@@ -439,22 +443,142 @@ public class BlobFileFormatTest {
@Test
public void testDuplicateMapBlobKeyLastWinsInline() throws IOException {
- assertDuplicateMapBlobKeyLastWins(false);
+ assertDuplicateMapBlobKeyLastWins(
+ false,
+ DataTypes.STRING(),
+ BinaryString.fromString("a"),
+ BinaryString.fromString("b"),
+ BinaryString.fromString("a"),
+ (byte) 'a');
}
@Test
public void testDuplicateMapBlobKeyLastWinsAsDescriptor() throws
IOException {
- assertDuplicateMapBlobKeyLastWins(true);
+ assertDuplicateMapBlobKeyLastWins(
+ true,
+ DataTypes.STRING(),
+ BinaryString.fromString("a"),
+ BinaryString.fromString("b"),
+ BinaryString.fromString("a"),
+ (byte) 'a');
}
- private void assertDuplicateMapBlobKeyLastWins(boolean blobAsDescriptor)
throws IOException {
+ @Test
+ public void testDuplicateBinaryMapBlobKeyLastWinsInline() throws
IOException {
+ assertDuplicateMapBlobKeyLastWins(
+ false,
+ DataTypes.BINARY(1),
+ new byte[] {1},
+ new byte[] {2},
+ new byte[] {1},
+ (byte) 1);
+ }
+
+ @Test
+ public void testDuplicateBinaryMapBlobKeyLastWinsAsDescriptor() throws
IOException {
+ assertDuplicateMapBlobKeyLastWins(
+ true,
+ DataTypes.BINARY(1),
+ new byte[] {1},
+ new byte[] {2},
+ new byte[] {1},
+ (byte) 1);
+ }
+
+ @Test
+ public void testBinaryMapBlobCopyPreservesLookup() throws IOException {
+ RowType rowType = writeBinaryMapBlobWithNullValue();
+ InternalRow row = readBinaryMapBlobRow(rowType);
+
+ InternalMap copiedMap =
+ new InternalMapSerializer(DataTypes.BYTES(),
DataTypes.BLOB()).copy(row.getMap(0));
+ assertBinaryKeyLookup(copiedMap);
+
+ InternalRow copiedRow = new InternalRowSerializer(rowType).copy(row);
+ assertBinaryKeyLookup(copiedRow.getMap(0));
+ }
+
+ @Test
+ public void testBinaryMapBlobJavaSerialization() throws Exception {
+ RowType rowType = writeBinaryMapBlobWithNullValue();
+ GenericMap map = (GenericMap) readBinaryMapBlobRow(rowType).getMap(0);
+
+ GenericMap restored = InstantiationUtil.clone(map);
+
+ assertBinaryKeyLookup(restored);
+ }
+
+ @Test
+ public void testBinaryMapBlobEqualityAcrossBackingRepresentations() throws
IOException {
+ RowType rowType = writeBinaryMapBlobWithNullValue();
+ InternalRow readRow = readBinaryMapBlobRow(rowType);
+ GenericMap readMap = (GenericMap) readRow.getMap(0);
+
Map<Object, Object> entries = new LinkedHashMap<>();
- entries.put(BinaryString.fromString("a"), new
BlobData("first".getBytes()));
- entries.put(BinaryString.fromString("b"), new
BlobData("second".getBytes()));
+ byte[] key = new byte[] {1};
+ entries.put(key, null);
+ GenericMap ordinaryMap = new GenericMap(entries);
+ GenericRow ordinaryRow = GenericRow.of(ordinaryMap);
+
+ assertThat(ordinaryMap.contains(key)).isTrue();
+ assertThat(ordinaryMap.contains(new byte[] {1})).isFalse();
+ assertThat(readMap).isEqualTo(ordinaryMap);
+ assertThat(ordinaryMap).isEqualTo(readMap);
+ assertThat(readMap.hashCode()).isEqualTo(ordinaryMap.hashCode());
+ assertThat(readRow).isEqualTo(ordinaryRow);
+ assertThat(ordinaryRow).isEqualTo(readRow);
+ assertThat(readRow.hashCode()).isEqualTo(ordinaryRow.hashCode());
+ }
+
+ private RowType writeBinaryMapBlobWithNullValue() throws IOException {
+ RowType rowType = RowType.of(DataTypes.MAP(DataTypes.BYTES(),
DataTypes.BLOB()));
+ Map<Object, Object> entries = new LinkedHashMap<>();
+ entries.put(new byte[] {1}, null);
+ try (PositionOutputStream out = fileIO.newOutputStream(file, false)) {
+ FormatWriter writer =
+ new BlobFileFormat(false,
BlobFormatWriter.DEFAULT_COPY_BUFFER_SIZE)
+ .createWriterFactory(rowType)
+ .create(out, null);
+ writer.addElement(GenericRow.of(new GenericMap(entries)));
+ writer.close();
+ }
+ return rowType;
+ }
+
+ private InternalRow readBinaryMapBlobRow(RowType rowType) throws
IOException {
+ FormatReaderFactory readerFactory =
+ new BlobFileFormat(false,
BlobFormatWriter.DEFAULT_COPY_BUFFER_SIZE)
+ .createReaderFactory(null, rowType, null);
+ FormatReaderContext context =
+ new FormatReaderContext(fileIO, file,
fileIO.getFileSize(file));
+ try (FileRecordReader<InternalRow> reader =
readerFactory.createReader(context)) {
+ return reader.readBatch().next();
+ }
+ }
+
+ private void assertBinaryKeyLookup(InternalMap map) {
+ assertThat(map).isInstanceOf(GenericMap.class);
+ GenericMap genericMap = (GenericMap) map;
+ assertThat(genericMap.contains(new byte[] {1})).isTrue();
+ assertThat(genericMap.get(new byte[] {1})).isNull();
+ assertThat(genericMap.keyArray().getBinary(0)).isEqualTo(new byte[]
{1});
+ }
+
+ private void assertDuplicateMapBlobKeyLastWins(
+ boolean blobAsDescriptor,
+ DataType keyType,
+ Object firstKey,
+ Object secondKey,
+ Object lookupKey,
+ byte duplicateKeyByte)
+ throws IOException {
+ Map<Object, Object> entries = new LinkedHashMap<>();
+ entries.put(firstKey, new BlobData("first".getBytes()));
+ entries.put(secondKey, new BlobData("second".getBytes()));
BlobFileFormat format =
new BlobFileFormat(blobAsDescriptor,
BlobFormatWriter.DEFAULT_COPY_BUFFER_SIZE);
- RowType rowType = RowType.of(DataTypes.MAP(DataTypes.STRING(),
DataTypes.BLOB()));
+ RowType rowType = RowType.of(DataTypes.MAP(keyType, DataTypes.BLOB()));
try (PositionOutputStream out = fileIO.newOutputStream(file, false)) {
FormatWriter writer =
format.createWriterFactory(rowType).create(out, null);
writer.addElement(GenericRow.of(new GenericMap(entries)));
@@ -469,7 +593,7 @@ public class BlobFileFormatTest {
java.nio.file.Path localFile = Paths.get(file.toUri());
byte[] bytes = Files.readAllBytes(localFile);
- bytes[payloadPosition + 10] = 'a';
+ bytes[payloadPosition + 10] = duplicateKeyByte;
Files.write(localFile, bytes);
FormatReaderFactory readerFactory = format.createReaderFactory(null,
rowType, null);
@@ -483,8 +607,11 @@ public class BlobFileFormatTest {
assertThat(rows).hasSize(1);
GenericMap result = (GenericMap) rows.get(0).getMap(0);
assertThat(result.size()).isOne();
- assertMapBlob(
- result.get(BinaryString.fromString("a")), blobAsDescriptor,
"second".getBytes());
+ assertThat(result.contains(lookupKey)).isTrue();
+ assertMapBlob(result.get(lookupKey), blobAsDescriptor,
"second".getBytes());
+ if (lookupKey instanceof byte[]) {
+ assertThat(result.keyArray().getBinary(0)).isEqualTo(lookupKey);
+ }
}
@Test
@@ -500,6 +627,9 @@ public class BlobFileFormatTest {
DataTypes.DECIMAL(20, 2),
DataTypes.DATE(),
DataTypes.TIME(3),
+ DataTypes.BINARY(4),
+ DataTypes.VARBINARY(8),
+ DataTypes.BYTES(),
DataTypes.CHAR(10),
DataTypes.VARCHAR(10)
};
@@ -514,6 +644,9 @@ public class BlobFileFormatTest {
Decimal.fromBigDecimal(new
BigDecimal("123456789012345678.90"), 20, 2),
-1,
45_296_789,
+ new byte[] {0, (byte) 0xff},
+ new byte[0],
+ new byte[] {1, 2, 3},
BinaryString.fromString("char"),
BinaryString.fromString("varchar")
};
@@ -538,6 +671,9 @@ public class BlobFileFormatTest {
},
{(byte) 0xff, (byte) 0xff, (byte) 0xff, (byte) 0xff},
{(byte) 0x95, 0x2c, (byte) 0xb3, 0x02},
+ {0, (byte) 0xff},
+ {},
+ {1, 2, 3},
"char".getBytes(),
"varchar".getBytes()
};
@@ -584,6 +720,9 @@ public class BlobFileFormatTest {
GenericMap result = (GenericMap) rows.get(0).getMap(0);
assertThat(result.contains(keys[i])).isTrue();
assertThat(((Blob)
result.get(keys[i])).toData()).isEqualTo("value".getBytes());
+ if (keys[i] instanceof byte[]) {
+ assertThat(result.keyArray().getBinary(0)).isEqualTo(keys[i]);
+ }
}
}
diff --git a/paimon-python/pypaimon/common/map_blob_key_serializer.py
b/paimon-python/pypaimon/common/map_blob_key_serializer.py
index 58eaa02dd4..53ccd4f5d1 100644
--- a/paimon-python/pypaimon/common/map_blob_key_serializer.py
+++ b/paimon-python/pypaimon/common/map_blob_key_serializer.py
@@ -88,6 +88,23 @@ class BooleanMapBlobKeySerializer(MapBlobKeySerializer):
raise ValueError("Invalid MAP<X, BLOB> boolean key.")
+class BinaryMapBlobKeySerializer(MapBlobKeySerializer):
+
+ def __init__(self, type_name: str):
+ self._type_name = type_name
+ self.fixed_length = -1
+
+ def serialize(self, key) -> bytes:
+ if not isinstance(key, bytes):
+ raise ValueError(
+ f"MAP<X, BLOB> {self._type_name} key must be bytes."
+ )
+ return key
+
+ def deserialize(self, data: bytes):
+ return data
+
+
class DecimalMapBlobKeySerializer(MapBlobKeySerializer):
def __init__(self, type_name: str, precision: int, scale: int):
@@ -211,6 +228,12 @@ def create_map_blob_key_serializer(data_type: DataType) ->
MapBlobKeySerializer:
return DateMapBlobKeySerializer()
if type_name == 'TIME' or type_name.startswith('TIME('):
return TimeMapBlobKeySerializer(type_name)
+ if (
+ type_name == 'BYTES'
+ or type_name.startswith('BINARY')
+ or type_name.startswith('VARBINARY')
+ ):
+ return BinaryMapBlobKeySerializer(type_name)
if type_name == 'STRING' or type_name.startswith('CHAR') or
type_name.startswith('VARCHAR'):
return MapBlobKeySerializer(type_name)
raise ValueError(f"Unsupported key type for MAP<X, BLOB>: {data_type}")
diff --git a/paimon-python/pypaimon/tests/blob_test.py
b/paimon-python/pypaimon/tests/blob_test.py
index 70100d2d1a..f661faebf2 100644
--- a/paimon-python/pypaimon/tests/blob_test.py
+++ b/paimon-python/pypaimon/tests/blob_test.py
@@ -3060,6 +3060,9 @@ class BlobEndToEndTest(unittest.TestCase):
),
(AtomicType("DATE"), datetime.date(1969, 12, 31)),
(AtomicType("TIME(3)"), datetime.time(12, 34, 56, 789000)),
+ (AtomicType("BINARY(4)"), bytes([0, 255])),
+ (AtomicType("VARBINARY(8)"), b""),
+ (AtomicType("BYTES"), b"bytes"),
(AtomicType("STRING"), "string"),
(AtomicType("CHAR(3)"), "abc"),
(AtomicType("VARCHAR(10)"), "varchar"),
@@ -3075,6 +3078,9 @@ class BlobEndToEndTest(unittest.TestCase):
b"\x00\xab\x54\xa9\x8c\xeb\x1f\x0a\xd2",
b"\xff\xff\xff\xff",
b"\x95\x2c\xb3\x02",
+ b"\x00\xff",
+ b"",
+ b"bytes",
b"string",
b"abc",
b"varchar",
@@ -3175,6 +3181,19 @@ class BlobEndToEndTest(unittest.TestCase):
RowKind.INSERT,
))
+ invalid_binary_key_writer = BlobFormatWriter(io.BytesIO())
+ binary_key_field = DataField(
+ 0,
+ "blob_map",
+ MapType(True, AtomicType("BINARY(4)"), AtomicType("BLOB")),
+ )
+ with self.assertRaisesRegex(ValueError, "key must be bytes"):
+ invalid_binary_key_writer.add_element(GenericRow(
+ [{"not-bytes": BlobData(b"value")}],
+ [binary_key_field],
+ RowKind.INSERT,
+ ))
+
invalid_time_key_writer = BlobFormatWriter(io.BytesIO())
time_key_field = DataField(
0,
diff --git a/paimon-python/pypaimon/tests/e2e/java_py_read_write_test.py
b/paimon-python/pypaimon/tests/e2e/java_py_read_write_test.py
index ce7b9d6038..fbf619a6f7 100644
--- a/paimon-python/pypaimon/tests/e2e/java_py_read_write_test.py
+++ b/paimon-python/pypaimon/tests/e2e/java_py_read_write_test.py
@@ -28,7 +28,7 @@ from parameterized import parameterized
from pypaimon.catalog.catalog_factory import CatalogFactory
from pypaimon.data.generic_variant import GenericVariant
from pypaimon.globalindex.data_evolution_global_index_scanner import
DataEvolutionGlobalIndexScanner
-from pypaimon.schema.data_types import VectorType
+from pypaimon.schema.data_types import PyarrowFieldParser, VectorType
from pypaimon.schema.schema import Schema
from pypaimon.read.read_builder import ReadBuilder
@@ -1567,6 +1567,12 @@ class JavaPyReadWriteTest(unittest.TestCase):
'time_payloads': {
datetime.time(12, 34, 56, 789000): b'java-time',
},
+ 'binary_payloads': {
+ bytes([0, 255, 1, 2]): b'java-binary',
+ },
+ 'varbinary_payloads': {
+ b'': b'java-varbinary',
+ },
}
for name, expected in expected_additional_payloads.items():
self.assertEqual(
@@ -1584,6 +1590,8 @@ class JavaPyReadWriteTest(unittest.TestCase):
pa.decimal128(20, 2), pa.large_binary())
date_map_blob_type = pa.map_(pa.date32(), pa.large_binary())
time_map_blob_type = pa.map_(pa.time32('ms'), pa.large_binary())
+ binary_schema_type = pa.map_(pa.binary(4), pa.large_binary())
+ varbinary_schema_type = pa.map_(pa.binary(), pa.large_binary())
pa_schema = pa.schema([
('id', pa.int32()),
('payloads', map_blob_type),
@@ -1592,6 +1600,8 @@ class JavaPyReadWriteTest(unittest.TestCase):
('high_decimal_payloads', high_decimal_map_blob_type),
('date_payloads', date_map_blob_type),
('time_payloads', time_map_blob_type),
+ ('binary_payloads', binary_schema_type),
+ ('varbinary_payloads', varbinary_schema_type),
])
schema = Schema.from_pyarrow_schema(
pa_schema,
@@ -1601,6 +1611,9 @@ class JavaPyReadWriteTest(unittest.TestCase):
'bucket': '-1',
},
)
+ pa_schema = PyarrowFieldParser.from_paimon_schema(schema.fields)
+ binary_map_blob_type = pa_schema.field('binary_payloads').type
+ varbinary_map_blob_type = pa_schema.field('varbinary_payloads').type
table_name = 'default.map_blob_python_test'
self.catalog.drop_table(table_name, True)
self.catalog.create_table(table_name, schema, False)
@@ -1647,6 +1660,22 @@ class JavaPyReadWriteTest(unittest.TestCase):
)], None, None, None],
type=time_map_blob_type,
),
+ 'binary_payloads': pa.array(
+ [
+ [
+ (bytes([0, 255, 1, 2]), b'python-binary-first'),
+ (bytes([0, 255, 1, 2]), b'python-binary'),
+ ],
+ None,
+ None,
+ None,
+ ],
+ type=binary_map_blob_type,
+ ),
+ 'varbinary_payloads': pa.array(
+ [[(b'', b'python-varbinary')], None, None, None],
+ type=varbinary_map_blob_type,
+ ),
}, schema=pa_schema)
write_builder = table.new_batch_write_builder()
table_write = write_builder.new_write()
@@ -1685,6 +1714,12 @@ class JavaPyReadWriteTest(unittest.TestCase):
'time_payloads': {
datetime.time(12, 34, 56, 789000): b'python-time',
},
+ 'binary_payloads': {
+ bytes([0, 255, 1, 2]): b'python-binary',
+ },
+ 'varbinary_payloads': {
+ b'': b'python-varbinary',
+ },
}
for name, expected in expected_additional_payloads.items():
self.assertEqual(