raminqaf commented on code in PR #29370:
URL: https://github.com/apache/flink/pull/29370#discussion_r4183138296


##########
flink-core/src/main/java/org/apache/flink/types/variant/BinaryVariantInternalBuilder.java:
##########
@@ -101,6 +101,9 @@ public class BinaryVariantInternalBuilder {
     // UTF-32 hint and could turn input that is not valid JSON into a value.
     private static final JsonFactory JSON_FACTORY =
             new 
JsonFactoryBuilder().disable(JsonFactory.Feature.CHARSET_DETECTION).build();
+    // The smallest unscaled magnitudes with more digits than a decimal4 or a 
decimal8 holds.
+    private static final long DECIMAL4_UNSCALED_LIMIT = 1_000_000_000L;
+    private static final long DECIMAL8_UNSCALED_LIMIT = 
1_000_000_000_000_000_000L;

Review Comment:
   The check needs the bound 10^p rather than p itself, so I kept the two 
constants but derived them from `MAX_DECIMAL4_PRECISION` and 
`MAX_DECIMAL8_PRECISION`. They are computed once when the class loads, so the 
per-record check stays two comparisons.



##########
flink-core/src/main/java/org/apache/flink/types/variant/BinaryVariantInternalBuilder.java:
##########
@@ -101,6 +101,9 @@ public class BinaryVariantInternalBuilder {
     // UTF-32 hint and could turn input that is not valid JSON into a value.
     private static final JsonFactory JSON_FACTORY =
             new 
JsonFactoryBuilder().disable(JsonFactory.Feature.CHARSET_DETECTION).build();
+    // The smallest unscaled magnitudes with more digits than a decimal4 or a 
decimal8 holds.

Review Comment:
   Fixed. The comment and the Javadoc of `appendDecimal(long, int)` now say 
`DECIMAL4` and `DECIMAL8`.



##########
flink-core/src/main/java/org/apache/flink/types/variant/BinaryVariantInternalBuilder.java:
##########
@@ -185,18 +188,25 @@ public BinaryVariant build() {
     }
 
     public void appendString(String str) {
-        byte[] text = str.getBytes(StandardCharsets.UTF_8);
-        boolean longStr = text.length > MAX_SHORT_STR_SIZE;
-        checkCapacity((longStr ? 1 + U32_SIZE : 1) + text.length);
+        appendString(str.getBytes(StandardCharsets.UTF_8));
+    }
+
+    /**
+     * Appends a string given as UTF-8 bytes. The bytes are copied as they 
are, so the caller must
+     * make sure they are valid UTF-8, which the variant spec requires of 
every string.
+     */
+    public void appendString(byte[] utf8) {
+        boolean longStr = utf8.length > MAX_SHORT_STR_SIZE;
+        checkCapacity((longStr ? 1 + U32_SIZE : 1) + utf8.length);

Review Comment:
   Done, as `1 + (longStr ? U32_SIZE : 0) + length`.



##########
flink-core/src/test/java/org/apache/flink/types/variant/BinaryVariantInternalBuilderTest.java:
##########
@@ -201,4 +208,58 @@ void testAppendFloat() {
 
         assertThatCode(() -> 
floatList.forEach(builder::appendFloat)).doesNotThrowAnyException();
     }
+
+    @ParameterizedTest
+    @ValueSource(
+            strings = {
+                "",
+                "Grüße, 世界 🚀",
+                "A string longer than 63 bytes is stored with a 4-byte length 
header."
+            })
+    void testAppendStringStoresUtf8BytesAsTheyAre(final String str) {
+        final byte[] utf8 = str.getBytes(UTF_8);
+        final BinaryVariantInternalBuilder builder = new 
BinaryVariantInternalBuilder(false);
+        builder.appendString(utf8);
+        final BinaryVariant variant = builder.build();
+
+        final byte[] value = variant.getValue();
+        final int headerSize = utf8.length > MAX_SHORT_STR_SIZE ? 1 + U32_SIZE 
: 1;

Review Comment:
   As Arvid wrote, the suggestion drops the 1-byte header and moves the 
threshold by one byte. I checked, and the test fails with it. I used the same 
shape as in the builder instead: `1 + (length > MAX_SHORT_STR_SIZE ? U32_SIZE : 
0)`.



##########
flink-table/flink-table-runtime/src/test/java/org/apache/flink/table/runtime/functions/VariantCastUtilsTest.java:
##########
@@ -65,6 +79,49 @@ void testCastToVariantRejectsValuesOverTheSizeLimit() {
                 .hasMessageStartingWith("Cannot cast a string value of 
16777212 bytes to VARIANT.");
     }
 
+    @ParameterizedTest
+    @ValueSource(strings = {"", "hello", "Grüße, 世界 🚀"})
+    void testCastStringToVariantStoresItsUtf8Bytes(final String str) {
+        
assertThat(fromString(binaryString(str.getBytes(UTF_8)))).isEqualTo(BUILDER.of(str));
+    }
+
+    @ParameterizedTest
+    @MethodSource("invalidUtf8")
+    void testCastStringToVariantReplacesInvalidUtf8(final byte[] invalid) {
+        final Variant variant = fromString(binaryString(invalid));
+
+        assertThat(variant.getString()).contains(REPLACEMENT_CHARACTER);
+        assertThat(variant).isEqualTo(BUILDER.of(new String(invalid, UTF_8)));
+    }
+
+    private static Stream<byte[]> invalidUtf8() {
+        return Stream.of(
+                new byte[] {'a', (byte) 0xFF, 'b'},
+                new byte[] {'a', (byte) 0xC3},
+                new byte[] {(byte) 0xC0, (byte) 0xAF},
+                new byte[] {(byte) 0xED, (byte) 0xA0, (byte) 0x80});
+    }
+
+    @ParameterizedTest(name = "{0} as DECIMAL({1}, {2})")
+    @CsvSource({
+        "0, 1, 0",
+        "1.50, 10, 2",
+        "999999999, 9, 0",
+        "-0.999999999, 9, 9",
+        "1000000000, 10, 0",
+        "0.0000000001, 10, 10",
+        "-999999999999999999, 18, 0",

Review Comment:
   As Arvid wrote, `DecimalType` only allows a scale between 0 and the 
precision, so `fromDecimal` never sees a negative one. The builder still 
handles it by passing it to `appendDecimal(BigDecimal)`, which rescales it to 
0. `BinaryVariantInternalBuilderTest` covers scale -1.



##########
flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/functions/VariantCastUtils.java:
##########
@@ -593,14 +593,31 @@ public static Variant fromDouble(double value) {
     }
 
     public static Variant fromDecimal(DecimalData value) {
-        return BUILDER.of(value.toBigDecimal());
+        if (!value.isCompact()) {
+            return BUILDER.of(value.toBigDecimal());
+        }
+        final BinaryVariantInternalBuilder builder = new 
BinaryVariantInternalBuilder(false);
+        builder.appendDecimal(value.toUnscaledLong(), value.scale());
+        return builder.build();
     }
 
+    /**
+     * Stores the UTF-8 bytes of the string as they are. A {@link StringData} 
may hold invalid
+     * UTF-8, which the variant spec does not allow, so such a value is 
decoded first and every
+     * malformed sequence is stored as the U+FFFD replacement character.
+     */
     public static Variant fromString(StringData value) {
+        final byte[] utf8 = value.toBytes();

Review Comment:
   Yes. `BinaryStringData` is the only `StringData` implementation, and the 
runtime already casts to it in many places. `fromString` now reads a string 
that lies in the first heap segment in place, through a new 
`appendString(byte[], int, int)`, so it no longer copies the bytes out with 
`toBytes()`. Otherwise it still takes the copy, and a string that only exists 
as a Java object is encoded straight from it, as on master. 
`VariantCastUtilsTest` covers each layout.
   
   Against the previous revision, JMH shows the cast 6 to 17% faster, and 
allocation drops by a quarter from 1 KiB up. A 64 KiB string now allocates 197 
KB per call instead of 262 KB.



##########
flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/functions/VariantCastUtils.java:
##########
@@ -593,14 +594,43 @@ public static Variant fromDouble(double value) {
     }
 
     public static Variant fromDecimal(DecimalData value) {
-        return BUILDER.of(value.toBigDecimal());
+        if (!value.isCompact()) {
+            return BUILDER.of(value.toBigDecimal());
+        }
+        final BinaryVariantInternalBuilder builder = new 
BinaryVariantInternalBuilder(false);
+        builder.appendDecimal(value.toUnscaledLong(), value.scale());
+        return builder.build();
     }
 
+    /**
+     * Stores the UTF-8 bytes of the string as they are. A {@link StringData} 
may hold invalid
+     * UTF-8, which the variant spec does not allow, so such a value is 
decoded first and every
+     * malformed sequence is stored as the U+FFFD replacement character.
+     */
     public static Variant fromString(StringData value) {

Review Comment:
   Thanks, this was more than a missing test. A `StringData.fromString` value 
was up to 3.7 times slower than on master: `getSegments()` encodes it with 
`StringUtf8Utils#encodeUTF8`, which is much slower than `String#getBytes`, and 
then the check reads bytes that are valid by construction. `LazyBinaryFormat` 
does expose the form through `getBinarySection()` and `getJavaObject()`, so 
`appendString` now encodes such a string straight from the Java object, as 
master did. It is on par with master again. `VariantCastUtilsTest` has a 
`JAVA_OBJECT` layout, and a test with an unpaired surrogate that also checks 
that the string keeps no binary form.



##########
flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/functions/VariantCastUtils.java:
##########
@@ -593,14 +594,43 @@ public static Variant fromDouble(double value) {
     }
 
     public static Variant fromDecimal(DecimalData value) {
-        return BUILDER.of(value.toBigDecimal());
+        if (!value.isCompact()) {
+            return BUILDER.of(value.toBigDecimal());
+        }
+        final BinaryVariantInternalBuilder builder = new 
BinaryVariantInternalBuilder(false);
+        builder.appendDecimal(value.toUnscaledLong(), value.scale());
+        return builder.build();
     }
 
+    /**
+     * Stores the UTF-8 bytes of the string as they are. A {@link StringData} 
may hold invalid
+     * UTF-8, which the variant spec does not allow, so such a value is 
decoded first and every
+     * malformed sequence is stored as the U+FFFD replacement character.
+     */
     public static Variant fromString(StringData value) {
+        final BinaryStringData string = (BinaryStringData) value;
+        final MemorySegment[] segments = string.getSegments();
+        final int length = string.getSizeInBytes();
+        final byte[] utf8;
+        final int offset;
+        // Read a string that lies in one heap segment in place, which saves 
copying it out.
+        if (segments.length == 1 && !segments[0].isOffHeap()) {

Review Comment:
   Done, with your condition. The layout test also has a string in the first of 
two segments now.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to