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 a3aed62cc130c8..c3ab8d29be811b 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 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 static String getString(byte[] value, int pos) { 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 static String getMetadataKey(byte[] metadata, int id) { 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 cec12149ceb5f8..3ce271896a2f82 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 @@ void testParseJsonObject() throws IOException { 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 77235e968eebce..4468fcbe8d8e17 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.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 @@ void testToJsonRejectsNonFiniteFloat(final float nonFinite) { .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]))