Skip to content

[Python] Refactor MatchContinuously onto the Watch transform - #39461

Open
Eliaaazzz wants to merge 1 commit into
apache:masterfrom
Eliaaazzz:matchcontinuously-on-watch
Open

[Python] Refactor MatchContinuously onto the Watch transform#39461
Eliaaazzz wants to merge 1 commit into
apache:masterfrom
Eliaaazzz:matchcontinuously-on-watch

Conversation

@Eliaaazzz

Copy link
Copy Markdown
Contributor

Routes fileio.MatchContinuously through the Watch transform when deduplication is enabled.

Stacked on #39023. That PR adds Watch itself 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 _RemoveDuplicates and _RemoveOldDuplicates are removed, since Watch performs the deduplication.

has_deduplication=False keeps the previous PeriodicImpulse behaviour, 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.

Completed checkpoint 3 for job 15502720... (56814 bytes)
Job beam-watch-matchcontinuously switched from state RUNNING to RESTARTING
Job beam-watch-matchcontinuously switched from state RESTARTING to RUNNING
Restoring job 15502720... from Checkpoint 3

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.py 28 passed, fileio_test.py MatchContinuously 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:

  • Choose reviewer(s) and mention them in a comment (R: @username).
  • Mention the appropriate issue in your description (for example: 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, comment fixes #<ISSUE NUMBER> instead.
  • Update CHANGES.md with noteworthy changes.
  • If this contribution is large, please file an Apache Individual Contributor License Agreement.

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

@gemini-code-assist

Copy link
Copy Markdown
Contributor

Caution

The consumer version of Gemini Code Assist on GitHub has been sunset. All code review activity has officially ceased.

@github-actions

Copy link
Copy Markdown
Contributor

Assigning reviewers:

R: @tvalentyn for label python.

Note: If you would like to opt out of this review, comment assign to next reviewer.

Available commands:

  • stop reviewer notifications - opt out of the automated review tooling
  • remind me after tests pass - tag the comment author after tests pass
  • waiting on author - shift the attention set back to the author (any comment or push by the author will return the attention set to the reviewers)

The PR bot will only process comments in the main thread (not review comments).

@tvalentyn

Copy link
Copy Markdown
Contributor

R: @Abacn

@github-actions

Copy link
Copy Markdown
Contributor

Stopping reviewer notifications for this pull request: review requested by someone other than the bot, ceding control. If you'd like to restart, comment assign set of reviewers

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.
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants