gustavodemorais opened a new pull request, #29339: URL: https://github.com/apache/flink/pull/29339
To be rebased after https://github.com/apache/flink/pull/29296 is merged. ## What is the purpose of the change With upsert inputs, a record that replaces a stored record with the same unique key may match different rows than the replaced one if the join condition has a non-equi part. The streaming join never retracted the joined rows and null padding that only the replaced record produced, which led to stale rows and wrong null padding. ## Brief change log - Join cases with non-equi join keys also work now The summary of the price to pay for this correctness 1. **Extra state read.** One lookup per accumulate. - Equi: has-unique-key side (upsert key contains the join key), other side outer. - Non-equi: any side with unique key. - Retract/append pay it too (returns null). Extra cost for a case that already works. - Can drop for retract/append: planner flag "input is upsert". Needs optional plan field. 2. **TTL refresh.** Upsert on R that keeps same match no longer rewrites L count -> L TTL not refreshed -> L may expire earlier. - It'd possible to change this and always write count (`count + (isNewMatch ? 1 : 0)`). This would be more write for this path. ## Verifying this change - StreamingJoinOperatorTest - StreamingMiniBatchJoinOperatorTest ## 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)`: no - The serializers: no - The runtime per-record code paths (performance sensitive): yes - one extra state lookup per accumulate record for joins with a non-equi condition - 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? no - If yes, how is the feature documented? not applicable --- ##### Was generative AI tooling used to co-author this PR? - [x] Yes (please specify the tool below) 2.1.282 (Claude Code) with Opus 5.5 -- 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]
