[Python] Refactor MatchContinuously onto the Watch transform - #39461
[Python] Refactor MatchContinuously onto the Watch transform#39461Eliaaazzz wants to merge 1 commit into
Conversation
|
Caution The consumer version of Gemini Code Assist on GitHub has been sunset. All code review activity has officially ceased. |
|
Assigning reviewers: R: @tvalentyn for label python. Note: If you would like to opt out of this review, comment Available commands:
The PR bot will only process comments in the main thread (not review comments). |
|
R: @Abacn |
|
Stopping reviewer notifications for this pull request: review requested by someone other than the bot, ceding control. If you'd like to restart, comment |
Route MatchContinuously through Watch when deduplication is enabled. The polling loop and the set of already-matched file ids now live in the splittable DoFn restriction, replacing the per-key state DoFns. Because the matched ids are part of the restriction, a runner with checkpointing enabled restores them after a restart and does not reprocess files. The docstring is updated accordingly. Verified on Flink 1.20: after killing the TaskManager mid stream the job restored from a checkpoint and every file was still emitted exactly once. has_deduplication=False keeps the previous PeriodicImpulse behaviour.
1361ca1 to
3760698
Compare
| termination = _WatchWindowTermination(clock, start_ts.micros, max_polls) | ||
| if self.match_upd: | ||
| output_key_fn = _file_path_and_mtime_key | ||
| output_key_coder = TupleCoder([StrUtf8Coder(), FloatCoder()]) |
There was a problem hiding this comment.
Please clean up the PR against the merged version of watch transform. In particular, please check #39023 (comment)
These explicit coder settings shouldn't be needed, otherwise it suggests the coder inference in #39023 still has some gaps
Routes
fileio.MatchContinuouslythrough theWatchtransform when deduplication is enabled.Stacked on #39023. That PR adds
Watchitself and is still open, so its three commits currently show up here too. Only the last commit, "[Python] Refactor MatchContinuously onto the Watch transform", is new in this PR. Once #39023 merges this reduces to a single commit.What changes
The polling loop and the set of already-matched file ids move into the splittable DoFn restriction. The per-key state DoFns
_RemoveDuplicatesand_RemoveOldDuplicatesare removed, sinceWatchperforms the deduplication.has_deduplication=Falsekeeps the previousPeriodicImpulsebehaviour, so that path is unchanged.Behaviour change worth calling out
Because the matched ids are part of the restriction, a runner with checkpointing enabled restores them after a restart and does not reprocess files. The class docstring previously stated the opposite, that already processed files are reprocessed on restart, which was accurate for the earlier memory-only implementation. The docstring is updated in this PR.
Validation
Fault tolerance on Flink 1.20 with checkpointing enabled: two files present at start, two added while running, then the TaskManager was killed mid stream. The JobManager restored the job from checkpoint 3 and every file was still emitted exactly once, with no reprocessing.
Also exercised on Dataflow Runner v2 reading a real GCS prefix: files present at startup and files added to the bucket mid run were each emitted exactly once.
Unit tests:
watch_test.py28 passed,fileio_test.pyMatchContinuously tests 8 passed. Formatted with yapf 0.43.0 and isort 7.0.0.Thank you for your contribution! Follow this checklist to help us incorporate your contribution quickly and easily:
R: @username).addresses #123), if applicable. This will automatically add a link to the pull request in the issue. If you would like the issue to automatically close on merging the pull request, commentfixes #<ISSUE NUMBER>instead.CHANGES.mdwith noteworthy changes.See the Contributor Guide for more tips on how to make review process smoother.
To check the build health, please visit https://github.com/apache/beam/blob/master/.test-infra/BUILD_STATUS.md