Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -66,6 +66,8 @@
@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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,15 +20,21 @@

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 {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

why don't we use Value?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Value extends IOReadableWritable, whose read(DataInputView) deserializes into this. ValueSerializer also instantiates via a public nullary constructor, which BinaryVariant cannot offer, since its constructor validates the version byte and the size limit.
It would also add a second serialization path. VARIANT already has VariantTypeInfo and VariantSerializer, and TypeExtractor would resolve to ValueTypeInfo instead, because it checks Value first.


/** Returns true if the variant is a primitive typed value, such as INT, DOUBLE, STRING, etc. */
boolean isPrimitive();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -285,4 +287,19 @@ void testGetThrowException() {
.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);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,7 @@
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;
Expand Down Expand Up @@ -889,7 +890,56 @@ private static Stream<TestSpec> functionSpecs() {
"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())),
// ---
TestSpec.forAsyncScalarFunction(

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

only test scalar function should be enough.

"Variant in async scalar function", VariantTypeAsyncFunction.class)
.expectStaticArgument(
StaticArgument.scalar("v", DataTypes.VARIANT(), false))
.expectOutput(TypeStrategies.explicit(DataTypes.VARIANT())),
// ---
TestSpec.forAggregateFunction(
"Variant in aggregate function", VariantTypeAggFunction.class)
.expectStaticArgument(
StaticArgument.scalar("v", DataTypes.VARIANT(), false))
.expectAccumulator(TypeStrategies.explicit(VariantState.TYPE))
.expectOutput(TypeStrategies.explicit(DataTypes.VARIANT())),
// ---
TestSpec.forTableFunction(
"Variant in table function", VariantTypeTableFunction.class)
.expectStaticArgument(
StaticArgument.scalar("v", DataTypes.VARIANT(), false))
.expectOutput(
TypeStrategies.explicit(
DataTypes.ROW(DataTypes.FIELD("v", DataTypes.VARIANT())))),
// ---
TestSpec.forProcessTableFunction(VariantProcessTableFunction.class)
.expectStaticArgument(
StaticArgument.scalar("v", DataTypes.VARIANT(), false))
.expectState("s", TypeStrategies.explicit(VariantState.TYPE))
.expectOutput(TypeStrategies.explicit(DataTypes.VARIANT())),
// ---
TestSpec.forProcessTableFunction(InvalidVariantStateProcessTableFunction.class)
.expectErrorMessage(
"State entries must use a mutable, composite data type. But was: VARIANT"));
}

private static Stream<TestSpec> procedureSpecs() {
Expand Down Expand Up @@ -2723,6 +2773,55 @@ public Bitmap[] call(Object procedureContext, Bitmap bitmap) {
}
}

@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;
}
}

private static class VariantTypeAsyncFunction extends AsyncScalarFunction {
public void eval(CompletableFuture<Variant> f, Variant v) {}
}

private static class VariantTypeAggFunction extends AggregateFunction<Variant, VariantState> {
public void accumulate(VariantState accumulator, Variant v) {}

@Override
public VariantState createAccumulator() {
return null;
}

@Override
public Variant getValue(VariantState accumulator) {
return null;
}
}

@FunctionHint(output = @DataTypeHint("ROW<v VARIANT>"))
private static class VariantTypeTableFunction extends TableFunction<Row> {
public void eval(Variant v) {}
}

private static class VariantProcessTableFunction extends ProcessTableFunction<Variant> {
public void eval(@StateHint VariantState s, Variant v) {}
}

private static class InvalidVariantStateProcessTableFunction
extends ProcessTableFunction<Variant> {
public void eval(@StateHint Variant s, Variant v) {}
}

public static class VariantState {
static final DataType TYPE =
DataTypes.STRUCTURED(VariantState.class, DataTypes.FIELD("v", DataTypes.VARIANT()));
public Variant v;
}

@FunctionHint(input = @DataTypeHint(value = "BITMAP", bridgedTo = CustomBitmap.class))
private static class InvalidCustomBitmapTypeFunction1 extends ScalarFunction {
public Bitmap eval(Bitmap bitmap) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 =>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -97,6 +97,8 @@ public List<TableTestProgram> programs() {
ProcessTableFunctionTestPrograms.PROCESS_ORDER_BY,
ProcessTableFunctionTestPrograms.PROCESS_MULTI_INPUT_ORDER_BY,
ProcessTableFunctionTestPrograms.PROCESS_ORDER_BY_TABLE_API,
ProcessTableFunctionTestPrograms.PROCESS_IMPLICIT_CASTS);
ProcessTableFunctionTestPrograms.PROCESS_IMPLICIT_CASTS,
ProcessTableFunctionTestPrograms.PROCESS_VARIANT,
ProcessTableFunctionTestPrograms.PROCESS_VARIANT_STATE);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -72,6 +72,8 @@
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.test.program.SinkTestStep;
import org.apache.flink.table.test.program.SourceTestStep;
import org.apache.flink.table.test.program.TableTestProgram;
Expand Down Expand Up @@ -928,6 +930,40 @@ 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, "

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

add a table arg with a variant column as well

+ "variant3 => PARSE_JSON('{\"a\":[1,\"b\"]}'))")
.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",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,7 @@
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;
Expand Down Expand Up @@ -549,6 +550,24 @@ public void eval(
}
}

/** 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 RequiredTimeFunction extends AppendProcessTableFunctionBase {
public void eval(@ArgumentHint({ArgumentTrait.ROW_SEMANTIC_TABLE, REQUIRE_ON_TIME}) Row r) {
Expand Down Expand Up @@ -1236,6 +1255,16 @@ public String toString() {
}
}

/** 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"),
Expand Down
Loading