raminqaf commented on code in PR #29369:
URL: https://github.com/apache/flink/pull/29369#discussion_r4193547767
##########
docs/content.zh/docs/sql/reference/data-types.md:
##########
@@ -1717,7 +1719,46 @@ CAST(NULL AS VARIANT) -- NULL
CAST(CAST('NaN' AS DOUBLE) AS VARIANT) -- NaN, stored as a DOUBLE
CAST(INTERVAL '2' DAY AS VARIANT) -- fails at validation
CAST(ARRAY[1, NULL] AS ARRAY<VARIANT>) -- [1, NULL], each element a
VARIANT, the NULL stays SQL NULL
-CAST(ARRAY[1, 2] AS VARIANT) -- fails at validation, not
supported yet
+```
+
+A whole `ARRAY`, `MAP`, `ROW`, or `STRUCTURED` value can also be cast into a
single `VARIANT`. An
+`ARRAY` becomes a variant array, and a `MAP`, `ROW`, or `STRUCTURED` value a
variant object. Each
+leaf is stored by the rules above, so the cast is supported only when every
leaf type casts to
+`VARIANT`.
+
+- A `ROW` or `STRUCTURED` value is keyed by its field names. The SQL `ROW`
constructor names its
+ fields `EXPR$0`, `EXPR$1`, and so on, and the Table API `row()` names them
`f0`, `f1`, and so on.
+ To choose the keys, cast to a `ROW` with named fields first, or name each
field with `as()` in the
+ Table API.
Review Comment:
Done. I also removed the two `EXPR$0` examples, since they show the same
thing.
##########
docs/content.zh/docs/sql/reference/data-types.md:
##########
@@ -1717,7 +1719,46 @@ CAST(NULL AS VARIANT) -- NULL
CAST(CAST('NaN' AS DOUBLE) AS VARIANT) -- NaN, stored as a DOUBLE
CAST(INTERVAL '2' DAY AS VARIANT) -- fails at validation
CAST(ARRAY[1, NULL] AS ARRAY<VARIANT>) -- [1, NULL], each element a
VARIANT, the NULL stays SQL NULL
-CAST(ARRAY[1, 2] AS VARIANT) -- fails at validation, not
supported yet
+```
+
+A whole `ARRAY`, `MAP`, `ROW`, or `STRUCTURED` value can also be cast into a
single `VARIANT`. An
+`ARRAY` becomes a variant array, and a `MAP`, `ROW`, or `STRUCTURED` value a
variant object. Each
+leaf is stored by the rules above, so the cast is supported only when every
leaf type casts to
+`VARIANT`.
+
+- A `ROW` or `STRUCTURED` value is keyed by its field names. The SQL `ROW`
constructor names its
+ fields `EXPR$0`, `EXPR$1`, and so on, and the Table API `row()` names them
`f0`, `f1`, and so on.
+ To choose the keys, cast to a `ROW` with named fields first, or name each
field with `as()` in the
+ Table API.
+- A `MAP` needs a character string key, which becomes the object key. A `NULL`
key fails the cast.
+ If a key appears twice, the last value is kept. The order of the entries is
part of the encoding,
+ so two `MAP`s with the same entries in a different order cast to `VARIANT`
values that are not
Review Comment:
Both are sorted. A ROW and a MAP both become a variant object, and the spec
orders object fields by key. The docs wrongly suggested that a MAP keeps its
order. I fixed that, and the example is now `MAP['b', 2, 'a', 1]`, which comes
back as `{"a": 1, "b": 2}`.
The MAP note was only about equality. The entry order still shapes the
bytes, so the same entries in a different order give VARIANTs that are not
equal. Sorting the entries first would fix that, at the cost of a sort per
value. I would leave that for a follow-up.
##########
docs/content.zh/docs/sql/reference/data-types.md:
##########
@@ -1949,6 +1990,7 @@ COALESCE(TRY_CAST('non-number' AS INT), 0) --- 结果返回数字 0 的
INT 格
5. 支持转换,当且仅当用使用 `INTERVAL` 做“月”到“年”的转换。
6. 支持转换,当且仅当用使用 `INTERVAL` 做“天”到“时间”的转换。
7. 仅支持转换到无界的 `VARBINARY`(`BYTES`),因为裁剪或填充会破坏序列化的位图数据。
+8. 支持转换,当且仅当所有叶子类型都支持转换为 `VARIANT`,并且所有 `MAP` 的键都是字符串类型。`ARRAY` 和 `MAP`
的转换总是可能会失败。`ROW` 和 `STRUCTURED` 的转换可能会失败,当且仅当某个字段可能转换失败、某个字段是
`VARIANT`,或者字段声明的总大小可能超过 `VARIANT` 的 16 MiB 上限。
Review Comment:
Done, reverted.
##########
flink-table/flink-table-common/src/main/java/org/apache/flink/table/data/utils/ToVariantConverter.java:
##########
@@ -0,0 +1,437 @@
+/*
+ * 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.data.utils;
+
+import org.apache.flink.annotation.Internal;
+import org.apache.flink.core.memory.MemorySegment;
+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.StringData;
+import org.apache.flink.table.data.TimestampData;
+import org.apache.flink.table.data.binary.BinaryStringData;
+import org.apache.flink.table.data.binary.StringUtf8Utils;
+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.LogicalTypeFamily;
+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.BinaryVariantUtil;
+import org.apache.flink.types.variant.Variant;
+import org.apache.flink.types.variant.VariantTypeException;
+
+import java.io.Serializable;
+import java.nio.charset.StandardCharsets;
+import java.util.ArrayList;
+import java.util.List;
+
+import static org.apache.flink.table.utils.DateTimeUtils.MILLIS_PER_DAY;
+
+/**
+ * Converts internal data of a SQL type into a {@link Variant}. It implements
{@code CAST(x AS
+ * VARIANT)}, and a format can use it to store a value in a {@code VARIANT}
the same way.
+ *
+ * <p>A leaf keeps the kind of its SQL type. For example, a {@code BIGINT} is
stored as a {@code
+ * BIGINT}, even when it would fit a smaller kind. A {@code TIMESTAMP(p)} or
{@code
+ * TIMESTAMP_LTZ(p)} is stored with microseconds up to {@code p = 6} and with
nanoseconds above,
+ * whatever the digits of the value. A {@code NaN} or infinite {@code FLOAT}
or {@code DOUBLE} is
+ * stored as is, although {@code PARSE_JSON} rejects it.
+ *
+ * <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} becomes a variant null, and a nested {@code VARIANT}
is copied in as is.
+ *
+ * <p>The converter is created once per type and writes a whole value into one
builder. A nested
+ * {@code ARRAY}, {@code MAP} or {@code ROW} is written in place rather than
built as a separate
+ * variant first. The converter keeps no state between calls.
+ */
+@Internal
+public final class ToVariantConverter implements Serializable {
+
+ private static final long serialVersionUID = 1L;
+
+ /** The highest timestamp precision a {@code VARIANT} stores with
microseconds. */
+ public static final int TIMESTAMP_PRECISION = 6;
+
+ /**
+ * The largest string or binary value a {@code VARIANT} holds, in bytes. A
{@code VARIANT} is
+ * limited to 16 MiB, and a long string or binary value spends 5 bytes of
that on its header.
+ */
+ public static final int MAX_PAYLOAD_BYTES =
+ BinaryVariantUtil.SIZE_LIMIT - 1 - BinaryVariantUtil.U32_SIZE;
+
+ 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.
+ *
+ * @throws IllegalArgumentException if the type cannot be cast to VARIANT
+ */
+ public static ToVariantConverter create(LogicalType type) {
+ return new ToVariantConverter(type.asSummaryString(),
createNullableWriter(type));
+ }
+
+ /**
+ * Converts a value into a new {@code VARIANT}. A {@code null} value
becomes a variant null.
+ *
+ * @throws TableRuntimeException if the value has a {@code NULL} map key,
a time or timestamp
+ * outside the range of its variant kind, does not fit into the 16 MiB
of a {@code VARIANT},
+ * or is nested too deeply
+ * @throws VariantTypeException if a nested {@code VARIANT} is malformed
+ */
+ 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 {
+ writeTo(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));
+ } catch (StackOverflowError e) {
+ // The builder copies an embedded VARIANT with one call per
nesting level, so a deeply
+ // nested one overflows the stack. A bare StackOverflowError would
not say which cast
+ // failed or why.
+ throw new TableRuntimeException(
+ String.format(
+ "Cannot cast a value of type %s to VARIANT because
it is nested too "
+ + "deeply.",
+ sourceType));
+ }
+ }
+
+ /**
+ * Writes a value into a builder that the caller owns, for example as a
field of a larger {@code
+ * VARIANT}. A {@code null} value becomes a variant null.
+ *
+ * <p>It fails like {@link #convert} for a {@code NULL} map key, a time or
timestamp out of
+ * range, and a string or binary value that can never fit. The caller
builds the result, so it
+ * handles {@link
BinaryVariantInternalBuilder#VARIANT_SIZE_LIMIT_EXCEPTION} and a {@link
+ * StackOverflowError} from a deeply nested {@code VARIANT} itself. A
repeated {@code MAP} key
+ * keeps the last value if the caller's builder allows duplicate keys, and
fails otherwise. The
+ * builder holds partial data after a failure, so the caller must not use
it any further.
+ */
+ public void writeTo(BinaryVariantInternalBuilder builder, Object value) {
+ writer.write(builder, value);
+ }
+
+ @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) -> appendDecimal(builder,
(DecimalData) value);
+ case CHAR:
+ case VARCHAR:
+ return (builder, value) -> appendString(builder, (StringData)
value);
+ case BINARY:
+ case VARBINARY:
+ return (builder, value) -> appendBinary(builder, (byte[])
value);
+ case DATE:
+ return (builder, value) -> builder.appendDate((Integer) value);
+ case TIME_WITHOUT_TIME_ZONE:
+ return (builder, value) ->
builder.appendTime(timeMicros((Integer) value));
+ case TIMESTAMP_WITHOUT_TIME_ZONE:
+ {
+ final int precision = LogicalTypeChecks.getPrecision(type);
+ if (precision <= TIMESTAMP_PRECISION) {
+ return (builder, value) ->
+
builder.appendTimestamp(timestampMicros((TimestampData) value));
+ }
+ final String timestampType = "TIMESTAMP(" + precision +
")";
+ return (builder, value) ->
+ builder.appendTimestampNanos(
+ timestampNanos((TimestampData) value,
timestampType));
+ }
+ case TIMESTAMP_WITH_LOCAL_TIME_ZONE:
+ {
+ final int precision = LogicalTypeChecks.getPrecision(type);
+ if (precision <= TIMESTAMP_PRECISION) {
+ return (builder, value) ->
+
builder.appendTimestampLtz(timestampMicros((TimestampData) value));
+ }
+ final String timestampType = "TIMESTAMP_LTZ(" + precision
+ ")";
+ return (builder, value) ->
+ builder.appendTimestampLtzNanos(
+ timestampNanos((TimestampData) value,
timestampType));
+ }
+ case UUID:
+ return (builder, value) ->
builder.appendUuid(UuidUtils.fromBytes((byte[]) value));
+ case VARIANT:
+ return (builder, value) ->
builder.appendVariant((BinaryVariant) value);
+ case ARRAY:
+ return createArrayWriter((ArrayType) type);
+ case MAP:
+ return createMapWriter((MapType) type);
+ case ROW:
+ case STRUCTURED_TYPE:
+ return createRowWriter(type);
+ case DISTINCT_TYPE:
+ return createWriter(((DistinctType) type).getSourceType());
+ default:
+ throw new IllegalArgumentException("Cannot cast " + type + "
to VARIANT.");
+ }
+ }
+
+ // A compact decimal is written from its unscaled long, which saves
building a BigDecimal.
+ private static void appendDecimal(BinaryVariantInternalBuilder builder,
DecimalData value) {
+ if (value.isCompact()) {
+ builder.appendDecimal(value.toUnscaledLong(), value.scale());
+ } else {
+ builder.appendDecimal(value.toBigDecimal());
+ }
+ }
+
+ private static void appendString(BinaryVariantInternalBuilder builder,
StringData value) {
+ final BinaryStringData string = (BinaryStringData) value;
+ try {
+ appendUtf8(builder, string);
+ } catch (VariantTypeException e) {
+ throw payloadTooLarge(e, "string", string.getSizeInBytes());
+ }
+ }
+
+ /**
+ * 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.
+ */
+ private static void appendUtf8(BinaryVariantInternalBuilder builder,
BinaryStringData string) {
+ if (string.getBinarySection() == null) {
Review Comment:
Moved it to `StringUtf8Utils#toValidUtf8Bytes(StringData)`, built on
`StringData#toBytes()`. It keeps the shortcut for Java-object strings from
#29370, but drops the in-place segment read, so a string in a binary row is
copied once more.
##########
flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/CodeGeneratorContext.scala:
##########
@@ -834,6 +839,31 @@ class CodeGeneratorContext(
addReusableObjectWithName(obj, newName(this, fieldNamePrefix),
fieldTypeTerm)
}
+ /**
+ * Adds a reusable Object to the member area of the generated class, unless
one was added for an
+ * equal key before. The factory only runs when the object is added.
+ *
+ * @param key
+ * the key that identifies the object
+ * @param factory
+ * creates the object to be added to the generated class
+ * @param fieldNamePrefix
+ * prefix field name of the generated member field term
+ * @param fieldTypeTerm
+ * field type class name
+ * @return
+ * the field term of the object for the key
+ */
+ def addReusableObjectIfAbsent(
Review Comment:
You're right, `addReusableObject` works. The new method only shared one
converter per type, which saves little. I removed it, so `CodeGeneratorContext`
is unchanged.
--
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]