MartijnVisser commented on code in PR #29296:
URL: https://github.com/apache/flink/pull/29296#discussion_r4131605514
##########
flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/operators/join/stream/StreamingJoinOperator.java:
##########
@@ -332,6 +338,28 @@ protected void processElement(
}
}
+ /**
+ * Returns whether the record adds a new match for other-side records,
which indicates that we
+ * have to increase the number of associations for a specific key match.
+ *
+ * <p>Matches are only counted when the other side is outer, so the result
is false otherwise
+ * and no lookup happens. If the join key contains the unique key, there
is at most one record
+ * per join key, so a match is never additional and no lookup is needed
either.
+ *
+ * <p>A suppressed retraction in mini-batch mode removes the matches of a
record but keeps it in
+ * state, so the suppressed accumulate message that follows is always an
additional match.
+ */
+ private boolean isAdditionalMatch(
+ RowData record, JoinRecordStateView stateView, boolean isLeft,
boolean isSuppress)
+ throws Exception {
+ final boolean otherIsOuter = isLeft ? rightIsOuter : leftIsOuter;
+ final JoinInputSideSpec inputSideSpec = isLeft ? leftInputSideSpec :
rightInputSideSpec;
+ return otherIsOuter
+ && (isSuppress
+ || (!inputSideSpec.joinKeyContainsUniqueKey()
+ && !stateView.hasRecord(record)));
Review Comment:
For a non-equi condition this assumes the old version also matched `other`,
so an update from non-match to match is no longer counted. Fine for the
non-equi subtask, but please note the assumption here.
##########
flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/runtime/stream/sql/JoinITCase.scala:
##########
@@ -809,6 +816,81 @@ class JoinITCase(miniBatch: MiniBatchMode, state:
StateBackendMode, enableAsyncS
assertThat(sink.getRetractResults.sorted).isEqualTo(expected.sorted)
}
+ /** Only a line with a new unique key adds a match for the order, an updated
line does not. */
+ @TestTemplate
+ def testLeftJoinOnRowKey(): Unit = {
+ // Disable test for miniBatch and asyncState
Review Comment:
Perhaps we should say why in there, since mini-batch folds the `+U`/`-D`
pair within a bundle and async state isn't fixed yet.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]