abrarsher23 opened a new issue, #18262:
URL: https://github.com/apache/iceberg/issues/18262

   **Apache Iceberg version:** 1.10.2 (reproduced). The relevant code — 
`DynamicWriter`'s writer map keyed
   on `WriteTarget` including `schemaId`, cleared only in `prepareCommit`, and 
`DynamicCommitter` building
   one `RowDelta` per (table, branch, checkpoint) — is **byte-identical in 
1.11.0, in 1.12.0-rc1 and on
   `main` (`flink/v2.3`)**, so the defect is expected to be present there too; 
only 1.10.2 has been run.
   
   **Query engine:** Flink (`DynamicIcebergSink`, DataStream API, upsert mode)
   
   **Please describe the bug 🐞**
   
   ### Summary
   
   In upsert mode, `DynamicIcebergSink` can emit an **equality delete** for a 
row it wrote **in the same
   commit**. Because all files a checkpoint produces for one table land in a 
single `RowDelta`, the delete and
   its target share a sequence number, and by spec an equality delete only 
applies to data with a *strictly
   lower* sequence number. The delete never binds and **both versions of the 
key remain visible**. A
   subsequent `RewriteDataFiles` merges the two data files and makes the 
duplicate permanent.
   
   The trigger is a schema evolution that happens *in the middle of a 
checkpoint*: the same table gets two
   `TaskWriter`s, and the second one starts with an empty `insertedRowMap`, so 
the upsert path's
   position-delete safety net is bypassed.
   
   ### Why this is a bug and not expected behaviour
   
   `FlinkSink`/`IcebergSink` maintain the invariant that a key written and 
re-written inside one checkpoint is
   retired with a **position** delete, never an equality delete, precisely 
because same-commit equality deletes
   cannot bind. `DynamicIcebergSink` inherits the delta-writer machinery but 
breaks the invariant, because it
   keys writers per record on a `WriteTarget` that includes the schema id.
   
   The sink itself documents the intended contract, in 
`DynamicCommitter#commitDeltaTxn`:
   
   > *"Position deletes committed to the table in this path are used only to 
delete rows from data files that
   > are being added in this commit."*
   
   That is true only while one `TaskWriter` per subtask covers the whole commit.
   
   ### Root cause, with pointers (line numbers at tag `apache-iceberg-1.10.2`, 
`flink/v2.0`; the same code is on `main` in `flink/v2.3`)
   
   1. `DynamicWriter#write` creates one writer per `WriteTarget`, per record —
      `DynamicWriter.java:92-100`:
   
      ```java
      writers.computeIfAbsent(
          new WriteTarget(
              element.tableName(), element.branch(), 
element.schema().schemaId(),
              element.spec().specId(), element.upsertMode(), 
element.equalityFields()),
      ```
   
      `WriteTarget` includes `schemaId` in `equals`/`hashCode` 
(`WriteTarget.java:32-37, 111-132`).
   
   2. `element.schema()` is the **resolved table schema** 
(`DynamicRecordProcessor.java:145-152, 155-173`).
      When an input schema returns `SCHEMA_UPDATE_NEEDED`, 
`TableUpdater#findOrCreateSchema` evolves the table
      inline and the table gains a new schema id (`TableUpdater.java:135-148`). 
With
      `immediateTableUpdate(true)` this happens in 
`DynamicRecordProcessor#collect`
      (`DynamicRecordProcessor.java:111-123`).
   
   3. The next record for the same table therefore has a different 
`WriteTarget`, so `computeIfAbsent` creates
      a **second `TaskWriter`**. `writers` is an unbounded `Maps.newHashMap()` 
(`DynamicWriter.java:84`) that is
      cleared only in `prepareCommit()` (`:193`), so both writers stay live for 
the rest of the checkpoint.
   
   4. Each `TaskWriter` owns a private `insertedRowMap`
      (`core .../io/BaseTaskWriter.java:121`, allocated `:143`). The upsert 
path chooses position vs equality
      delete solely from that map (`BaseTaskWriter.java:216-220`, `:186-196`; 
call site
      `BaseDeltaTaskWriter.java:82-88`). The second writer's map is empty ⇒ 
**equality** delete.
   
   5. `DynamicCommitter` groups every committable by `(tableName, branch)` and 
builds **one** `RowDelta` per
      checkpoint (`DynamicCommitter.java:125-134, 309-320, 327-328`) ⇒ one 
snapshot ⇒ one sequence number
      (`SnapshotProducer.java:257`) ⇒ the equality delete cannot apply to the 
