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
41 changes: 41 additions & 0 deletions docs/content.zh/docs/dev/table/tableApi.md
Original file line number Diff line number Diff line change
Expand Up @@ -1152,6 +1152,47 @@ val result = orders
{{< /tab >}}
{{< /tabs >}}

#### Lateral Snapshot Join

{{< label "Batch" >}} {{< label "Streaming" >}}

Enriches an append-only table with the current state of an updating table by calling the built-in `SNAPSHOT` [process table function]({{< ref "docs/sql/functions/built-in-functions" >}}) as the build (right) side of a `LATERAL` join. Each probe-side row is joined with the build-side state that is current when the row is processed. Both inner and left outer joins are supported, and the predicate must contain at least one equality condition.

See [LATERAL SNAPSHOT join]({{< ref "docs/sql/reference/queries/joins" >}}#lateral-snapshot-join) for the full list of `SNAPSHOT` arguments and the join semantics.

{{< tabs "snapshotjoin" >}}
{{< tab "Java" >}}
```java
Table orders = tableEnv.from("Orders");
Table currencyRates = tableEnv.from("CurrencyRates");

// inner join: enrich every order with the current conversion rate
Table result = orders
.joinLateral(
call(
"SNAPSHOT",
currencyRates.asArgument("input"),
descriptor("update_time").asArgument("on_time")),
$("currency").isEqual($("r_currency")));

// left outer join: unmatched orders are preserved and padded with nulls
Table leftResult = orders
.leftOuterJoinLateral(
call(
"SNAPSHOT",
currencyRates.asArgument("input"),
descriptor("update_time").asArgument("on_time")),
$("currency").isEqual($("r_currency")));
```
{{< /tab >}}
{{< tab "Scala" >}}
Currently not supported in Scala Table API.
{{< /tab >}}
{{< tab "Python" >}}
Currently not supported in Python Table API.
{{< /tab >}}
{{< /tabs >}}

{{< top >}}

### Set 操作
Expand Down
41 changes: 41 additions & 0 deletions docs/content/docs/dev/table/tableApi.md
Original file line number Diff line number Diff line change
Expand Up @@ -1151,6 +1151,47 @@ Currently not supported in Python Table API.
{{< /tab >}}
{{< /tabs >}}

#### Lateral Snapshot Join

{{< label "Batch" >}} {{< label "Streaming" >}}

Enriches an append-only table with the current state of an updating table by calling the built-in `SNAPSHOT` [process table function]({{< ref "docs/sql/functions/built-in-functions" >}}) as the build (right) side of a `LATERAL` join. Each probe-side row is joined with the build-side state that is current when the row is processed. Both inner and left outer joins are supported, and the predicate must contain at least one equality condition.

See [LATERAL SNAPSHOT join]({{< ref "docs/sql/reference/queries/joins" >}}#lateral-snapshot-join) for the full list of `SNAPSHOT` arguments and the join semantics.

{{< tabs "snapshotjoin" >}}
{{< tab "Java" >}}
```java
Table orders = tableEnv.from("Orders");
Table currencyRates = tableEnv.from("CurrencyRates");

// inner join: enrich every order with the current conversion rate
Table result = orders
.joinLateral(
call(
"SNAPSHOT",
currencyRates.asArgument("input"),
descriptor("update_time").asArgument("on_time")),
$("currency").isEqual($("r_currency")));

// left outer join: unmatched orders are preserved and padded with nulls
Table leftResult = orders
.leftOuterJoinLateral(
call(
"SNAPSHOT",
currencyRates.asArgument("input"),
descriptor("update_time").asArgument("on_time")),
$("currency").isEqual($("r_currency")));
```
{{< /tab >}}
{{< tab "Scala" >}}
Currently not supported in Scala Table API.
{{< /tab >}}
{{< tab "Python" >}}
Currently not supported in Python Table API.
{{< /tab >}}
{{< /tabs >}}

{{< top >}}

### Set Operations
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@
import org.apache.flink.table.expressions.ExpressionUtils;
import org.apache.flink.table.expressions.ResolvedExpression;
import org.apache.flink.table.expressions.utils.ResolvedExpressionDefaultVisitor;
import org.apache.flink.table.functions.BuiltInFunctionDefinitions;
import org.apache.flink.table.functions.FunctionDefinition;
import org.apache.flink.table.functions.FunctionKind;
import org.apache.flink.table.operations.CorrelatedFunctionQueryOperation;
Expand Down Expand Up @@ -58,6 +59,14 @@ QueryOperation create(ResolvedExpression callExpr, List<String> leftTableFieldNa
return callExpr.accept(calculatedTableCreator);
}

/**
* Whether the given function is {@code SNAPSHOT}, currently the only {@link
* FunctionKind#PROCESS_TABLE} function that is valid as the build side of a lateral join.
*/
static boolean isSnapshot(FunctionDefinition definition) {
return definition.equals(BuiltInFunctionDefinitions.SNAPSHOT);
}

private static class FunctionTableCallVisitor
extends ResolvedExpressionDefaultVisitor<CorrelatedFunctionQueryOperation> {
private final List<String> leftTableFieldNames;
Expand Down Expand Up @@ -91,7 +100,11 @@ private CorrelatedFunctionQueryOperation unwrapFromAlias(CallExpression call) {
+ alias)))
.collect(toList());

if (!isFunctionOfKind(children.get(0), FunctionKind.TABLE)) {
// SNAPSHOT is PROCESS_TABLE rather than TABLE but is valid in a LATERAL join.
Expression firstChild = children.get(0);
if (!isFunctionOfKind(children.get(0), FunctionKind.TABLE)

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.

Suggested change
if (!isFunctionOfKind(children.get(0), FunctionKind.TABLE)
if (!isFunctionOfKind(firstChild, FunctionKind.TABLE)

?

&& !(firstChild instanceof CallExpression
&& isSnapshot(((CallExpression) firstChild).getFunctionDefinition()))) {
throw fail();
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -80,15 +80,29 @@ private void validateCondition(
JoinType joinType,
ResolvedExpression condition,
boolean correlated) {
boolean alwaysTrue = ExpressionUtils.extractValue(condition, Boolean.class).orElse(false);

if (alwaysTrue) {
final boolean alwaysTrue =
ExpressionUtils.extractValue(condition, Boolean.class).orElse(false);

// A lateral SNAPSHOT join is rewritten into a dedicated join that supports an ON predicate
// but requires at least one equi-join predicate (for both INNER and LEFT OUTER). An
// always-true/empty condition is therefore not allowed; fall through to the equi-join check
// below so the user gets a clear error instead of a later planner failure.
final boolean isLateralSnapshot =
correlated
&& right instanceof CorrelatedFunctionQueryOperation
&& CorrelatedFunctionTableFactory.isSnapshot(
((CorrelatedFunctionQueryOperation) right)
.getResolvedFunction()
.getDefinition());

if (alwaysTrue && !isLateralSnapshot) {
return;
}

Boolean equiJoinExists = condition.accept(equiJoinExistsChecker);
if (correlated
&& right instanceof CorrelatedFunctionQueryOperation
&& !isLateralSnapshot
&& joinType != JoinType.INNER) {
throw new ValidationException(
"Predicate for lateral left outer join with table function can only be empty or literal true.");
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -284,7 +284,17 @@ private AggCall getAggCall(Expression aggregateExpression) {
@Override
public RelNode visit(JoinQueryOperation join) {
final Set<CorrelationId> corSet;
if (join.isCorrelated()) {
// Check if this is a LATERAL SNAPSHOT join. The SNAPSHOT function never references the
// outer row. Emit a plain LogicalJoin instead of a Correlate, so that
// LogicalJoinToLateralSnapshotJoinRule matches.
final QueryOperation right = join.getChildren().get(1);
final boolean isLateralSnapshotJoin =
right instanceof CorrelatedFunctionQueryOperation
&& isSnapshot(
((CorrelatedFunctionQueryOperation) right)
.getResolvedFunction()
.getDefinition());
if (join.isCorrelated() && !isLateralSnapshotJoin) {
corSet = Collections.singleton(relBuilder.peek().getCluster().createCorrel());
} else {
corSet = Collections.emptySet();
Expand Down Expand Up @@ -352,44 +362,11 @@ public RelNode visit(FunctionQueryOperation functionTable) {
.map(
resolvedArg -> {
if (resolvedArg instanceof TableReferenceExpression) {
final TableReferenceExpression tableRef =
(TableReferenceExpression) resolvedArg;
final LogicalType tableArgType =
tableRef.getOutputDataType().getLogicalType();
final RelDataType rowType =
typeFactory.buildRelNodeRowType(
(RowType) tableArgType);
final int[] partitionKeys;
final int[] orderKeys;
final SortOrder[] sortOrders;
if (tableRef.getQueryOperation()
instanceof PartitionQueryOperation) {
final PartitionQueryOperation partitionOperation =
(PartitionQueryOperation)
tableRef.getQueryOperation();
partitionKeys =
partitionOperation.getPartitionKeys();
orderKeys = partitionOperation.getOrderKeys();
final SortDirection[] directions =
partitionOperation.getOrderDirections();
sortOrders = new SortOrder[directions.length];
for (int i = 0; i < directions.length; i++) {
sortOrders[i] =
SortOrder.fromSortDirection(
directions[i]);
}
} else {
partitionKeys = new int[0];
orderKeys = new int[0];
sortOrders = new SortOrder[0];
}
final RexTableArgCall tableArgCall =
new RexTableArgCall(
rowType,
buildTableArgCall(
(TableReferenceExpression) resolvedArg,
inputStack.size(),
partitionKeys,
orderKeys,
sortOrders);
typeFactory);
inputStack.add(relBuilder.build());
return tableArgCall;
}
Expand All @@ -400,22 +377,9 @@ public RelNode visit(FunctionQueryOperation functionTable) {
// relBuilder.build() works in LIFO fashion, this restores the original input order
Collections.reverse(inputStack);

final BridgingSqlFunction sqlFunction =
BridgingSqlFunction.of(relBuilder.getCluster(), contextFunction);

final RexNode call =
relBuilder
.getRexBuilder()
.makeCall(outputRelDataType, sqlFunction, rexNodeArgs);
final RelNode functionScan =
LogicalTableFunctionScan.create(
relBuilder.getCluster(),
inputStack,
call,
null,
outputRelDataType,
Collections.emptySet());
relBuilder.push(functionScan);
relBuilder.push(
createTableFunctionScan(
contextFunction, rexNodeArgs, inputStack, outputRelDataType));
return relBuilder.build();
}

Expand All @@ -430,6 +394,11 @@ public RelNode visit(PartitionQueryOperation partition) {
public RelNode visit(CorrelatedFunctionQueryOperation correlatedFunction) {
final ContextResolvedFunction contextFunction =
correlatedFunction.getResolvedFunction();
// SNAPSHOT is currently the only PTF valid in a lateral join and hence the only
// function accepting table args.
if (isSnapshot(contextFunction.getDefinition())) {
return convertCorrelatedFunctionWithTableArgs(correlatedFunction);
}
final List<RexNode> parameters = convertToRexNodes(correlatedFunction.getArguments());

final FunctionDefinition functionDefinition = contextFunction.getDefinition();
Expand All @@ -455,6 +424,87 @@ public RelNode visit(CorrelatedFunctionQueryOperation correlatedFunction) {
return relBuilder.build();
}

private RelNode convertCorrelatedFunctionWithTableArgs(
CorrelatedFunctionQueryOperation correlatedFunction) {
final FlinkTypeFactory typeFactory = ShortcutUtils.unwrapTypeFactory(relBuilder);
final List<RelNode> inputs = new ArrayList<>();
final List<RexNode> args = new ArrayList<>();
for (ResolvedExpression arg : correlatedFunction.getArguments()) {
if (arg instanceof TableReferenceExpression) {
final TableReferenceExpression tableRef = (TableReferenceExpression) arg;
args.add(buildTableArgCall(tableRef, inputs.size(), typeFactory));
inputs.add(tableRef.getQueryOperation().accept(QueryOperationConverter.this));
} else {
args.add(convertExprToRexNode(arg));
}
}
// Derive the output row type like the FROM-clause PTF path
// (FunctionQueryOperation#getOutputDataType).
final RelDataType outputType =
typeFactory.buildRelNodeRowType(
(RowType)
DataTypeUtils.fromResolvedSchemaPreservingTimeAttributes(
correlatedFunction.getResolvedSchema())
.getLogicalType());
return createTableFunctionScan(
correlatedFunction.getResolvedFunction(), args, inputs, outputType);
}

private boolean isSnapshot(FunctionDefinition definition) {

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.

there is existing LateralSnapshotJoinUtil.isSnapshotFunction
is there a reason we don't reuse it?

// SNAPSHOT is currently the only PTF valid as the build side of a lateral join; the
// lateral-join-specific handling in this converter keys off this single predicate.
return definition.equals(BuiltInFunctionDefinitions.SNAPSHOT);
}

/**
* Converts a {@link TableReferenceExpression} table argument into a {@link
* RexTableArgCall}, carrying PARTITION BY / ORDER BY keys when the referenced operation is
* a {@link PartitionQueryOperation}.
*/
private RexTableArgCall buildTableArgCall(
TableReferenceExpression tableRef, int inputIndex, FlinkTypeFactory typeFactory) {
final RelDataType rowType =
typeFactory.buildRelNodeRowType(
(RowType) tableRef.getOutputDataType().getLogicalType());
final int[] partitionKeys;
final int[] orderKeys;
final SortOrder[] sortOrders;
if (tableRef.getQueryOperation() instanceof PartitionQueryOperation) {
final PartitionQueryOperation partitionOperation =
(PartitionQueryOperation) tableRef.getQueryOperation();
partitionKeys = partitionOperation.getPartitionKeys();
orderKeys = partitionOperation.getOrderKeys();
final SortDirection[] directions = partitionOperation.getOrderDirections();
sortOrders = new SortOrder[directions.length];
for (int i = 0; i < directions.length; i++) {
sortOrders[i] = SortOrder.fromSortDirection(directions[i]);
}
} else {
partitionKeys = new int[0];
orderKeys = new int[0];
sortOrders = new SortOrder[0];
}
return new RexTableArgCall(rowType, inputIndex, partitionKeys, orderKeys, sortOrders);
}

/** Builds a {@link LogicalTableFunctionScan} for a function call with table arguments. */
private RelNode createTableFunctionScan(
ContextResolvedFunction function,
List<RexNode> args,
List<RelNode> inputs,
RelDataType outputType) {
final BridgingSqlFunction sqlFunction =
BridgingSqlFunction.of(relBuilder.getCluster(), function);
final RexNode call = relBuilder.getRexBuilder().makeCall(outputType, sqlFunction, args);
return LogicalTableFunctionScan.create(
relBuilder.getCluster(),
inputs,
call,
null,
outputType,
Collections.emptySet());
}

private RelNode convertLegacyTableFunction(
CorrelatedFunctionQueryOperation correlatedFunction,
TableFunctionDefinition functionDefinition,
Expand Down
Loading