[Python] Refactor MatchContinuously onto the Watch transform - #39461
[Python] Refactor MatchContinuously onto the Watch transform#39461Eliaaazzz wants to merge 3 commits 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
There was a problem hiding this comment.
You are right, there was a real gap: registry.get_coder does not convert typing or native generic annotations, so tuple[str, float] silently fell back to pickling, and the explicit coder here was masking exactly that. Watch now converts hints with convert_to_beam_type before the registry lookup, the explicit settings are gone, and the key functions infer the same StrUtf8Coder and TupleCoder as before, with a new watch_test case covering the native generic annotation. Also trimmed the docstrings this PR adds and updated the stale stacked PR note in the description. This round is one follow up commit on top of the reviewed one, commit 4a804a4.
There was a problem hiding this comment.
One more follow up commit, 4587416, from the same self review pass: the poll return type is annotated PollResult[FileMetadata], typing.Tuple key inference is covered next to the native form, the third party test imports are ordered, and the remaining Java references in test comments are gone.
registry.get_coder receives typing and native generic annotations such as tuple[str, float] unconverted and falls back to pickling. Watch now converts hints with convert_to_beam_type before the registry lookup, so MatchContinuously's annotated key functions infer the same StrUtf8Coder and TupleCoder the explicit settings supplied. Also trims the docstrings and comments this PR adds.
…ports Annotates _MatchContinuouslyPollFn with PollResult[FileMetadata], covers typing.Tuple key inference alongside the native form, reorders the third-party test imports, and drops the remaining Java references from test comments.
Routes
fileio.MatchContinuouslythrough theWatchtransform when deduplication is enabled.Built on the
Watchtransform, which merged in #39023.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