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]

Reply via email to