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

Reply via email to