Skip to content

[design risk] DELETE is silently dropped when +I(k) and -D(k) arrive in the same checkpoint #74

Description

@fightBoxing

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:

  1. delete(k) — remove k from Lance (even though it may not exist yet).
  2. 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)

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

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

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

  • A test seeds a dataset, then feeds +I(k) followed by -D(k) in one checkpoint; final
    dataset must not contain k.
  • A test feeds -D(k) followed by +I(k) in one checkpoint; final dataset must contain k.
  • A test feeds two consecutive +U(k1) then +U(k2) in one checkpoint; final dataset must
    contain only the k2 state.
  • The disjoint-key case (delete k1, insert k2) remains correct (regression guard).
  • Replay / exactly-once: re-flushing the same checkpoint produces the same final state.

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

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions