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]