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

Reply via email to