abrarsher23 commented on PR #18288: URL: https://github.com/apache/iceberg/pull/18288#issuecomment-5918387868
Following up on my earlier comment. We needed this fix before an upstream release, so we implemented it in our fork: https://github.com/Affirm/iceberg/pull/17 (on 1.11.0, `flink/v2.0`). It's the same basic idea as this PR, so here's what we learned along the way in case it helps. Happy to help port any of it. **What we verified.** Beyond unit tests, we ran the real sink end to end against a live CDC stream. In one transaction a row is inserted, a column is added, and the same row is updated, all inside one checkpoint. Without the change the key was visible twice (two equality deletes, no position delete). With it, the second writer retired the first writer's row with a position delete and the key was visible once. **Suggestions:** 1. **Key-type comparison** (from my earlier comment, still applies). `PositionDeleteTrackerKey` uses `Types.StructType` equality, which includes field names, docs and defaults. A doc-only change on a key column, which the dynamic sink applies on its own, gives the second writer a new tracker and the duplicate comes back. Comparing field ID, type and optionality by position, which is what `StructLikeMap` compares, avoids that. In our PR that's `InsertedRowTracker.acceptsKeyType`. 2. **CDC without upsert** (also from before). Trackers are only created when `upsertMode()` is true, but `UPDATE_BEFORE` / `DELETE` go through `delete(row)` and the same inserted-row lookup, so a changelog written without upsert hits the same bug. Passing the tracker regardless of upsert mode fixes it, and single-writer behaviour doesn't change. 3. **One tracker per partition.** `BaseDeltaTaskWriter` hands the same tracker to every `RowDataDeltaWriter`, so in `PartitionedDeltaWriter` all partitions of a dynamic-sink table now share one map, even without a schema change. `writePosDelete` writes the delete under the *current* writer's `partitionKey`, so if a key were ever found through another partition's entry, the position delete would land in a different partition from the data file it names. With the partition fields in the upsert key that shouldn't happen today, but keeping one tracker per partition (the existing per-partition map, shared across writers) removes the dependency on it. 4. **Extra `loadTable`.** When a record carries no equality fields, the tracker path calls `catalog.loadTable` again for each new writer, although the factory was just built from the same table. Reusing the equality fields the factory already resolved saves a catalog round trip per writer. 5. **Checkstyle.** `build-checks` fails on `PositionDeleteTracker.PathOffset`: `path` and `rowOffset` must be private with accessors (`VisibilityModifier`). 6. **API shape.** +1 to @pvary's `create` / `track` / `untrack` suggestion. We ended up with something close: `InsertedRowTracker.create(keyType)`, a key-type check, and put/remove kept package-private. A writer only clears a tracker it created itself; the sink owns the shared ones and drops them at `prepareCommit` and `close`. 7. **Tests.** The current test covers the unpartitioned upsert + add-column case. Other cases we found worth pinning: a partitioned table; a `DELETE` through the second writer; alternating old/new schema versions within one checkpoint (newest value wins); a key-column doc change; a CDC update without upsert; trackers dropped at the checkpoint, so the next checkpoint falls back to an equality delete; and trackers scoped per table. A quick mutation check was useful too: giving each writer its own tracker again makes 7 of the 19 `TestDynamicWriter` tests in our PR fail. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
