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 905e0eb9e7 [format] Honor JSON null map key options (#10031)
905e0eb9e7 is described below
commit 905e0eb9e72ca7a3754d24ae9b3d51a3cca46a23
Author: Akash Reddy Jammula <[email protected]>
AuthorDate: Sun Sep 20 20:26:35 2026 -0700
[format] Honor JSON null map key options (#10031)
---
.../apache/paimon/format/json/JsonFileReader.java | 31 +---
.../paimon/format/json/JsonFormatWriter.java | 41 ++++-
.../paimon/format/json/JsonMapKeyConverter.java | 60 +++++++
.../paimon/format/json/JsonFileFormatTest.java | 183 +++++++++++++++++----
4 files changed, 252 insertions(+), 63 deletions(-)
diff --git
a/paimon-format/src/main/java/org/apache/paimon/format/json/JsonFileReader.java
b/paimon-format/src/main/java/org/apache/paimon/format/json/JsonFileReader.java
index 5004396871..e776170816 100644
---
a/paimon-format/src/main/java/org/apache/paimon/format/json/JsonFileReader.java
+++
b/paimon-format/src/main/java/org/apache/paimon/format/json/JsonFileReader.java
@@ -18,9 +18,6 @@
package org.apache.paimon.format.json;
-import org.apache.paimon.casting.CastExecutor;
-import org.apache.paimon.casting.CastExecutors;
-import org.apache.paimon.data.BinaryString;
import org.apache.paimon.data.BinaryVector;
import org.apache.paimon.data.GenericArray;
import org.apache.paimon.data.GenericMap;
@@ -215,33 +212,7 @@ public class JsonFileReader extends AbstractTextFileReader
{
private Object convertPrimitiveStringToType(
String str, DataType dataType, JsonOptions options) {
try {
- switch (dataType.getTypeRoot()) {
- case TINYINT:
- return Byte.parseByte(str);
- case SMALLINT:
- return Short.parseShort(str);
- case INTEGER:
- return Integer.parseInt(str);
- case BIGINT:
- return Long.parseLong(str);
- case FLOAT:
- return Float.parseFloat(str);
- case DOUBLE:
- return Double.parseDouble(str);
- case BOOLEAN:
- return Boolean.parseBoolean(str);
- case CHAR:
- case VARCHAR:
- return BinaryString.fromString(str);
- default:
- CastExecutor cast =
CastExecutors.resolve(DataTypes.STRING(), dataType);
- if (cast == null) {
- // resolve returns null when no rule matches the
target type.
- throw new UnsupportedOperationException(
- "Unsupported data type for JSON format: " +
dataType);
- }
- return cast.cast(BinaryString.fromString(str));
- }
+ return JsonMapKeyConverter.convert(str, dataType);
} catch (Exception e) {
return handleParseError(e);
}
diff --git
a/paimon-format/src/main/java/org/apache/paimon/format/json/JsonFormatWriter.java
b/paimon-format/src/main/java/org/apache/paimon/format/json/JsonFormatWriter.java
index 2aeb9a5c1b..32d28bc36d 100644
---
a/paimon-format/src/main/java/org/apache/paimon/format/json/JsonFormatWriter.java
+++
b/paimon-format/src/main/java/org/apache/paimon/format/json/JsonFormatWriter.java
@@ -53,6 +53,8 @@ public class JsonFormatWriter extends AbstractTextFileWriter {
private static final Base64.Encoder BASE64_ENCODER = Base64.getEncoder();
private final String lineDelimiter;
+ private final JsonOptions.MapNullKeyMode mapNullKeyMode;
+ private final String mapNullKeyLiteral;
public JsonFormatWriter(
PositionOutputStream outputStream,
@@ -62,6 +64,8 @@ public class JsonFormatWriter extends AbstractTextFileWriter {
throws IOException {
super(outputStream, rowType, compression);
this.lineDelimiter = options.getLineDelimiter();
+ this.mapNullKeyMode = options.getMapNullKeyMode();
+ this.mapNullKeyLiteral = options.getMapNullKeyLiteral();
}
@Override
@@ -156,12 +160,47 @@ public class JsonFormatWriter extends
AbstractTextFileWriter {
for (int i = 0; i < size; i++) {
Object key = InternalRowUtils.get(keyArray, i, keyType);
+ String keyString;
+ if (key == null) {
+ switch (mapNullKeyMode) {
+ case DROP:
+ continue;
+ case LITERAL:
+ validateMapNullKeyLiteral(keyType);
+ keyString = mapNullKeyLiteral;
+ break;
+ case FAIL:
+ throw new IllegalArgumentException(
+ "JSON format does not support null map keys
when "
+ + "'json.map-null-key-mode' is
'FAIL'.");
+ default:
+ throw new IllegalStateException(
+ "Unknown map null key mode: " +
mapNullKeyMode);
+ }
+ } else {
+ keyString = convertToString(key, keyType);
+ }
Object value = InternalRowUtils.get(valueArray, i, valueType);
- result.put(convertToString(key, keyType), convertRowValue(value,
valueType));
+ result.put(keyString, convertRowValue(value, valueType));
}
return result;
}
+ private void validateMapNullKeyLiteral(DataType keyType) {
+ try {
+ Object converted = JsonMapKeyConverter.convert(mapNullKeyLiteral,
keyType);
+ if (converted == null) {
+ throw new IllegalArgumentException("Converted map key is
null.");
+ }
+ } catch (RuntimeException e) {
+ throw new IllegalArgumentException(
+ String.format(
+ "Cannot convert 'json.map-null-key-literal' value
'%s' to map key type %s.",
+ mapNullKeyLiteral, keyType),
+ e);
+ }
+ }
+
private String convertToString(Object value, DataType dataType) {
if (value == null) {
return null;
diff --git
a/paimon-format/src/main/java/org/apache/paimon/format/json/JsonMapKeyConverter.java
b/paimon-format/src/main/java/org/apache/paimon/format/json/JsonMapKeyConverter.java
new file mode 100644
index 0000000000..d638c5e684
--- /dev/null
+++
b/paimon-format/src/main/java/org/apache/paimon/format/json/JsonMapKeyConverter.java
@@ -0,0 +1,60 @@
+/*
+ * 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.format.json;
+
+import org.apache.paimon.casting.CastExecutor;
+import org.apache.paimon.casting.CastExecutors;
+import org.apache.paimon.data.BinaryString;
+import org.apache.paimon.types.DataType;
+import org.apache.paimon.types.DataTypes;
+
+/** Converts JSON object field names to Paimon map keys. */
+final class JsonMapKeyConverter {
+
+ private JsonMapKeyConverter() {}
+
+ static Object convert(String value, DataType dataType) {
+ switch (dataType.getTypeRoot()) {
+ case TINYINT:
+ return Byte.parseByte(value);
+ case SMALLINT:
+ return Short.parseShort(value);
+ case INTEGER:
+ return Integer.parseInt(value);
+ case BIGINT:
+ return Long.parseLong(value);
+ case FLOAT:
+ return Float.parseFloat(value);
+ case DOUBLE:
+ return Double.parseDouble(value);
+ case BOOLEAN:
+ return Boolean.parseBoolean(value);
+ case CHAR:
+ case VARCHAR:
+ return BinaryString.fromString(value);
+ default:
+ CastExecutor cast = CastExecutors.resolve(DataTypes.STRING(),
dataType);
+ if (cast == null) {
+ throw new UnsupportedOperationException(
+ "Unsupported data type for JSON format: " +
dataType);
+ }
+ return cast.cast(BinaryString.fromString(value));
+ }
+ }
+}
diff --git
a/paimon-format/src/test/java/org/apache/paimon/format/json/JsonFileFormatTest.java
b/paimon-format/src/test/java/org/apache/paimon/format/json/JsonFileFormatTest.java
index 3778f2a8f2..0629c3fe83 100644
---
a/paimon-format/src/test/java/org/apache/paimon/format/json/JsonFileFormatTest.java
+++
b/paimon-format/src/test/java/org/apache/paimon/format/json/JsonFileFormatTest.java
@@ -23,6 +23,7 @@ import org.apache.paimon.data.BinaryString;
import org.apache.paimon.data.BinaryVector;
import org.apache.paimon.data.GenericMap;
import org.apache.paimon.data.GenericRow;
+import org.apache.paimon.data.InternalMap;
import org.apache.paimon.data.InternalRow;
import org.apache.paimon.data.serializer.InternalRowSerializer;
import org.apache.paimon.format.FileFormat;
@@ -44,6 +45,8 @@ import org.junit.jupiter.api.Test;
import java.io.IOException;
import java.nio.charset.StandardCharsets;
+import java.nio.file.Files;
+import java.nio.file.Paths;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.List;
@@ -329,25 +332,15 @@ public class JsonFileFormatTest extends
FormatReadWriteTest {
DataTypes.ROW(
DataTypes.INT().notNull(),
DataTypes.MAP(DataTypes.STRING(), DataTypes.INT()));
-
- // Test JSON_MAP_NULL_KEY_MODE = FAIL with actual data
Options options = new Options();
options.set(JsonOptions.JSON_MAP_NULL_KEY_MODE,
JsonOptions.MapNullKeyMode.FAIL);
-
- // Create test data with valid maps
List<InternalRow> testData =
- Arrays.asList(
- GenericRow.of(1, new GenericMap(createTestMap("key1",
1, "key2", 2))),
- GenericRow.of(2, new GenericMap(createTestMap("name",
100, "value", 200))));
-
- List<InternalRow> result = writeThenRead(options, rowType, testData,
"test_fail_mode");
+ Arrays.asList(GenericRow.of(1, new
GenericMap(createTestMap(null, 1, "key", 2))));
- // Verify results
- assertThat(result).hasSize(2);
- assertThat(result.get(0).getInt(0)).isEqualTo(1);
- assertThat(result.get(0).getMap(1).size()).isEqualTo(2);
- assertThat(result.get(1).getInt(0)).isEqualTo(2);
- assertThat(result.get(1).getMap(1).size()).isEqualTo(2);
+ assertThatThrownBy(() -> writeThenRead(options, rowType, testData,
"test_fail_mode"))
+ .isInstanceOf(IllegalArgumentException.class)
+ .hasMessageContaining("json.map-null-key-mode")
+ .hasMessageContaining("FAIL");
}
@Test
@@ -357,22 +350,18 @@ public class JsonFileFormatTest extends
FormatReadWriteTest {
DataTypes.INT().notNull(),
DataTypes.MAP(DataTypes.STRING(), DataTypes.INT()));
- // Test JSON_MAP_NULL_KEY_MODE = DROP with actual data
Options options = new Options();
options.set(JsonOptions.JSON_MAP_NULL_KEY_MODE,
JsonOptions.MapNullKeyMode.DROP);
-
- // Create test data
List<InternalRow> testData =
Arrays.asList(
- GenericRow.of(
- 1, new GenericMap(createTestMap("key1", 1,
"key2", 2, "key3", 3))));
+ GenericRow.of(1, new GenericMap(createTestMap(null, 1,
"remaining", 2))));
List<InternalRow> result = writeThenRead(options, rowType, testData,
"test_drop_mode");
- // Verify results
assertThat(result).hasSize(1);
- assertThat(result.get(0).getInt(0)).isEqualTo(1);
- assertThat(result.get(0).getMap(1).size()).isEqualTo(3);
+ InternalMap map = result.get(0).getMap(1);
+ assertThat(map.size()).isEqualTo(1);
+ assertThat(map.valueArray().getInt(findMapKey(map,
"remaining"))).isEqualTo(2);
}
@Test
@@ -384,13 +373,11 @@ public class JsonFileFormatTest extends
FormatReadWriteTest {
String[] literals = {"EMPTY", "MISSING", "UNDEFINED", "NULL_VALUE"};
- // Create test data once (reused for all literals)
List<InternalRow> testData =
Arrays.asList(
GenericRow.of(
1,
- new GenericMap(
- createTestMap("name", "Alice", "city",
"New York"))));
+ new GenericMap(createTestMap(null, "missing",
"name", "Alice"))));
for (String literal : literals) {
Options options = new Options();
@@ -400,13 +387,122 @@ public class JsonFileFormatTest extends
FormatReadWriteTest {
List<InternalRow> result =
writeThenRead(options, rowType, testData, "test_literal_"
+ literal);
- // Verify results
assertThat(result).hasSize(1);
- assertThat(result.get(0).getInt(0)).isEqualTo(1);
- assertThat(result.get(0).getMap(1).size()).isEqualTo(2);
+ InternalMap map = result.get(0).getMap(1);
+ assertThat(map.size()).isEqualTo(2);
+ assertThat(map.valueArray().getString(findMapKey(map,
literal)).toString())
+ .isEqualTo("missing");
+ assertThat(map.valueArray().getString(findMapKey(map,
"name")).toString())
+ .isEqualTo("Alice");
}
}
+ @Test
+ public void testMapNullKeyLiteralInNestedMap() throws IOException {
+ RowType rowType =
+ DataTypes.ROW(
+ DataTypes.MAP(
+ DataTypes.STRING(),
+ DataTypes.MAP(DataTypes.STRING(),
DataTypes.INT())));
+ Options options = new Options();
+ options.set(JsonOptions.JSON_MAP_NULL_KEY_MODE,
JsonOptions.MapNullKeyMode.LITERAL);
+ options.set(JsonOptions.JSON_MAP_NULL_KEY_LITERAL, "null-key");
+ GenericMap nestedMap = new GenericMap(createTestMap(null, 1, "key",
2));
+ List<InternalRow> testData =
+ Arrays.asList(GenericRow.of(new
GenericMap(createTestMap("outer", nestedMap))));
+
+ List<InternalRow> result = writeThenRead(options, rowType, testData,
"test_nested_literal");
+
+ InternalMap outerMap = result.get(0).getMap(0);
+ InternalMap innerMap =
outerMap.valueArray().getMap(findMapKey(outerMap, "outer"));
+ assertThat(innerMap.size()).isEqualTo(2);
+ assertThat(innerMap.valueArray().getInt(findMapKey(innerMap,
"null-key"))).isEqualTo(1);
+ assertThat(innerMap.valueArray().getInt(findMapKey(innerMap,
"key"))).isEqualTo(2);
+ }
+
+ @Test
+ public void testMapNullKeySerializedOutput() throws IOException {
+ RowType rowType = DataTypes.ROW(DataTypes.MAP(DataTypes.STRING(),
DataTypes.INT()));
+ InternalRow row = GenericRow.of(new GenericMap(createTestMap(null, 1,
"key", 2)));
+
+ Options dropOptions = new Options();
+ dropOptions.set(JsonOptions.JSON_MAP_NULL_KEY_MODE,
JsonOptions.MapNullKeyMode.DROP);
+ assertThat(writeToJson(dropOptions, rowType, row, "test_drop_output"))
+ .isEqualTo("{\"f0\":{\"key\":\"2\"}}\n");
+
+ Options defaultLiteralOptions = new Options();
+ defaultLiteralOptions.set(
+ JsonOptions.JSON_MAP_NULL_KEY_MODE,
JsonOptions.MapNullKeyMode.LITERAL);
+ assertThat(writeToJson(defaultLiteralOptions, rowType, row,
"test_default_literal_output"))
+ .isEqualTo("{\"f0\":{\"null\":\"1\",\"key\":\"2\"}}\n");
+
+ Options customLiteralOptions = new Options();
+ customLiteralOptions.set(
+ JsonOptions.JSON_MAP_NULL_KEY_MODE,
JsonOptions.MapNullKeyMode.LITERAL);
+ customLiteralOptions.set(JsonOptions.JSON_MAP_NULL_KEY_LITERAL,
"missing");
+ assertThat(writeToJson(customLiteralOptions, rowType, row,
"test_custom_literal_output"))
+ .isEqualTo("{\"f0\":{\"missing\":\"1\",\"key\":\"2\"}}\n");
+ }
+
+ @Test
+ public void testMapNullKeyLiteralCollisionUsesLastEntry() throws
IOException {
+ RowType rowType = DataTypes.ROW(DataTypes.MAP(DataTypes.STRING(),
DataTypes.INT()));
+ Options options = new Options();
+ options.set(JsonOptions.JSON_MAP_NULL_KEY_MODE,
JsonOptions.MapNullKeyMode.LITERAL);
+ options.set(JsonOptions.JSON_MAP_NULL_KEY_LITERAL, "key");
+
+ InternalRow realKeyLast = GenericRow.of(new
GenericMap(createTestMap(null, 1, "key", 2)));
+ assertThat(writeToJson(options, rowType, realKeyLast,
"test_real_key_last"))
+ .isEqualTo("{\"f0\":{\"key\":\"2\"}}\n");
+
+ InternalRow nullKeyLast = GenericRow.of(new
GenericMap(createTestMap("key", 2, null, 1)));
+ assertThat(writeToJson(options, rowType, nullKeyLast,
"test_null_key_last"))
+ .isEqualTo("{\"f0\":{\"key\":\"1\"}}\n");
+ }
+
+ @Test
+ public void testMapNullKeyLiteralMustMatchKeyType() {
+ RowType rowType = DataTypes.ROW(DataTypes.MAP(DataTypes.INT(),
DataTypes.INT()));
+ Options options = new Options();
+ options.set(JsonOptions.JSON_MAP_NULL_KEY_MODE,
JsonOptions.MapNullKeyMode.LITERAL);
+ java.util.Map<Integer, Integer> map = new java.util.LinkedHashMap<>();
+ map.put(null, 1);
+ List<InternalRow> testData = Arrays.asList(GenericRow.of(new
GenericMap(map)));
+
+ assertThatThrownBy(
+ () ->
+ writeThenRead(
+ options, rowType, testData,
"test_invalid_literal_type"))
+ .isInstanceOf(IllegalArgumentException.class)
+ .hasMessageContaining("json.map-null-key-literal")
+ .hasMessageContaining("INT");
+ }
+
+ @Test
+ public void testMapNullKeyLiteralWithNonStringKeyType() throws IOException
{
+ RowType rowType = DataTypes.ROW(DataTypes.MAP(DataTypes.INT(),
DataTypes.INT()));
+ Options options = new Options();
+ options.set(JsonOptions.JSON_MAP_NULL_KEY_MODE,
JsonOptions.MapNullKeyMode.LITERAL);
+ options.set(JsonOptions.JSON_MAP_NULL_KEY_LITERAL, "0");
+ java.util.Map<Integer, Integer> map = new java.util.LinkedHashMap<>();
+ map.put(null, 1);
+ map.put(2, 2);
+
+ List<InternalRow> result =
+ writeThenRead(
+ options,
+ rowType,
+ Arrays.asList(GenericRow.of(new GenericMap(map))),
+ "test_integer_literal");
+
+ InternalMap resultMap = result.get(0).getMap(0);
+ java.util.Map<Integer, Integer> actual = new java.util.HashMap<>();
+ for (int i = 0; i < resultMap.size(); i++) {
+ actual.put(resultMap.keyArray().getInt(i),
resultMap.valueArray().getInt(i));
+ }
+ assertThat(actual).containsEntry(0, 1).containsEntry(2, 2).hasSize(2);
+ }
+
@Test
public void testWithCustomLineDelimiters() throws IOException {
RowType rowType =
@@ -479,19 +575,42 @@ public class JsonFileFormatTest extends
FormatReadWriteTest {
throw new IllegalArgumentException("Key-value pairs must be even
number of arguments");
}
- java.util.Map<BinaryString, Object> map = new java.util.HashMap<>();
+ java.util.Map<BinaryString, Object> map = new
java.util.LinkedHashMap<>();
for (int i = 0; i < keyValuePairs.length; i += 2) {
String key = (String) keyValuePairs[i];
Object value = keyValuePairs[i + 1];
if (value instanceof String) {
- map.put(BinaryString.fromString(key),
BinaryString.fromString((String) value));
+ map.put(
+ key == null ? null : BinaryString.fromString(key),
+ BinaryString.fromString((String) value));
} else {
- map.put(BinaryString.fromString(key), value);
+ map.put(key == null ? null : BinaryString.fromString(key),
value);
}
}
return map;
}
+ private int findMapKey(InternalMap map, String key) {
+ for (int i = 0; i < map.size(); i++) {
+ if (map.keyArray().getString(i).toString().equals(key)) {
+ return i;
+ }
+ }
+ throw new AssertionError("Map key not found: " + key);
+ }
+
+ private String writeToJson(Options options, RowType rowType, InternalRow
row, String testPrefix)
+ throws IOException {
+ FileFormat format =
+ new JsonFileFormat(new
FileFormatFactory.FormatContext(options, 1024, 1024));
+ Path testFile = new Path(parent, testPrefix + "_" + UUID.randomUUID()
+ ".json");
+ try (PositionOutputStream out = fileIO.newOutputStream(testFile,
false);
+ FormatWriter writer =
format.createWriterFactory(rowType).create(out, "none")) {
+ writer.addElement(row);
+ }
+ return new String(Files.readAllBytes(Paths.get(testFile.toUri())),
StandardCharsets.UTF_8);
+ }
+
private List<InternalRow> writeThenRead(
Options options, RowType rowType, List<InternalRow> testData,
String testPrefix)
throws IOException {