Skip to content

[FLINK-40885][table] Support LATERAL SNAPSHOT join in the Table API - #29368

Open
fhueske wants to merge 3 commits into
apache:masterfrom
confluentinc:fhueske-FLINK-40885-Add-LSJ-to-TableAPI
Open

fhueske wants to merge 3 commits into
apache:masterfrom
confluentinc:fhueske-FLINK-40885-Add-LSJ-to-TableAPI

Conversation

@fhueske

@fhueske fhueske commented Oct 2, 2026 •

Copy link
Copy Markdown
Contributor

What is the purpose of the change

Enable the LATERAL SNAPSHOT processing-time temporal join to be expressed through the programmatic Table API (Table.joinLateral / Table.leftOuterJoinLateral with a call("SNAPSHOT", …) build side), not just SQL. The Table API surface produces the same optimized plan and operator (StreamExecLateralSnapshotJoin) as the equivalent SQL.

Brief changelog

  • QueryOperationConverter: convert a SNAPSHOT TABLE argument into a RexTableArgCall + scan input and emit a plain LogicalJoin (not a Correlate) so LogicalJoinToLateralSnapshotJoinRule matches, incl. INNER joins; deduplicate the table-argument conversion shared with the FROM-clause PTF path.
  • JoinOperationFactory: allow an ON predicate for a lateral SNAPSHOT join and require an equi-join predicate (reject always-true/empty ON with a clear error).
  • CorrelatedFunctionTableFactory: accept .as(...) aliasing on SNAPSHOT (a PROCESS_TABLE function).
  • Tests: Table API ↔ SQL plan-parity (INNER/LEFT, transformed build side, aliased columns, optional args, scrambled named-arg order, default compile-time, hidden METADATA VIRTUAL rowtime across column-expansion strategies); Table API validation negatives as unit tests; execution smokes; plus matching SQL coverage for the metadata-rowtime case.
  • Docs: extended Table API join docs

Adding docs is still to be done. Will update the PR later.

Verifying this change

  • LateralSnapshotJoinTableApiTest (streaming)— plan-parity + validation negatives (15 tests)
  • LateralSnapshotJoinTableApiTest (batch) — plan-parity with SQL
  • LateralSnapshotJoinSemanticTests — SQL + Table API end-to-end execution smokes (15 tests)
  • LateralSnapshotJoinTest (stream + batch) — SQL plan/golden + negative cases

Does this pull request potentially affect one of the following parts:

  • Dependencies (does it add or upgrade a dependency): no
  • The public API, i.e., is any changed class annotated with @Public(Evolving): yes
  • The serializers: no
  • The runtime per-record code paths (performance sensitive): no
  • Anything that affects deployment or recovery: JobManager (and its components), Checkpointing, Kubernetes/Yarn, ZooKeeper: no
  • The S3 file system connector: no

Documentation

  • Does this pull request introduce a new feature? yes
  • If yes, how is the feature documented? extended Table API join docs

Was generative AI tooling used to co-author this PR?
  • Yes (please specify the tool below)

Generated-by: Claude Code, Opus 4.8

@flinkbot

flinkbot commented Oct 2, 2026 •

Copy link
Copy Markdown
Collaborator

CI report:

Bot commands The @flinkbot bot supports the following commands:
  • @flinkbot run azure re-run the last Azure build

Enable Table.joinLateral / leftOuterJoinLateral with a call("SNAPSHOT", ...)
build side, producing the same plan as the equivalent SQL.

Release note: The LATERAL SNAPSHOT join can now be expressed through the
programmatic Table API, not just SQL, via Table.joinLateral /
Table.leftOuterJoinLateral with a call("SNAPSHOT", ...) build side.

Generated-by: Claude Code (Claude Opus 4.8)
…OT join

Verifies that the Table API surface produces the same optimized physical plan
as the equivalent SQL in batch mode, where the join degenerates to a regular join.

Generated-by: Claude Code (Claude Opus 4.8)
Add a Lateral Snapshot Join subsection to the Table API joins documentation
(English and Chinese) with a Java example; Scala and Python are marked as not
supported.

Generated-by: Claude Code (Claude Opus 4.8)
@airlock-confluentinc
airlock-confluentinc Bot force-pushed the fhueske-FLINK-40885-Add-LSJ-to-TableAPI branch from c7b880c to 7993432 Compare October 2, 2026 15:52
@fhueske

fhueske commented Oct 2, 2026

Copy link
Copy Markdown
Contributor Author

@flinkbot run azure

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?

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)

?

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants