Adds drain support to Spanner CDC source. - #39994
Open
acrites wants to merge 2 commits into
Open
Conversation
…ion on SDFs allows terminating immediately and changing to DROP TABLE IF EXISTS makes deletion idempotent.
Contributor
|
Checks are failing. Will not request review until checks are succeeding. If you'd like to override that behavior, comment |
Contributor
Author
|
The remaining failing PreCommit Spotless check says it's failing at |
Contributor
|
2026-09-04T08:36:38.6399796Z [ant:checkstyle] [ERROR] /runner/_work/beam/beam/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CdcRecord.java:20: You are using raw guava, please use vendored guava classes. [ForbidNonVendoredGuava] #40015 forward fix under way |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Before this change, we were calling the default implementation of truncateRestriction on drain, which, for bounded restrictions tries to finish evaluating the entire restriction. Both SDF's in the Spanner CDC source use bounded restrictions that are effectively unbounded. Overriding truncateRestriction to return null on these SDFs allows terminating immediately.
In addition, we had to change the cleanup action from DROP TABLE to DROP TABLE IF EXISTS to make deletion idempotent. Otherwise, if that work item gets retried by the backend, it will start fail-looping, preventing drain from completing.
We don't have a great way to test drain here, but I ran a test pipeline on Dataflow and verified that it can now drain successfully.
Thank you for your contribution! Follow this checklist to help us incorporate your contribution quickly and easily:
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
GitHub Actions Tests Status (on master branch)
See CI.md for more information about GitHub Actions CI or the workflows README to see a list of phrases to trigger workflows.