Background
This is a derived issue from the DELETE-via-PK epic (#63). It captures a correctness bug that
@LuciferYang flagged during review of #63 and asked to be raised separately rather than mixed into
the epic thread. This document is the standalone tracking ticket for that bug.
Positioning: this is a pre-registered design risk for #63 Phase B/C, not a defect in the
current main branch. Today LanceDynamicTableSink#getChangelogMode() declares insert-only and
LanceSink has no DELETE / RowKind-bucketing logic, so the bug cannot manifest yet. It becomes
live only once Phase B (RowKind bucketing) and Phase C (delete-then-insert flush) land. Raised
now, ahead of implementation, so the Phase B/C design accounts for it.
The epic's Phase B and Phase C introduce two independent decisions that, when composed, produce a
wrong final dataset state:
- Phase B buckets the buffered
RowData by RowKind: DELETE rows go into a delete bucket,
INSERT/UPDATE_AFTER rows go into an insert bucket.
- Phase C fixes the checkpoint flush order as deletes first, then inserts.
The bug
Consider a single checkpoint window in which both messages below are buffered for the same
primary key k:
+I(k) -- insert row with PK k
-D(k) -- delete row with PK k
In a real CDC stream these two have a meaningful temporal order:
| Real upstream order |
Correct final state of k |
+I(k) then -D(k) |
k must not exist (inserted, then deleted) |
-D(k) then +I(k) |
k must exist (deleted, then re-inserted) |
But because the sink flushes deletes first, then inserts regardless of the actual per-key
order, both real orders are collapsed into the same sequence:
delete(k) — remove k from Lance (even though it may not exist yet).
insert(k) — append k to Lance.
Final state in both cases: k exists. For the +I(k)-then--D(k) case this is wrong — the
delete was effectively dropped, leaving a row that upstream already removed. This violates the
exactly-once / last-write-wins contract that a CDC sink must honor.
Root cause (abstraction level)
The ordering bug is a symptom, not the root cause. The real defect is that bucketing by
RowKind destroys the temporal ordering between multiple events that share the same key. Once
+I(k) lands in the insert bucket and -D(k) lands in the delete bucket, the "who came first"
information is irretrievably lost from the bucket structure. The fixed delete-before-insert order
merely makes that lost information visible as a wrong final state.
This is why "just flip the order" is not a fix: insert-before-delete would instead break the
-D(k)-then-+I(k) case. Any correct fix must either preserve per-key order (don't bucket by
RowKind; process keyed and in arrival order) or collapse to a final state before bucketing
(reduce per key first).
Scope is broader than the +I/-D example
+I(k) / -D(k) is the minimal illustration used in the epic review, but the same defect recurs
for any same-key multi-event window. Given epic Phase B's handling of the four RowKinds, the
affected combinations are:
| Same-key events within one checkpoint |
Expected final state of k |
Bucket + delete-first result |
+I(k) then -D(k) |
absent |
present ❌ |
-D(k) then +I(k) |
present |
present ✅ (accidentally correct) |
+U(k) then -D(k) |
absent |
present ❌ |
+U(k1) then +U(k2) |
only k2 |
both k1 and k2 present ❌ |
(+U = UPDATE_AFTER, which epic Phase B expands to delete-by-PK + append; two consecutive
+Us therefore produce two inserts unless collapsed.)
Why "deletes first, then inserts" exists in the first place
The fixed order is safe only when delete and insert targets are disjoint keys (the common case:
an UPDATE_AFTER on key k is implemented as delete(k) + insert(k')). There, ordering matters
for a different reason: you must delete before re-inserting so the delete predicate does not also
match the freshly inserted row. The Phase C order was chosen for that disjoint-key reason, but it is
not sufficient for the same-key +I/-D collision.
Scope / affected code (once Phase A/B/C land)
LanceSink (or its successor): the RowKind-bucketing + flushDeletes-then-flushInserts
checkpoint logic.
LanceDynamicTableSink#getChangelogMode: declares the full +I/-U/+U/-D set when a PK is
present; this bug only manifests once DELETE is actually accepted.
Candidate fixes (to be evaluated during implementation)
-
Per-key last-write-wins collapse (the reduction logic). Before flushing, reduce the buffered
rows per PK to a single final action derived from the last RowKind event seen for that key
within the checkpoint. A same-key +I/-D pair collapses to delete-only or insert-only,
eliminating the ordering ambiguity. This matches how Paimon/Iceberg CDC sinks treat keyed
changelog.
⚠️ Precondition — this does NOT come for free. Collapse only works if every event for a given
PK is visible to the same collapsing buffer. The current LanceSink extends RichSinkFunction
(non-keyed) and SinkFunctionProvider.of(...) (no keyBy) provide no such guarantee — a
+I(k) and -D(k) pair can be partitioned onto different subtasks, each seeing only half the
pair. Therefore option 1 depends on option 2's keyed guarantee (or an equivalent explicit
partitioner / key-grouping). It is not "topology-neutral" as the original wording implied.
-
Keyed ordering guarantee (the enabling prerequisite). Use Flink's keyed state / keyed sink
(KeyedProcessFunction) so all events for a given PK are processed by the same subtask in
arrival order, then apply them in arrival order instead of a global delete-before-insert order.
This is the structural fix; option 1 is the reduction logic layered on top of it.
-
Reject mixed same-key operations within a checkpoint with a clear error. Pragmatic and
cheap, but unrealistic for real CDC traffic and therefore not recommended as a final design.
Recommended direction: option 2 (keyed sink) as the foundation + option 1 (per-key collapse) as
the reduction semantics on top. Treating them as independent alternatives is the trap; they are
layers of the same solution, not competing options.
RowKind reduction matrix (needed before implementation)
"Last-write-wins collapse" is precise only once the reduction table is pinned down. Below is the
proposed reduction of a per-key event stream to a single final action (to be validated during
Phase B):
Last observed RowKind for k |
Final action |
INSERT (+I) |
append k |
UPDATE_AFTER (+U) |
delete k + append k (upsert) |
DELETE (-D) |
delete k |
UPDATE_BEFORE (-U) |
drop (only valid if followed by its paired +U) |
Open edge cases to resolve:
- A lone
-U(k) with no following +U(k) — dropping it loses the intent of the surrounding
sequence; needs an explicit decision.
-D(k) then +U(k) in the same window — reduces to a single upsert, not two actions.
Idempotency interaction with Phase C
Collapse fixes ordering, not idempotency. Deletes are naturally idempotent (re-applying the same
PK predicate is a no-op), but appends are not — replaying a checkpoint that contains an insert
would duplicate the row. The exactly-once replay guarantee still requires Phase C's transaction /
affected_rows (or equivalent) work; this issue must not imply that collapse alone delivers it.
Acceptance criteria
Minimal reproduction
Without waiting for the full Phase B/C sink, the bug is reproducible by driving the sink's flush
directly in a unit test: open a Lance dataset, invoke() an INSERT RowData(k) then a DELETE
RowData(k) with the same PK, force a checkpoint (or call the delete-then-insert flush), and assert
countRows() is 0 (it will be 1 today). This is the fastest way to confirm the failure before the
fix lands.
Severity
P0 / data-correctness blocker for full CDC support: it silently produces a wrong final dataset
state (a row that upstream deleted remains present), with no error and no log. It does not affect
insert-only pipelines (no PK), which is why it does not block Phase 0 or Phase A.
Relationship to #63
Background
This is a derived issue from the DELETE-via-PK epic (#63). It captures a correctness bug that
@LuciferYang flagged during review of #63 and asked to be raised separately rather than mixed into
the epic thread. This document is the standalone tracking ticket for that bug.
The epic's Phase B and Phase C introduce two independent decisions that, when composed, produce a
wrong final dataset state:
RowDatabyRowKind:DELETErows go into a delete bucket,INSERT/UPDATE_AFTERrows go into an insert bucket.The bug
Consider a single checkpoint window in which both messages below are buffered for the same
primary key
k:In a real CDC stream these two have a meaningful temporal order:
k+I(k)then-D(k)kmust not exist (inserted, then deleted)-D(k)then+I(k)kmust exist (deleted, then re-inserted)But because the sink flushes deletes first, then inserts regardless of the actual per-key
order, both real orders are collapsed into the same sequence:
delete(k)— removekfrom Lance (even though it may not exist yet).insert(k)— appendkto Lance.Final state in both cases:
kexists. For the+I(k)-then--D(k)case this is wrong — thedelete was effectively dropped, leaving a row that upstream already removed. This violates the
exactly-once / last-write-wins contract that a CDC sink must honor.
Root cause (abstraction level)
The ordering bug is a symptom, not the root cause. The real defect is that bucketing by
RowKinddestroys the temporal ordering between multiple events that share the same key. Once+I(k)lands in the insert bucket and-D(k)lands in the delete bucket, the "who came first"information is irretrievably lost from the bucket structure. The fixed delete-before-insert order
merely makes that lost information visible as a wrong final state.
This is why "just flip the order" is not a fix: insert-before-delete would instead break the
-D(k)-then-+I(k)case. Any correct fix must either preserve per-key order (don't bucket byRowKind; process keyed and in arrival order) or collapse to a final state before bucketing
(reduce per key first).
Scope is broader than the
+I/-Dexample+I(k)/-D(k)is the minimal illustration used in the epic review, but the same defect recursfor any same-key multi-event window. Given epic Phase B's handling of the four
RowKinds, theaffected combinations are:
k+I(k)then-D(k)-D(k)then+I(k)+U(k)then-D(k)+U(k1)then+U(k2)k2k1andk2present ❌(
+U=UPDATE_AFTER, which epic Phase B expands to delete-by-PK + append; two consecutive+Us therefore produce two inserts unless collapsed.)Why "deletes first, then inserts" exists in the first place
The fixed order is safe only when delete and insert targets are disjoint keys (the common case:
an
UPDATE_AFTERon keykis implemented asdelete(k)+insert(k')). There, ordering mattersfor a different reason: you must delete before re-inserting so the delete predicate does not also
match the freshly inserted row. The Phase C order was chosen for that disjoint-key reason, but it is
not sufficient for the same-key
+I/-Dcollision.Scope / affected code (once Phase A/B/C land)
LanceSink(or its successor): theRowKind-bucketing +flushDeletes-then-flushInsertscheckpoint logic.
LanceDynamicTableSink#getChangelogMode: declares the full+I/-U/+U/-Dset when a PK ispresent; this bug only manifests once
DELETEis actually accepted.Candidate fixes (to be evaluated during implementation)
Per-key last-write-wins collapse (the reduction logic). Before flushing, reduce the buffered
rows per PK to a single final action derived from the last
RowKindevent seen for that keywithin the checkpoint. A same-key
+I/-Dpair collapses to delete-only or insert-only,eliminating the ordering ambiguity. This matches how Paimon/Iceberg CDC sinks treat keyed
changelog.
PK is visible to the same collapsing buffer. The current
LanceSink extends RichSinkFunction(non-keyed) and
SinkFunctionProvider.of(...)(nokeyBy) provide no such guarantee — a+I(k)and-D(k)pair can be partitioned onto different subtasks, each seeing only half thepair. Therefore option 1 depends on option 2's keyed guarantee (or an equivalent explicit
partitioner / key-grouping). It is not "topology-neutral" as the original wording implied.
Keyed ordering guarantee (the enabling prerequisite). Use Flink's keyed state / keyed sink
(
KeyedProcessFunction) so all events for a given PK are processed by the same subtask inarrival order, then apply them in arrival order instead of a global delete-before-insert order.
This is the structural fix; option 1 is the reduction logic layered on top of it.
Reject mixed same-key operations within a checkpoint with a clear error. Pragmatic and
cheap, but unrealistic for real CDC traffic and therefore not recommended as a final design.
Recommended direction: option 2 (keyed sink) as the foundation + option 1 (per-key collapse) as
the reduction semantics on top. Treating them as independent alternatives is the trap; they are
layers of the same solution, not competing options.
RowKind reduction matrix (needed before implementation)
"Last-write-wins collapse" is precise only once the reduction table is pinned down. Below is the
proposed reduction of a per-key event stream to a single final action (to be validated during
Phase B):
RowKindforkINSERT(+I)kUPDATE_AFTER(+U)k+ appendk(upsert)DELETE(-D)kUPDATE_BEFORE(-U)+U)Open edge cases to resolve:
-U(k)with no following+U(k)— dropping it loses the intent of the surroundingsequence; needs an explicit decision.
-D(k)then+U(k)in the same window — reduces to a single upsert, not two actions.Idempotency interaction with Phase C
Collapse fixes ordering, not idempotency. Deletes are naturally idempotent (re-applying the same
PK predicate is a no-op), but appends are not — replaying a checkpoint that contains an insert
would duplicate the row. The exactly-once replay guarantee still requires Phase C's transaction /
affected_rows(or equivalent) work; this issue must not imply that collapse alone delivers it.Acceptance criteria
+I(k)followed by-D(k)in one checkpoint; finaldataset must not contain
k.-D(k)followed by+I(k)in one checkpoint; final dataset must containk.+U(k1)then+U(k2)in one checkpoint; final dataset mustcontain only the
k2state.k1, insertk2) remains correct (regression guard).Minimal reproduction
Without waiting for the full Phase B/C sink, the bug is reproducible by driving the sink's flush
directly in a unit test: open a Lance dataset,
invoke()anINSERTRowData(k)then aDELETERowData(k)with the same PK, force a checkpoint (or call the delete-then-insert flush), and assertcountRows()is 0 (it will be 1 today). This is the fastest way to confirm the failure before thefix lands.
Severity
P0 / data-correctness blocker for full CDC support: it silently produces a wrong final dataset
state (a row that upstream deleted remains present), with no error and no log. It does not affect
insert-only pipelines (no PK), which is why it does not block Phase 0 or Phase A.
Relationship to #63