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 d8b82e6a546 [FLINK-37977][table] Support `VARIANT` in user-defined
functions and process table functions
d8b82e6a546 is described below
commit d8b82e6a546aa47b50913feffe029f625f080e9b
Author: Ramin Gharib <[email protected]>
AuthorDate: Tue Aug 25 17:21:34 2026 +0200
[FLINK-37977][table] Support `VARIANT` in user-defined functions and
process table functions
---
.../apache/flink/types/variant/BinaryVariant.java | 2 +
.../org/apache/flink/types/variant/Variant.java | 10 +++-
.../flink/types/variant/BinaryVariantTest.java | 17 +++++++
.../flink/table/types/logical/VariantType.java | 3 +-
.../table/types/utils/ClassDataTypeConverter.java | 2 -
.../table/types/utils/ValueDataTypeConverter.java | 4 +-
.../table/types/ClassDataTypeConverterTest.java | 2 -
.../apache/flink/table/types/LogicalTypesTest.java | 5 +-
.../table/types/ValueDataTypeConverterTest.java | 5 +-
.../extraction/TypeInferenceExtractorTest.java | 32 ++++++++++++-
.../flink/table/planner/codegen/CodeGenUtils.scala | 2 +
.../stream/ProcessTableFunctionSemanticTests.java | 5 +-
.../stream/ProcessTableFunctionTestPrograms.java | 54 ++++++++++++++++++++++
.../exec/stream/ProcessTableFunctionTestUtils.java | 36 +++++++++++++++
.../nodes/exec/stream/VariantSemanticTest.java | 36 +++++++++++++++
15 files changed, 196 insertions(+), 19 deletions(-)
diff --git
a/flink-core/src/main/java/org/apache/flink/types/variant/BinaryVariant.java
b/flink-core/src/main/java/org/apache/flink/types/variant/BinaryVariant.java
index c5dd0bf3efe..97c11eedf52 100644
--- a/flink-core/src/main/java/org/apache/flink/types/variant/BinaryVariant.java
+++ b/flink-core/src/main/java/org/apache/flink/types/variant/BinaryVariant.java
@@ -66,6 +66,8 @@ import static
org.apache.flink.types.variant.BinaryVariantUtil.variantConstructo
@Internal
public final class BinaryVariant implements Variant {
+ private static final long serialVersionUID = 1L;
+
private final byte[] value;
private final byte[] metadata;
// The variant value doesn't use the whole `value` binary, but starts from
its `pos` index and
diff --git
a/flink-core/src/main/java/org/apache/flink/types/variant/Variant.java
b/flink-core/src/main/java/org/apache/flink/types/variant/Variant.java
index dd7ac96b6fc..1ff663caaf1 100644
--- a/flink-core/src/main/java/org/apache/flink/types/variant/Variant.java
+++ b/flink-core/src/main/java/org/apache/flink/types/variant/Variant.java
@@ -20,15 +20,21 @@ package org.apache.flink.types.variant;
import org.apache.flink.annotation.PublicEvolving;
+import java.io.Serializable;
import java.math.BigDecimal;
import java.time.Instant;
import java.time.LocalDate;
import java.time.LocalDateTime;
import java.util.List;
-/** Variant represent a semi-structured data. */
+/**
+ * Variant represent a semi-structured data.
+ *
+ * <p>Instances are serializable so that they can be held as member variables
of user-defined
+ * functions or passed into their constructors.
+ */
@PublicEvolving
-public interface Variant {
+public interface Variant extends Serializable {
/** Returns true if the variant is a primitive typed value, such as INT,
DOUBLE, STRING, etc. */
boolean isPrimitive();
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 4468fcbe8d8..4b4a23ad34d 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
@@ -18,6 +18,8 @@
package org.apache.flink.types.variant;
+import org.apache.flink.core.testutils.CommonTestUtils;
+
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.params.ParameterizedTest;
@@ -334,4 +336,19 @@ class BinaryVariantTest {
.isInstanceOf(VariantTypeException.class)
.hasMessage("Expected type DOUBLE but got FLOAT");
}
+
+ @Test
+ void testJavaSerialization() throws Exception {
+ Variant variant =
+ builder.object()
+ .add("i", builder.of(1))
+ .add("nested",
builder.array().add(builder.of("v")).build())
+ .build();
+
assertThat(CommonTestUtils.createCopySerializable(variant)).isEqualTo(variant);
+
+ // a sub-variant is addressed by a position into the value binary of
the enclosing document
+ Variant subVariant = variant.getField("nested");
+ assertThat(((BinaryVariant) subVariant).getPos()).isGreaterThan(0);
+
assertThat(CommonTestUtils.createCopySerializable(subVariant)).isEqualTo(subVariant);
+ }
}
diff --git
a/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/logical/VariantType.java
b/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/logical/VariantType.java
index 45a2f064fb5..623a0b5f709 100644
---
a/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/logical/VariantType.java
+++
b/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/logical/VariantType.java
@@ -19,7 +19,6 @@
package org.apache.flink.table.types.logical;
import org.apache.flink.annotation.PublicEvolving;
-import org.apache.flink.types.variant.BinaryVariant;
import org.apache.flink.types.variant.Variant;
import java.util.Collections;
@@ -40,7 +39,7 @@ import java.util.Set;
public final class VariantType extends LogicalType {
private static final Set<String> INPUT_OUTPUT_CONVERSION =
- conversionSet(Variant.class.getName(),
BinaryVariant.class.getName());
+ conversionSet(Variant.class.getName());
public VariantType(boolean isNullable) {
super(isNullable, LogicalTypeRoot.VARIANT);
diff --git
a/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/utils/ClassDataTypeConverter.java
b/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/utils/ClassDataTypeConverter.java
index efb21532d2e..f0fca13ea29 100644
---
a/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/utils/ClassDataTypeConverter.java
+++
b/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/utils/ClassDataTypeConverter.java
@@ -30,7 +30,6 @@ import org.apache.flink.types.ColumnList;
import org.apache.flink.types.Row;
import org.apache.flink.types.bitmap.Bitmap;
import org.apache.flink.types.bitmap.RoaringBitmapData;
-import org.apache.flink.types.variant.BinaryVariant;
import org.apache.flink.types.variant.Variant;
import java.math.BigDecimal;
@@ -80,7 +79,6 @@ public final class ClassDataTypeConverter {
addDefaultDataType(
java.time.Period.class, DataTypes.INTERVAL(DataTypes.YEAR(4),
DataTypes.MONTH()));
addDefaultDataType(ColumnList.class, DataTypes.DESCRIPTOR());
- addDefaultDataType(BinaryVariant.class, DataTypes.VARIANT());
addDefaultDataType(Variant.class, DataTypes.VARIANT());
addDefaultDataType(Bitmap.class, DataTypes.BITMAP());
addDefaultDataType(RoaringBitmapData.class, DataTypes.BITMAP());
diff --git
a/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/utils/ValueDataTypeConverter.java
b/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/utils/ValueDataTypeConverter.java
index 347f562d6d3..b63e83e039d 100644
---
a/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/utils/ValueDataTypeConverter.java
+++
b/flink-table/flink-table-common/src/main/java/org/apache/flink/table/types/utils/ValueDataTypeConverter.java
@@ -92,7 +92,9 @@ public final class ValueDataTypeConverter {
return convertToArrayType((Object[]) value)
.map(dt -> dt.notNull().bridgedTo(value.getClass()));
} else if (value instanceof Variant) {
- convertedDataType = DataTypes.VARIANT();
+ // BinaryVariant is internal, so the conversion class is the
Variant interface rather
+ // than the runtime class of the value.
+ return Optional.of(DataTypes.VARIANT().notNull());
} else if (value instanceof RoaringBitmapData) {
convertedDataType = DataTypes.BITMAP();
}
diff --git
a/flink-table/flink-table-common/src/test/java/org/apache/flink/table/types/ClassDataTypeConverterTest.java
b/flink-table/flink-table-common/src/test/java/org/apache/flink/table/types/ClassDataTypeConverterTest.java
index 0f52a927c00..ebf629a0f72 100644
---
a/flink-table/flink-table-common/src/test/java/org/apache/flink/table/types/ClassDataTypeConverterTest.java
+++
b/flink-table/flink-table-common/src/test/java/org/apache/flink/table/types/ClassDataTypeConverterTest.java
@@ -25,7 +25,6 @@ import
org.apache.flink.table.types.utils.ClassDataTypeConverter;
import org.apache.flink.types.Row;
import org.apache.flink.types.bitmap.Bitmap;
import org.apache.flink.types.bitmap.RoaringBitmapData;
-import org.apache.flink.types.variant.BinaryVariant;
import org.apache.flink.types.variant.Variant;
import org.junit.jupiter.params.ParameterizedTest;
@@ -96,7 +95,6 @@ class ClassDataTypeConverterTest {
new AtomicDataType(new
SymbolType<>()).bridgedTo(TimeIntervalUnit.class)),
of(Row.class, null),
of(Variant.class, DataTypes.VARIANT()),
- of(BinaryVariant.class,
DataTypes.VARIANT().bridgedTo(BinaryVariant.class)),
of(Bitmap.class, DataTypes.BITMAP().bridgedTo(Bitmap.class)),
of(RoaringBitmapData.class,
DataTypes.BITMAP().bridgedTo(RoaringBitmapData.class)));
}
diff --git
a/flink-table/flink-table-common/src/test/java/org/apache/flink/table/types/LogicalTypesTest.java
b/flink-table/flink-table-common/src/test/java/org/apache/flink/table/types/LogicalTypesTest.java
index ee40463cf0b..8b06ceb5118 100644
---
a/flink-table/flink-table-common/src/test/java/org/apache/flink/table/types/LogicalTypesTest.java
+++
b/flink-table/flink-table-common/src/test/java/org/apache/flink/table/types/LogicalTypesTest.java
@@ -65,7 +65,6 @@ import
org.apache.flink.table.types.logical.ZonedTimestampType;
import org.apache.flink.types.Row;
import org.apache.flink.types.bitmap.Bitmap;
import org.apache.flink.types.bitmap.RoaringBitmapData;
-import org.apache.flink.types.variant.BinaryVariant;
import org.apache.flink.types.variant.Variant;
import org.assertj.core.api.ThrowingConsumer;
@@ -610,9 +609,7 @@ public class LogicalTypesTest {
.hasSerializableString("VARIANT")
.hasSummaryString("VARIANT")
.supportsOutputConversion(Variant.class)
- .supportsOutputConversion(BinaryVariant.class)
- .supportsInputConversion(Variant.class)
- .supportsInputConversion(BinaryVariant.class);
+ .supportsInputConversion(Variant.class);
}
@Test
diff --git
a/flink-table/flink-table-common/src/test/java/org/apache/flink/table/types/ValueDataTypeConverterTest.java
b/flink-table/flink-table-common/src/test/java/org/apache/flink/table/types/ValueDataTypeConverterTest.java
index 9f767c44afd..58e3b5bc38a 100644
---
a/flink-table/flink-table-common/src/test/java/org/apache/flink/table/types/ValueDataTypeConverterTest.java
+++
b/flink-table/flink-table-common/src/test/java/org/apache/flink/table/types/ValueDataTypeConverterTest.java
@@ -26,7 +26,6 @@ import org.apache.flink.table.types.logical.SymbolType;
import org.apache.flink.table.types.utils.ValueDataTypeConverter;
import org.apache.flink.types.bitmap.Bitmap;
import org.apache.flink.types.bitmap.RoaringBitmapData;
-import org.apache.flink.types.variant.BinaryVariant;
import org.apache.flink.types.variant.Variant;
import org.junit.jupiter.params.ParameterizedTest;
@@ -121,9 +120,7 @@ class ValueDataTypeConverterTest {
DataTypes.ARRAY(DataTypes.ARRAY(DataTypes.INT()))),
of(TimePointUnit.HOUR, new AtomicDataType(new SymbolType<>(),
TimePointUnit.class)),
of(new BigDecimal[0], null),
- of(
- Variant.newBuilder().of("hello"),
- DataTypes.VARIANT().bridgedTo(BinaryVariant.class)),
+ of(Variant.newBuilder().of("hello"), DataTypes.VARIANT()),
of(Bitmap.empty(),
DataTypes.BITMAP().bridgedTo(RoaringBitmapData.class)),
of(
Bitmap.fromArray(new int[] {1, 2}),
diff --git
a/flink-table/flink-table-common/src/test/java/org/apache/flink/table/types/extraction/TypeInferenceExtractorTest.java
b/flink-table/flink-table-common/src/test/java/org/apache/flink/table/types/extraction/TypeInferenceExtractorTest.java
index f253819ce2c..3b7c10ab8ac 100644
---
a/flink-table/flink-table-common/src/test/java/org/apache/flink/table/types/extraction/TypeInferenceExtractorTest.java
+++
b/flink-table/flink-table-common/src/test/java/org/apache/flink/table/types/extraction/TypeInferenceExtractorTest.java
@@ -50,6 +50,7 @@ import org.apache.flink.table.types.inference.TypeStrategy;
import org.apache.flink.table.types.utils.DataTypeFactoryMock;
import org.apache.flink.types.Row;
import org.apache.flink.types.bitmap.Bitmap;
+import org.apache.flink.types.variant.Variant;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.MethodSource;
@@ -889,7 +890,25 @@ class TypeInferenceExtractorTest {
"Logical type 'BITMAP' does not support a
conversion from or to class
'org.apache.flink.table.types.extraction.TypeInferenceExtractorTest$CustomBitmap'."),
TestSpec.forScalarFunction("Custom Bitmap",
InvalidCustomBitmapTypeFunction2.class)
.expectErrorMessage(
- "Could not extract a valid type inference for
function class
'org.apache.flink.table.types.extraction.TypeInferenceExtractorTest$InvalidCustomBitmapTypeFunction2'."));
+ "Could not extract a valid type inference for
function class
'org.apache.flink.table.types.extraction.TypeInferenceExtractorTest$InvalidCustomBitmapTypeFunction2'."),
+ // ---
+ TestSpec.forScalarFunction("Variant in scalar function",
VariantTypeFunction.class)
+ .expectStaticArgument(
+ StaticArgument.scalar("v",
DataTypes.VARIANT(), false))
+ .expectStaticArgument(
+ StaticArgument.scalar(
+ "array",
DataTypes.ARRAY(DataTypes.VARIANT()), false))
+ .expectStaticArgument(
+ StaticArgument.scalar(
+ "map",
+ DataTypes.MAP(DataTypes.INT(),
DataTypes.VARIANT()),
+ false))
+ .expectStaticArgument(
+ StaticArgument.scalar(
+ "row",
+ DataTypes.ROW(DataTypes.FIELD("a",
DataTypes.VARIANT())),
+ false))
+
.expectOutput(TypeStrategies.explicit(DataTypes.VARIANT())));
}
private static Stream<TestSpec> procedureSpecs() {
@@ -2723,6 +2742,17 @@ class TypeInferenceExtractorTest {
}
}
+ @FunctionHint(output = @DataTypeHint("VARIANT"))
+ private static class VariantTypeFunction extends ScalarFunction {
+ public Variant eval(
+ Variant v,
+ Variant[] array,
+ Map<Integer, Variant> map,
+ @DataTypeHint("ROW<a VARIANT>") Row row) {
+ return null;
+ }
+ }
+
@FunctionHint(input = @DataTypeHint(value = "BITMAP", bridgedTo =
CustomBitmap.class))
private static class InvalidCustomBitmapTypeFunction1 extends
ScalarFunction {
public Bitmap eval(Bitmap bitmap) {
diff --git
a/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/CodeGenUtils.scala
b/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/CodeGenUtils.scala
index aacb8db5a0d..739858689bf 100644
---
a/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/CodeGenUtils.scala
+++
b/flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/codegen/CodeGenUtils.scala
@@ -388,6 +388,8 @@ object CodeGenUtils {
}
val serTerm = ctx.addReusableObject(serializer, "serializer")
s"$term.toObject($serTerm).hashCode()"
+ case VARIANT =>
+ s"$term.hashCode()"
case BITMAP =>
s"$term.hashCode()"
case NULL | SYMBOL | UNRESOLVED =>
diff --git
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/ProcessTableFunctionSemanticTests.java
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/ProcessTableFunctionSemanticTests.java
index 4b6a105cb2e..45020941d47 100644
---
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/ProcessTableFunctionSemanticTests.java
+++
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/ProcessTableFunctionSemanticTests.java
@@ -98,6 +98,9 @@ public class ProcessTableFunctionSemanticTests extends
SemanticTestBase {
ProcessTableFunctionTestPrograms.PROCESS_MULTI_INPUT_ORDER_BY,
ProcessTableFunctionTestPrograms.PROCESS_ORDER_BY_TABLE_API,
ProcessTableFunctionTestPrograms.PROCESS_IMPLICIT_CASTS,
-
ProcessTableFunctionTestPrograms.PROCESS_ROW_DATA_CONVERSION_TABLE);
+
ProcessTableFunctionTestPrograms.PROCESS_ROW_DATA_CONVERSION_TABLE,
+ ProcessTableFunctionTestPrograms.PROCESS_VARIANT,
+ ProcessTableFunctionTestPrograms.PROCESS_VARIANT_TABLE_ARG,
+ ProcessTableFunctionTestPrograms.PROCESS_VARIANT_STATE);
}
}
diff --git
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/ProcessTableFunctionTestPrograms.java
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/ProcessTableFunctionTestPrograms.java
index 5a9b260c220..2fd78546406 100644
---
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/ProcessTableFunctionTestPrograms.java
+++
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/ProcessTableFunctionTestPrograms.java
@@ -73,6 +73,9 @@ import
org.apache.flink.table.planner.plan.nodes.exec.stream.ProcessTableFunctio
import
org.apache.flink.table.planner.plan.nodes.exec.stream.ProcessTableFunctionTestUtils.UpdatingJoinFunction;
import
org.apache.flink.table.planner.plan.nodes.exec.stream.ProcessTableFunctionTestUtils.UpdatingRetractFunction;
import
org.apache.flink.table.planner.plan.nodes.exec.stream.ProcessTableFunctionTestUtils.UpdatingUpsertFunction;
+import
org.apache.flink.table.planner.plan.nodes.exec.stream.ProcessTableFunctionTestUtils.VariantFunction;
+import
org.apache.flink.table.planner.plan.nodes.exec.stream.ProcessTableFunctionTestUtils.VariantStateFunction;
+import
org.apache.flink.table.planner.plan.nodes.exec.stream.ProcessTableFunctionTestUtils.VariantTableArgFunction;
import org.apache.flink.table.test.program.SinkTestStep;
import org.apache.flink.table.test.program.SourceTestStep;
import org.apache.flink.table.test.program.TableTestProgram;
@@ -929,6 +932,57 @@ public class ProcessTableFunctionTestPrograms {
"INSERT INTO sink SELECT * FROM f(columnList1 =>
NULL, columnList3 => DESCRIPTOR(a, b, c))")
.build();
+ public static final TableTestProgram PROCESS_VARIANT =
+ TableTestProgram.of(
+ "process-variant",
+ "takes nullable, optional, and not nullable
VARIANT arguments")
+ .setupTemporarySystemFunction("f", VariantFunction.class)
+ .setupSql(BASIC_VALUES)
+ .setupTableSink(
+ SinkTestStep.newBuilder("sink")
+ .addSchema(BASE_SINK_SCHEMA)
+ .consumedValues("+I[{null, null,
{\"a\":[1,\"b\"]}}]")
+ .build())
+ .runSql(
+ "INSERT INTO sink SELECT * FROM f("
+ + "variant1 => NULL, "
+ + "variant3 =>
PARSE_JSON('{\"a\":[1,\"b\"]}'))")
+ .build();
+
+ public static final TableTestProgram PROCESS_VARIANT_TABLE_ARG =
+ TableTestProgram.of("process-variant-table-arg", "table argument
with a VARIANT column")
+ .setupTemporarySystemFunction("f",
VariantTableArgFunction.class)
+ .setupSql(
+ "CREATE VIEW t AS SELECT * FROM "
+ + "(VALUES ('Bob',
PARSE_JSON('{\"a\":1}')), "
+ + "('Alice', PARSE_JSON('[1,\"b\"]'))) AS
T(name, v)")
+ .setupTableSink(
+ SinkTestStep.newBuilder("sink")
+ .addSchema(KEYED_BASE_SINK_SCHEMA)
+ .consumedValues(
+ "+I[Bob, {+I[Bob, {\"a\":1}]}]",
+ "+I[Alice, {+I[Alice,
[1,\"b\"]]}]")
+ .build())
+ .runSql("INSERT INTO sink SELECT * FROM f(r => TABLE t
PARTITION BY name)")
+ .build();
+
+ public static final TableTestProgram PROCESS_VARIANT_STATE =
+ TableTestProgram.of("process-variant-state", "state entry with a
VARIANT field")
+ .setupTemporarySystemFunction("f",
VariantStateFunction.class)
+ .setupSql(MULTI_VALUES)
+ .setupTableSink(
+ SinkTestStep.newBuilder("sink")
+ .addSchema(KEYED_BASE_SINK_SCHEMA)
+ .consumedValues(
+ "+I[Bob, {VariantScore(v=null),
+I[Bob, 12]}]",
+ "+I[Alice, {VariantScore(v=null),
+I[Alice, 42]}]",
+ "+I[Bob, {VariantScore(v=12),
+I[Bob, 99]}]",
+ "+I[Bob, {VariantScore(v=99),
+I[Bob, 100]}]",
+ "+I[Alice, {VariantScore(v=42),
+I[Alice, 400]}]")
+ .build())
+ .runSql("INSERT INTO sink SELECT * FROM f(r => TABLE t
PARTITION BY name)")
+ .build();
+
public static final TableTestProgram PROCESS_TIME_CONVERSIONS =
TableTestProgram.of(
"process-time-conversions",
diff --git
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/ProcessTableFunctionTestUtils.java
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/ProcessTableFunctionTestUtils.java
index 57d185fd78d..6bbcf9a926d 100644
---
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/ProcessTableFunctionTestUtils.java
+++
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/ProcessTableFunctionTestUtils.java
@@ -42,6 +42,7 @@ import org.apache.flink.table.types.inference.TypeInference;
import org.apache.flink.types.ColumnList;
import org.apache.flink.types.Row;
import org.apache.flink.types.RowKind;
+import org.apache.flink.types.variant.Variant;
import java.time.Duration;
import java.time.Instant;
@@ -550,6 +551,31 @@ public class ProcessTableFunctionTestUtils {
}
}
+ /** Testing function. */
+ public static class VariantFunction extends AppendProcessTableFunctionBase
{
+ public void eval(
+ Variant variant1,
+ @ArgumentHint(isOptional = true) Variant variant2,
+ @DataTypeHint("VARIANT NOT NULL") Variant variant3) {
+ collectObjects(variant1, variant2, variant3);
+ }
+ }
+
+ /** Testing function. */
+ public static class VariantStateFunction extends
AppendProcessTableFunctionBase {
+ public void eval(@StateHint VariantScore s,
@ArgumentHint(SET_SEMANTIC_TABLE) Row r) {
+ collectObjects(s, r);
+ s.v = Variant.newBuilder().of(r.<Integer>getFieldAs("score"));
+ }
+ }
+
+ /** Testing function. */
+ public static class VariantTableArgFunction extends
AppendProcessTableFunctionBase {
+ public void eval(@ArgumentHint(SET_SEMANTIC_TABLE) Row r) {
+ collectObjects(r);
+ }
+ }
+
/** Testing function. */
public static class RequiredTimeFunction extends
AppendProcessTableFunctionBase {
public void eval(@ArgumentHint({ArgumentTrait.ROW_SEMANTIC_TABLE,
REQUIRE_ON_TIME}) Row r) {
@@ -1244,6 +1270,16 @@ public class ProcessTableFunctionTestUtils {
}
}
+ /** POJO for state. */
+ public static class VariantScore {
+ public Variant v;
+
+ @Override
+ public String toString() {
+ return String.format("VariantScore(v=%s)", v);
+ }
+ }
+
private static final Map<String, String> MODE_SUMMARY =
Map.ofEntries(
Map.entry("[INSERT, UPDATE_BEFORE, UPDATE_AFTER]",
"retract-no-delete"),
diff --git
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/VariantSemanticTest.java
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/VariantSemanticTest.java
index af95ee23a65..d9f49877bdd 100644
---
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/VariantSemanticTest.java
+++
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/VariantSemanticTest.java
@@ -193,6 +193,35 @@ public class VariantSemanticTest extends SemanticTestBase {
.runSql("INSERT INTO sink_t SELECT udf(v) FROM t")
.build();
+ static final TableTestProgram VARIANT_IN_VIEW =
+ TableTestProgram.of("variant-in-view", "validates variant
unparsing through a view")
+ .setupTemporarySystemFunction("udf", VariantIdentity.class)
+ .setupTableSource(
+ SourceTestStep.newBuilder("t")
+ .addSchema("v VARIANT")
+ .producedValues(
+ Row.of(
+ BUILDER.object()
+ .add("k",
BUILDER.of(1))
+ .build()),
+
Row.of(BUILDER.array().add(BUILDER.of("x")).build()),
+ new Row(1))
+ .build())
+ .setupSql("CREATE VIEW variant_view AS SELECT udf(v) AS v
FROM t")
+ .setupTableSink(
+ SinkTestStep.newBuilder("sink_t")
+ .addSchema("v VARIANT")
+ .consumedValues(
+ Row.of(
+ BUILDER.object()
+ .add("k",
BUILDER.of(1))
+ .build()),
+
Row.of(BUILDER.array().add(BUILDER.of("x")).build()),
+ new Row(1))
+ .build())
+ .runSql("INSERT INTO sink_t SELECT v FROM variant_view")
+ .build();
+
static final TableTestProgram VARIANT_AS_UDAF_ARG =
TableTestProgram.of("variant-as-udaf-arg", "validates variant as
udaf argument")
.setupTemporarySystemFunction("udf", MyAggFunc.class)
@@ -458,6 +487,7 @@ public class VariantSemanticTest extends SemanticTestBase {
BUILTIN_AGG,
BUILTIN_AGG_WITH_RETRACTION,
VARIANT_AS_UDF_ARG,
+ VARIANT_IN_VIEW,
VARIANT_AS_UDAF_ARG,
VARIANT_AS_AGG_KEY,
VARIANT_ARRAY_ACCESS,
@@ -468,6 +498,12 @@ public class VariantSemanticTest extends SemanticTestBase {
VARIANT_OBJECT_ERROR_ACCESS);
}
+ public static class VariantIdentity extends ScalarFunction {
+ public Variant eval(Variant v) {
+ return v;
+ }
+ }
+
public static class MyUdf extends ScalarFunction {
public Integer eval(Variant v) {