This is an automated email from the ASF dual-hosted git repository.
snuyanzin pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/flink.git
The following commit(s) were added to refs/heads/master by this push:
new 33a198ae41c [FLINK-40305][core] Decode `VARIANT` strings and object
keys as UTF-8
33a198ae41c is described below
commit 33a198ae41ca1f11a8a8f313a5295551d1d591df
Author: Ramin Gharib <[email protected]>
AuthorDate: Wed Aug 5 12:27:54 2026 +0200
[FLINK-40305][core] Decode `VARIANT` strings and object keys as UTF-8
`BinaryVariantUtil` decoded string values and object field names with `new
String(byte[], int, int)`, which uses the JVM default charset, while
`BinaryVariantInternalBuilder` writes both as UTF-8. The two only agree on Java
18+, where JEP 400 made UTF-8 the default charset. On Java 11 and 17 a
non-UTF-8 platform charset corrupts any non-ASCII text.
Corrupted field names are the worse half of this. `getField(name)` silently
returns null, and `getFieldNames()` and `toJson()` return mangled keys.
Both call sites now pass `StandardCharsets.UTF_8` explicitly, matching
Spark's `VariantUtil`.
---
.../flink/types/variant/BinaryVariantUtil.java | 6 ++-
.../variant/BinaryVariantInternalBuilderTest.java | 12 ++++++
.../flink/types/variant/BinaryVariantTest.java | 49 ++++++++++++++++++++++
3 files changed, 65 insertions(+), 2 deletions(-)
diff --git
a/flink-core/src/main/java/org/apache/flink/types/variant/BinaryVariantUtil.java
b/flink-core/src/main/java/org/apache/flink/types/variant/BinaryVariantUtil.java
index a3aed62cc13..c3ab8d29be8 100644
---
a/flink-core/src/main/java/org/apache/flink/types/variant/BinaryVariantUtil.java
+++
b/flink-core/src/main/java/org/apache/flink/types/variant/BinaryVariantUtil.java
@@ -23,6 +23,7 @@ import org.apache.flink.types.variant.Variant.Type;
import java.math.BigDecimal;
import java.math.BigInteger;
+import java.nio.charset.StandardCharsets;
import java.time.format.DateTimeFormatter;
import java.time.format.DateTimeFormatterBuilder;
import java.util.Arrays;
@@ -528,7 +529,7 @@ public class BinaryVariantUtil {
length = readUnsigned(value, pos + 1, U32_SIZE);
}
checkIndex(start + length - 1, value.length);
- return new String(value, start, length);
+ return new String(value, start, length, StandardCharsets.UTF_8);
}
throw unexpectedType(Type.STRING);
}
@@ -625,6 +626,7 @@ public class BinaryVariantUtil {
throw malformedVariant();
}
checkIndex(stringStart + nextOffset - 1, metadata.length);
- return new String(metadata, stringStart + offset, nextOffset - offset);
+ return new String(
+ metadata, stringStart + offset, nextOffset - offset,
StandardCharsets.UTF_8);
}
}
diff --git
a/flink-core/src/test/java/org/apache/flink/types/variant/BinaryVariantInternalBuilderTest.java
b/flink-core/src/test/java/org/apache/flink/types/variant/BinaryVariantInternalBuilderTest.java
index cec12149ceb..3ce271896a2 100644
---
a/flink-core/src/test/java/org/apache/flink/types/variant/BinaryVariantInternalBuilderTest.java
+++
b/flink-core/src/test/java/org/apache/flink/types/variant/BinaryVariantInternalBuilderTest.java
@@ -122,6 +122,18 @@ class BinaryVariantInternalBuilderTest {
assertThat(variant.getField("k2").getDecimal()).isEqualTo(BigDecimal.valueOf(1.5));
}
+ @Test
+ void testParseJsonWithNonAsciiStringsAndKeys() throws IOException {
+ String json = "{\"schlüssel\":\"Grüße, 世界 🚀\",\"キー\":[\"äöü\"]}";
+
+ BinaryVariant variant = BinaryVariantInternalBuilder.parseJson(json,
false);
+
+
assertThat(variant.getFieldNames()).containsExactlyInAnyOrder("schlüssel",
"キー");
+
assertThat(variant.getField("schlüssel").getString()).isEqualTo("Grüße, 世界 🚀");
+
assertThat(variant.getField("キー").getElement(0).getString()).isEqualTo("äöü");
+ assertThat(variant.toJson()).isEqualTo(json);
+ }
+
@ParameterizedTest
@ValueSource(strings = {"NaN", "Infinity", "-Infinity", "1e400", "-1e400"})
void testParseJsonRejectsNonFiniteNumbers(final String nonFiniteNumber) {
diff --git
a/flink-core/src/test/java/org/apache/flink/types/variant/BinaryVariantTest.java
b/flink-core/src/test/java/org/apache/flink/types/variant/BinaryVariantTest.java
index 77235e968ee..4468fcbe8d8 100644
---
a/flink-core/src/test/java/org/apache/flink/types/variant/BinaryVariantTest.java
+++
b/flink-core/src/test/java/org/apache/flink/types/variant/BinaryVariantTest.java
@@ -24,10 +24,12 @@ import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.ValueSource;
import java.math.BigDecimal;
+import java.nio.charset.StandardCharsets;
import java.time.Instant;
import java.time.LocalDate;
import java.time.LocalDateTime;
import java.time.temporal.ChronoUnit;
+import java.util.Collections;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
@@ -254,6 +256,53 @@ class BinaryVariantTest {
.hasMessageContaining("cannot be serialized to JSON");
}
+ @Test
+ void testNonAsciiStringsAndFieldNames() {
+ // Multi-byte code points make the UTF-8 byte length differ from the
character count, so a
+ // charset mismatch between writing and reading mangles the text
instead of preserving it.
+ final String nestedKey = "キー";
+ final String shortValue = "Grüße, 世界 🚀";
+ final String longValue = String.join("", Collections.nCopies(20,
"äö🚀"));
+
+ assertThat(longValue.getBytes(StandardCharsets.UTF_8).length)
+ .as("long string must not fit into the short string encoding")
+ .isGreaterThan(BinaryVariantUtil.MAX_SHORT_STR_SIZE);
+
+ final BinaryVariant variant =
+ (BinaryVariant)
+ builder.object()
+ .add("schlüssel", builder.of(shortValue))
+ .add(
+ nestedKey,
+ builder.object()
+ .add("schlüssel",
builder.of(longValue))
+ .build())
+ .build();
+
+ // Reading through the raw binaries is what happens once a variant has
been serialized, and
+ // it is the only path that decodes the field names from the metadata.
+ final BinaryVariant decoded = new BinaryVariant(variant.getValue(),
variant.getMetadata());
+
+
assertThat(decoded.getFieldNames()).containsExactlyInAnyOrder("schlüssel",
nestedKey);
+
assertThat(decoded.getField("schlüssel").getString()).isEqualTo(shortValue);
+
assertThat(decoded.getField(nestedKey).getFieldNames()).containsExactly("schlüssel");
+
assertThat(decoded.getField(nestedKey).getField("schlüssel").getString())
+ .isEqualTo(longValue);
+ assertThat(decoded.toJson())
+ .isEqualTo(
+ "{\""
+ + "schlüssel"
+ + "\":\""
+ + shortValue
+ + "\",\""
+ + nestedKey
+ + "\":{\""
+ + "schlüssel"
+ + "\":\""
+ + longValue
+ + "\"}}");
+ }
+
@Test
void testVariantException() {
assertThatThrownBy(() -> new BinaryVariant(new byte[0], new byte[0]))