gustavodemorais commented on code in PR #29296:
URL: https://github.com/apache/flink/pull/29296#discussion_r4132902995
##########
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:
Yes, I'm aware of this one. I want to address that in another ticket/PR.
Added a comment with a TODO for
https://issues.apache.org/jira/browse/FLINK-40841
--
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]