[IcebergIO] CDC record and sort key - #39981
Conversation
Codecov Report✅ All modified and coverable lines are covered by tests. Additional details and impacted files@@ Coverage Diff @@
## master #39981 +/- ##
============================================
+ Coverage 55.95% 58.33% +2.37%
- Complexity 2272 13542 +11270
============================================
Files 1112 2566 +1454
Lines 175063 267300 +92237
Branches 1458 11054 +9596
============================================
+ Hits 97959 155926 +57967
- Misses 74645 105460 +30815
- Partials 2459 5914 +3455
Flags with carried forward coverage won't be shown. Click here to find out more. ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
|
Assigning reviewers: R: @kennknowles for label java. 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). |
| /** | ||
| * One change record carried through the CDC sink's shuffle. | ||
| * | ||
| * <p>{@link ValueKind} is reified because it's not preserved across a {@code GroupByKey}. |
There was a problem hiding this comment.
Also, once you are "inside" a sink, you don't (necessarily) need the implicit propagation of ValueKind as metadata, since it is more like an explicit field you will write to the sink.
| * level and cannot infer a dynamic {@link Row} field, so {@code @DefaultSchema} alone cannot | ||
| * produce a working coder for {@link CdcRecord}. | ||
| */ | ||
| final class CdcRecordCoder extends CustomCoder<CdcRecord> { |
There was a problem hiding this comment.
Is there any world in which Python will want to directly send these CdcRecord things over the wire? Just curious. You would want a StructuredCoder<CdcRecord> not a CustomCoder which I think translates to beam:coder:javasdk with java serialized payload, aka not intelligible to runners or other languages.
There was a problem hiding this comment.
CdcRecord is something internal to the sink. The user supplies Beam Rows (whether it's Java or Python) and they they get resolved and converted to CdcRecords
Part of #39979
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.