first writer's data file.
   
   ### Minimal reproduction
   
   Reproduced with an integration test that feeds `DynamicIcebergSink` the 
following inputs (parallelism 1,
   `DistributionMode.NONE`, `immediateTableUpdate(true)`, 
`setUpsertMode(true)`, equality field `order_id`,
   checkpointing disabled so a bounded run produces exactly one commit).
   
   - schema **v1** = `(order_id string, amount int)`, equality key `order_id`
   - schema **v2** = v1 plus a **nullable** `note string` (this is what makes 
`CompareSchemasVisitor` return
     `SCHEMA_UPDATE_NEEDED`)
   - records, in order, in the **same** commit, both routed to the same table 
and branch:
     1. `INSERT` of `{order_id="K", amount=100}` under **v1**
     2. `UPDATE_AFTER` of `{order_id="K", amount=200, note=null}` under **v2**
   
   **Expected** (and what happens when both records are under v1): `count(*) 
WHERE order_id = 'K'` → **1**;
   the commit contains 1 data file, **1 position delete** (`content = 1`, the 
update retiring the insert) and
   1 equality delete (`content = 2`, the insert's normal `deleteKey` miss, 
which retires the key in *earlier*
   commits).
   
   **Actual:** `count(*)` → **2**. The commit contains **2 data files** (one 
per writer), **2 equality
   deletes** and **0 position deletes**, all at the same `sequence_number`; 
snapshot summary
   `added-data-files=2, added-equality-delete-files=2, 
total-position-deletes=0`. The update's delete was
   written as an equality delete by the second writer, so it cannot remove the 
first writer's row.
   `DynamicWriteResultAggregator` logs `Emitted 2 commit message to downstream 
committer operator`
   (`DynamicWriteResultAggregator.java:120`) instead of 1 — the cheapest 
external signal that a table was
   split across two writers in one checkpoint.
   
   Two neighbouring cases that should stay green and are useful as regression 
guards:
   - both records under v1 (no evolution) → 1 row, 1 position delete;
   - the evolution happening between **two different keys** (`c(K1)`@v1, 
`c(K2)`@v2) → 2 rows, 2 data files,
     2 equality deletes, 0 position deletes — correct. Note this means *"a data 
file and an equality delete
     at the same sequence number"* is by itself **not** the symptom; the 
symptom is an in-commit **update**
     whose delete is an equality delete instead of a position delete.
   
   ### Impact
   
   Silent duplicate rows on upsert tables whenever a key is updated twice 
inside the checkpoint in which the
   input schema evolves — i.e. exactly when an upstream CDC source adds a 
column. The damage is invisible to
   the sink (no exception, no failed commit), it is not self-healing, and 
`RewriteDataFiles` makes it
   permanent by merging both versions into one file and leaving the equality 
delete dangling.
   
   ### Proposed fix
   
   Give all `TaskWriter`s of one table a **shared inserted-row tracker** for 
the lifetime of a checkpoint, so
   the second writer's `deleteKey` finds the first writer's `PathOffset` and 
emits a **position** delete
   naming the first writer's data file — which is legal at an equal sequence 
number and is the mechanism
   `DynamicCommitter` already documents.
   
   Sharing scope: the `WriteTarget` **with `schemaId` removed** — `(tableName, 
branch, specId, upsertMode,
   equalityFields)`. The map's key type is only the equality-field projection
   (`TypeUtil.select(schema, equalityFieldIds)`, 
`BaseDeltaTaskWriter.java:63`), which is unchanged by adding
   a column, so the two writers' trackers are type-compatible. Lifetime = one 
checkpoint: clear it next to
   `writers.clear()` (`DynamicWriter.java:193`).
   
   Implementation notes:
   - `BaseEqualityDeltaWriter` currently **clears and nulls** `insertedRowMap` 
in `close()`
     (`BaseTaskWriter.java:243-246`); an externally-owned tracker must not be 
cleared there.
   - `PathOffset` is `private static` (`BaseTaskWriter.java:267`); either 
expose it or introduce a small
     `InsertedRowTracker` interface so the Flink module can own the map without 
widening the API.
   - Guard the sharing when an equality field's **type** changed between schema 
versions (allowed promotions),
     and keep `specId` in the sharing scope, since position deletes are 
partition-scoped. A mid-checkpoint
     **partition-spec** evolution would remain unfixed by this change; it 
should be called out.
   
   Happy to open a PR with the fix and a regression test for this case.
   


-- 
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