AHeise commented on code in PR #29369:
URL: https://github.com/apache/flink/pull/29369#discussion_r4184726670


##########
flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/functions/casting/ToVariantCastRule.java:
##########
@@ -0,0 +1,155 @@
+/*
+ * 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.flink.table.planner.functions.casting;
+
+import org.apache.flink.table.runtime.functions.ToVariantConverter;
+import org.apache.flink.table.types.logical.DistinctType;
+import org.apache.flink.table.types.logical.LogicalType;
+import org.apache.flink.table.types.logical.LogicalTypeRoot;
+import org.apache.flink.table.types.logical.utils.LogicalTypeCasts;
+import org.apache.flink.table.types.logical.utils.LogicalTypeChecks;
+import org.apache.flink.types.variant.BinaryVariantUtil;
+import org.apache.flink.types.variant.Variant;
+
+import java.util.ArrayDeque;
+import java.util.Deque;
+
+import static 
org.apache.flink.table.planner.functions.casting.CastRuleUtils.methodCall;
+
+/**
+ * Cast rule from {@link LogicalTypeRoot#ARRAY}, {@link LogicalTypeRoot#MAP}, 
{@link
+ * LogicalTypeRoot#ROW} and {@link LogicalTypeRoot#STRUCTURED_TYPE} to {@link
+ * LogicalTypeRoot#VARIANT}.
+ *
+ * <p>The generated code calls a {@link ToVariantConverter} created once for 
the input type, which
+ * describes how a value is stored. A failure anywhere in the value fails the 
whole cast, so {@code
+ * TRY_CAST} returns {@code NULL} rather than a partial result. See {@link 
#canFail} for when a cast
+ * can fail.
+ */
+class ConstructedToVariantCastRule extends 
AbstractNullAwareCodeGeneratorCastRule<Object, Variant> {
+
+    static final ConstructedToVariantCastRule INSTANCE = new 
ConstructedToVariantCastRule();
+
+    /** The largest fixed-size value: a 16-byte decimal with its header and 
scale. */
+    private static final long MAX_FIXED_SIZE = 18;
+
+    /** A header byte and a 4-byte size or length, the largest header a value 
has. */
+    private static final long MAX_HEADER_SIZE = 1 + BinaryVariantUtil.U32_SIZE;
+
+    private ConstructedToVariantCastRule() {
+        
super(CastRulePredicate.builder().predicate(ConstructedToVariantCastRule::matches).build());
+    }
+
+    private static boolean matches(LogicalType input, LogicalType target) {
+        return target.is(LogicalTypeRoot.VARIANT)
+                && (input.is(LogicalTypeRoot.ARRAY)
+                        || input.is(LogicalTypeRoot.MAP)
+                        || input.is(LogicalTypeRoot.ROW)
+                        || input.is(LogicalTypeRoot.STRUCTURED_TYPE))
+                && LogicalTypeCasts.supportsExplicitCast(input, target);
+    }
+
+    /**
+     * Returns whether a value of the input type can fail the cast, so that 
{@code TRY_CAST} returns
+     * {@code NULL} for it instead of failing. A value can fail in two ways:
+     *
+     * <ul>
+     *   <li>A leaf fails like it does on its own, see {@link 
PrimitiveToVariantCastRule#canFail}.
+     *   <li>The whole value does not fit into the 16 MiB of a {@code 
VARIANT}. Any {@code ARRAY},
+     *       {@code MAP} or nested {@code VARIANT} can reach that size. A 
{@code ROW} can when the
+     *       declared sizes of its fields add up to more.
+     * </ul>
+     *
+     * <p>A {@code MAP} also fails on a {@code NULL} key, but every {@code 
MAP} can fail anyway.
+     */
+    @Override
+    public boolean canFail(LogicalType inputLogicalType, LogicalType 
targetLogicalType) {
+        // Without an ARRAY, MAP or VARIANT, the size of a value is bounded by 
the sum of its
+        // leaves and object headers.
+        long maxSize = 0;
+        final Deque<LogicalType> pending = new ArrayDeque<>();
+        pending.push(inputLogicalType);
+        while (!pending.isEmpty()) {
+            final LogicalType type = pending.pop();
+            if (type.is(LogicalTypeRoot.DISTINCT_TYPE)) {
+                pending.push(((DistinctType) type).getSourceType());
+                continue;
+            }
+            switch (type.getTypeRoot()) {
+                case ARRAY:
+                case MAP:
+                case VARIANT:
+                    return true;
+                case ROW:
+                case STRUCTURED_TYPE:
+                    
LogicalTypeChecks.getFieldTypes(type).forEach(pending::push);
+                    break;
+                default:
+                    if (PrimitiveToVariantCastRule.INSTANCE.canFail(type, 
targetLogicalType)) {
+                        return true;
+                    }
+            }
+            maxSize += maxOwnSize(type);
+        }
+        return maxSize > BinaryVariantUtil.SIZE_LIMIT;
+    }
+
+    private static long maxOwnSize(LogicalType type) {
+        switch (type.getTypeRoot()) {
+            case ROW:
+            case STRUCTURED_TYPE:
+                // Only the object header. Its fields are counted on their 
own. A 4-byte id per
+                // field, and a 4-byte offset per field plus one for the end.
+                return MAX_HEADER_SIZE
+                        + BinaryVariantUtil.U32_SIZE
+                                * (2L * LogicalTypeChecks.getFieldCount(type) 
+ 1);
+            case CHAR:
+            case VARCHAR:
+                return MAX_HEADER_SIZE
+                        + (long) 
PrimitiveToVariantCastRule.MAX_UTF8_BYTES_PER_CHAR
+                                * LogicalTypeChecks.getLength(type);
+            case BINARY:
+            case VARBINARY:
+                return MAX_HEADER_SIZE + LogicalTypeChecks.getLength(type);
+            default:
+                return MAX_FIXED_SIZE;
+        }
+    }
+
+    /* Example generated code for ARRAY<INT>, inside the null check of the 
base class. The
+    converter is a field, created once when the code is generated:
+
+    result$1 = toVariantConverter$2.convert(array$0);
+
+    */
+    @Override
+    protected String generateCodeBlockInternal(

Review Comment:
   Covered. The test comment says Calcite turns a plain SQL `CAST(ARRAY<ROW> AS 
ARRAY<VARIANT>)` into a MULTISET. What does a user get there, a validation 
error? If that applies to any `ARRAY<ROW>` target, is there a ticket?



##########
flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/functions/casting/ToVariantCastRule.java:
##########
@@ -0,0 +1,178 @@
+/*
+ * 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.flink.table.planner.functions.casting;
+
+import org.apache.flink.table.runtime.functions.ToVariantConverter;
+import org.apache.flink.table.runtime.functions.VariantCastUtils;
+import org.apache.flink.table.types.logical.DistinctType;
+import org.apache.flink.table.types.logical.LogicalType;
+import org.apache.flink.table.types.logical.LogicalTypeRoot;
+import org.apache.flink.table.types.logical.utils.LogicalTypeCasts;
+import org.apache.flink.table.types.logical.utils.LogicalTypeChecks;
+import org.apache.flink.types.variant.BinaryVariantUtil;
+import org.apache.flink.types.variant.Variant;
+
+import java.util.ArrayDeque;
+import java.util.Deque;
+
+import static 
org.apache.flink.table.planner.functions.casting.CastRuleUtils.box;
+import static 
org.apache.flink.table.planner.functions.casting.CastRuleUtils.methodCall;
+
+/**
+ * Cast rule to {@link LogicalTypeRoot#VARIANT} from every type that {@link 
LogicalTypeCasts}
+ * allows: a primitive type, or an {@link LogicalTypeRoot#ARRAY}, {@link 
LogicalTypeRoot#MAP},
+ * {@link LogicalTypeRoot#ROW} or {@link LogicalTypeRoot#STRUCTURED_TYPE} 
whose leaves all cast to
+ * {@code VARIANT} and whose {@code MAP} keys are character strings.
+ *
+ * <p>The generated code calls a {@link ToVariantConverter} created once for 
the input type, which
+ * describes how a value is stored. A failure anywhere in the value fails the 
whole cast, so {@code
+ * TRY_CAST} returns {@code NULL} rather than a partial result. See {@link 
#canFail} for when a cast
+ * can fail.
+ */
+class ToVariantCastRule extends AbstractNullAwareCodeGeneratorCastRule<Object, 
Variant> {
+
+    static final ToVariantCastRule INSTANCE = new ToVariantCastRule();
+
+    /** A character takes up to 4 bytes in UTF-8, which a declared length 
counts as one. */
+    private static final int MAX_UTF8_BYTES_PER_CHAR = 4;
+
+    /** The largest fixed-size value: a 16-byte decimal with its header and 
scale. */
+    private static final long MAX_FIXED_SIZE = 18;
+
+    /** A header byte and a 4-byte size or length, the largest header a value 
has. */
+    private static final long MAX_HEADER_SIZE = 1 + BinaryVariantUtil.U32_SIZE;
+
+    private ToVariantCastRule() {
+        
super(CastRulePredicate.builder().predicate(ToVariantCastRule::matches).build());
+    }
+
+    private static boolean matches(LogicalType input, LogicalType target) {
+        // No cast rule takes the NULL type, although LogicalTypeCasts allows 
it to cast to any
+        // type.
+        return target.is(LogicalTypeRoot.VARIANT)
+                && !input.is(LogicalTypeRoot.NULL)
+                && LogicalTypeCasts.supportsExplicitCast(input, target);
+    }
+
+    /**
+     * Returns whether a value of the input type can fail the cast, so that 
{@code TRY_CAST} returns
+     * {@code NULL} for it instead of failing. A value can fail in three ways:
+     *
+     * <ul>
+     *   <li>A {@code TIMESTAMP(p)} or {@code TIMESTAMP_LTZ(p)} with {@code p 
> 6} is stored with
+     *       nanoseconds, which only cover 1677-09-21 to 2262-04-11.
+     *   <li>The value does not fit into the 16 MiB of a {@code VARIANT}. Any 
{@code ARRAY}, {@code
+     *       MAP} or nested {@code VARIANT} can reach that size. A string or 
binary type can when
+     *       its declared length allows more, and a {@code ROW} when the 
declared sizes of its
+     *       fields add up to more.
+     *   <li>A {@code MAP} with a {@code NULL} key fails. Every {@code MAP} 
can already fail on its
+     *       size, so this makes no other type fallible.
+     * </ul>
+     *
+     * <p>Every other type never fails.
+     *
+     * <p>The size check trusts the declared length. Flink does not enforce 
that length on values
+     * from a source, so a longer value can still reach the cast. {@code 
TRY_CAST} then fails for it
+     * instead of returning {@code NULL}.
+     */
+    @Override
+    public boolean canFail(LogicalType inputLogicalType, LogicalType 
targetLogicalType) {
+        // Without an ARRAY, MAP or VARIANT, the size of a value is bounded by 
the sum of its
+        // leaves and object headers.
+        long maxSize = 0;
+        final Deque<LogicalType> pending = new ArrayDeque<>();
+        pending.push(inputLogicalType);
+        while (!pending.isEmpty()) {
+            final LogicalType type = pending.pop();
+            if (type.is(LogicalTypeRoot.DISTINCT_TYPE)) {
+                pending.push(((DistinctType) type).getSourceType());
+                continue;
+            }
+            switch (type.getTypeRoot()) {
+                case ARRAY:
+                case MAP:
+                case VARIANT:
+                    return true;
+                case ROW:
+                case STRUCTURED_TYPE:
+                    
LogicalTypeChecks.getFieldTypes(type).forEach(pending::push);
+                    break;
+                case TIMESTAMP_WITHOUT_TIME_ZONE:

Review Comment:
   `timeMicros` throws a `TableRuntimeException` for a TIME outside one day, 
but `canFail` treats TIME as infallible, so `TRY_CAST(t AS VARIANT)` fails for 
such a value instead of returning `NULL`. The Javadoc above ("Every other type 
never fails") and footnote 8 in the docs say the same. Such a value only comes 
from a source or a UDF, but FLINK-40845 shows it happens. Could TIME count as 
fallible here, like TIMESTAMP(7..9)? The try/catch it adds costs next to 
nothing.



##########
flink-core/src/main/java/org/apache/flink/types/variant/BinaryVariantInternalBuilder.java:
##########
@@ -537,7 +537,9 @@ public void finishWritingArray(int start, 
ArrayList<Integer> offsets) {
     // the
     // input variant, we can directly copy the binary slice.
     public void appendVariant(BinaryVariant v) {
-        appendVariantImpl(v.getValue(), v.getMetadata(), v.getPos());
+        // A nested variant, such as a field or an element of another variant, 
starts at its own
+        // position in the shared buffer. getValue() would copy it to position 
0 instead.
+        appendVariantImpl(v.rawValue(), v.getMetadata(), v.getPos());

Review Comment:
   FLINK-40911 has no affects or fix versions yet. Please set them so it can be 
picked for the release branches. The `BinaryVariant#hashCode` vs `equals` 
ticket is still missing.



##########
flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/functions/ToVariantConverter.java:
##########
@@ -0,0 +1,260 @@
+/*
+ * 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.flink.table.runtime.functions;
+
+import org.apache.flink.annotation.Internal;
+import org.apache.flink.table.api.TableRuntimeException;
+import org.apache.flink.table.data.ArrayData;
+import org.apache.flink.table.data.DecimalData;
+import org.apache.flink.table.data.MapData;
+import org.apache.flink.table.data.RowData;
+import org.apache.flink.table.data.TimestampData;
+import org.apache.flink.table.types.logical.ArrayType;
+import org.apache.flink.table.types.logical.DistinctType;
+import org.apache.flink.table.types.logical.LogicalType;
+import org.apache.flink.table.types.logical.LogicalTypeRoot;
+import org.apache.flink.table.types.logical.MapType;
+import org.apache.flink.table.types.logical.utils.LogicalTypeChecks;
+import org.apache.flink.table.types.logical.utils.UuidUtils;
+import org.apache.flink.types.variant.BinaryVariant;
+import org.apache.flink.types.variant.BinaryVariantInternalBuilder;
+import org.apache.flink.types.variant.BinaryVariantInternalBuilder.FieldEntry;
+import org.apache.flink.types.variant.Variant;
+import org.apache.flink.types.variant.VariantTypeException;
+
+import java.io.Serializable;
+import java.util.ArrayList;
+import java.util.List;
+
+/**
+ * Converts internal data of a SQL type into a {@link Variant}, like {@code 
CAST(x AS VARIANT)}.
+ *
+ * <p>An {@code ARRAY} becomes a variant array. A {@code ROW} or {@code 
STRUCTURED} value becomes a
+ * variant object keyed by its field names, and a {@code MAP} one keyed by its 
character string
+ * keys. A {@code NULL} element, field or map value becomes a variant null. A 
leaf is stored like
+ * the cast stores it on its own, see {@link VariantCastUtils}, and a nested 
{@code VARIANT} is
+ * embedded as is.
+ *
+ * <p>The converter is created once per type and writes a whole value into one 
builder, so a nested
+ * value is encoded in a single pass. It keeps no state between calls.
+ */
+@Internal
+public final class ToVariantConverter implements Serializable {
+
+    private static final long serialVersionUID = 1L;
+
+    private final String sourceType;
+    private final ValueWriter writer;
+
+    private ToVariantConverter(String sourceType, ValueWriter writer) {
+        this.sourceType = sourceType;
+        this.writer = writer;
+    }
+
+    /** Creates a converter for a type that {@code LogicalTypeCasts} allows to 
cast to VARIANT. */
+    public static ToVariantConverter create(LogicalType type) {
+        return new ToVariantConverter(type.asSummaryString(), 
createWriter(type));
+    }
+
+    /**
+     * Converts a value that is not {@code NULL}.
+     *
+     * @throws TableRuntimeException if the value has a {@code NULL} map key, 
a timestamp outside
+     *     the range of its variant kind, or does not fit into the 16 MiB of a 
{@code VARIANT}
+     */
+    public Variant convert(Object value) {
+        // A MAP can hold the same key twice at runtime. The builder keeps the 
last one, like the
+        // MAP constructor does.
+        final BinaryVariantInternalBuilder builder = new 
BinaryVariantInternalBuilder(true);
+        try {
+            writer.write(builder, value);
+            return builder.build();
+        } catch (VariantTypeException e) {
+            if (e != 
BinaryVariantInternalBuilder.VARIANT_SIZE_LIMIT_EXCEPTION) {
+                throw e;
+            }
+            throw new TableRuntimeException(
+                    String.format(
+                            "Cannot cast a value of type %s to VARIANT. A 
VARIANT is limited to "
+                                    + "16 MiB.",
+                            sourceType));
+        }
+    }
+
+    @FunctionalInterface
+    private interface ValueWriter extends Serializable {
+        void write(BinaryVariantInternalBuilder builder, Object value);
+    }
+
+    private static ValueWriter createNullableWriter(LogicalType type) {
+        final ValueWriter writer = createWriter(type);
+        return (builder, value) -> {
+            if (value == null) {
+                builder.appendNull();
+            } else {
+                writer.write(builder, value);
+            }
+        };
+    }
+
+    private static ValueWriter createWriter(LogicalType type) {
+        switch (type.getTypeRoot()) {
+            case NULL:
+                return (builder, value) -> builder.appendNull();
+            case BOOLEAN:
+                return (builder, value) -> builder.appendBoolean((Boolean) 
value);
+            case TINYINT:
+                return (builder, value) -> builder.appendByte((Byte) value);
+            case SMALLINT:
+                return (builder, value) -> builder.appendShort((Short) value);
+            case INTEGER:
+                return (builder, value) -> builder.appendInt((Integer) value);
+            case BIGINT:
+                return (builder, value) -> builder.appendLong((Long) value);
+            case FLOAT:
+                return (builder, value) -> builder.appendFloat((Float) value);
+            case DOUBLE:
+                return (builder, value) -> builder.appendDouble((Double) 
value);
+            case DECIMAL:
+                return (builder, value) ->
+                        builder.appendDecimal(((DecimalData) 
value).toBigDecimal());
+            case CHAR:
+            case VARCHAR:
+                return (builder, value) -> 
builder.appendString(value.toString());
+            case BINARY:
+            case VARBINARY:
+                return (builder, value) -> builder.appendBinary((byte[]) 
value);
+            case DATE:
+                return (builder, value) -> builder.appendDate((Integer) value);
+            case TIME_WITHOUT_TIME_ZONE:
+                // TIME is millisecond-of-day at runtime, the variant kind 
holds microseconds.
+                return (builder, value) -> builder.appendTime((Integer) value 
* 1_000L);
+            case TIMESTAMP_WITHOUT_TIME_ZONE:
+                {
+                    final int precision = LogicalTypeChecks.getPrecision(type);
+                    return (builder, value) ->
+                            VariantCastUtils.appendTimestamp(
+                                    builder, (TimestampData) value, precision);
+                }
+            case TIMESTAMP_WITH_LOCAL_TIME_ZONE:
+                {
+                    final int precision = LogicalTypeChecks.getPrecision(type);
+                    return (builder, value) ->
+                            VariantCastUtils.appendTimestampLtz(
+                                    builder, (TimestampData) value, precision);
+                }
+            case UUID:
+                return (builder, value) -> 
builder.appendUuid(UuidUtils.fromBytes((byte[]) value));
+            case VARIANT:

Review Comment:
   Verified with a 20,000-level VARIANT. Please link the follow-up ticket here 
once it is filed.



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