Conversation
7681a1c to
2bd8780
Compare
1f3861a to
c7a4fe2
Compare
|
@JingsongLi Can I get a review on this? |
|
Thanks for the feature. I reviewed the write and read paths end to end and found two blocking issues plus a consistency concern. 1. [P1] The option is applied for every changelog producer, but only the
|
c7a4fe2 to
5ec5070
Compare
@JingsongLi Let me explain this in detail below. Yeah I agree the changes are more complex that's why I initially went with minimal change approach. Use case and motivationPaimon’s lookup changelog producer reconstructs UPDATE_BEFORE and DELETE records by looking up the previous table state. As a result, regular columns in those records describe the old row, not necessarily the event that caused the change. For example: UPDATE_BEFORE: event_ts = 50, __internal__event_ts = 100 Downstream CDC consumers often need the event timestamp or source sequence to apply changes deterministically. Concurrent Sinks to a Single Cassandra tableMultiple sources can fan into the same Cassandra table while being processed by independent jobs. A common example is migrating from Kafka to Paimon: Kafka retains the original event payload for retract events, whereas Paimon’s lookup-generated changelog may represent UPDATE_BEFORE and DELETE records using the previous row image. This makes it difficult to replace a Kafka-based job with a Paimon-based job without changing downstream conflict-resolution behavior. Exposing selected fields as changelog metadata preserves the event-level values and gives the sink a consistent, globally comparable event version for conditional writes and tombstones, preventing stale updates from overwriting newer changes or resurrecting deleted records. This behavior allows creating Sinks from multiple Paimon tables merging into a single Cassandra column. |
JingsongLi
left a comment
There was a problem hiding this comment.
The Cassandra WRITETIME use case gives this feature concrete end-to-end value. The earlier lookup-only validation, physical before-image preservation, and Spark write-side filtering are present on this head. I found one remaining schema issue before this should ship:
P2 — Allocate synthetic metadata field IDs above nested field IDs (ChangelogEventMetadata.java:117). extraValueFields takes the maximum top-level field ID. A valid table can have id (0), payload ROW<nested INT> (top-level 1, nested field ID 3), and event_ts (2). Exposing event_ts assigns the new __internal__event_ts field ID 3, duplicating payload.nested. I added this case to ChangelogEventMetadataTest in an isolated checkout: appendMetadataFields(...).collectFieldIds(...) fails with Broken schema, field id 3 is duplicated. The generated row type is used by the changelog writer and both read integrations, so its IDs need to remain valid for nested schemas and field-ID based mapping. Please derive the next ID from the full recursive field-ID set (for example, RowType.currentHighestFieldId(valueType.getFields()) + 1) and test a nested table through a changelog read.
Validation: the added core tests passed 16/16 and Spark 3 PaimonCDCSourceTest passed 6/6 on 5ec5070; the strengthened nested-ID test fails as above. I attempted the 3 Flink integration tests locally, but all stopped on the local Avro classpath mismatch DataFileWriter.setEncoder before exercising this change. The PR's Flink and other CI jobs are green on this head.
5ec5070 to
b0fd1ea
Compare
Thanks a lot for pointing it out. I have made that fix. |
|
Re-reviewed head Local JDK 8 core suites passed: |
|
I took another look at head 1. [P1] Event metadata breaks retract semantics when used by normal SQL operators.
2. [P1] Changing The prefix is not marked immutable, so it can be changed after changelog files exist. The write path stored the old prefixed field name, but There is also a source/sink mismatch in the documentation example: Finally, the option names obscure the contract. |
e201136 to
c287adb
Compare
|
@JingsongLi Thanks for the review.
|
Preserve correct before-images in changelog retraction records and append event values as configurable metadata columns, accessible via Flink SupportsReadingMetadata and Spark table schema.
c287adb to
2b5f23f
Compare
Preserve correct before-images in changelog retraction records and append event values as configurable metadata columns, accessible via Flink SupportsReadingMetadata and Spark table schema.
Purpose
<prefix><field>metadata columnssuch as Cassandra
WRITETIME.Tests
paimon-docs verifypasses on Java 8 and Java 11.