Skip to content

SS-148 Iceberg sink, don't attempt commit unless it's definitely not a duplicate - #38441

Merged
ublubu merged 12 commits into
MaterializeInc:mainfrom
ublubu:iceberg-retry
Sep 15, 2026
Merged

ublubu merged 12 commits into
MaterializeInc:mainfrom
ublubu:iceberg-retry

Conversation

@ublubu

@ublubu ublubu commented Aug 24, 2026 •

Copy link
Copy Markdown
Contributor

replaces #38333

We rely on MaterializeInc/iceberg-rust#7 making the Transaction commit internals public.

Then we reimplement Transaction commit, but without the built-in rebase+retry.

Previously, iceberg-rust loaded the table from the Catalog on every commit attempt, applying the transaction's actions on top of the latest state of the table.

With this PR, Mz loads the table from the Catalog on every commit attempt, but it does not apply the transaction's action on top of the latest state of the table until after it inspects the table state (for funny business like a previously successful attempt or another writer taking over).

@ublubu
ublubu marked this pull request as ready for review August 26, 2026 20:57
@ublubu
ublubu requested a review from a team as a code owner August 26, 2026 20:57
@def-

def- commented Aug 26, 2026

Copy link
Copy Markdown
Contributor

QA LLM Review

1. MEDIUM -- Pre-commit guard only rejects an exactly-equal frontier, so a writer that committed a strict prefix of our batch still produces duplicate rows

src/storage/src/sink/iceberg.rs:883

The new guard proves "not a duplicate" only when last_frontier is exactly equal to frontier (idempotent retry) or >= it (fenced). It never checks last_frontier against batch_lower, so when another incarnation committed a snapshot that lands strictly inside our batch range (batch_lower < last_frontier < frontier) all three conditions fall through and we commit [batch_lower, frontier) on top of it, silently duplicating every row in [batch_lower, last_frontier) with no error surfaced.

Details

The reachable window is a sink version handoff, which is a normal ALTER SINK (see workflow_alter_commit_interval) or a replica restart where the old incarnation is still running. Version V reads the table, is not fenced (nothing newer is on it yet), and commits [F, G). Version V+1 had already read resume_upper = F at mint_batch_descriptions (src/storage/src/sink/iceberg.rs:1228) and minted its catch-up batch as [F, H) with H > G. At commit time V+1 sees last_version = V, so last_version > sink_version is false; last_frontier = G != H, so the idempotency arm is skipped; and PartialOrder::less_equal(H, G) is false, so the fencing arm at line 913 is skipped too. V+1 commits [F, H), duplicating [F, G). The same shape applies to two same-version incarnations racing after a replica restart, where mz-sink-version provides no fence at all.

last_frontier == batch_lower is an invariant of the minter in every single-writer case, so it can be asserted directly: no snapshots means a fresh table and batch_lower == as_of; the catch-up batch's lower is resume_upper, which is by construction the last committed frontier; and every subsequent batch's lower is the previous batch's upper. Third-party snapshots that copy our summary (as workflow_commit_conflict injects) and replace compactions both leave last_frontier at the sink's last committed frontier, so neither is disturbed.

Suggested fix, after the two existing fatal arms:

if last_frontier != *batch_lower {
    // Another writer committed inside our batch range. Committing now would
    // duplicate [batch_lower, last_frontier). Restart and re-resume instead.
    return (
        table,
        RetryResult::FatalErr(anyhow!(
            "Iceberg table '{}' advanced to frontier {:?} inside batch [{}, {}); \
             another writer is active.",
            conn_table,
            last_frontier,
            batch_lower.pretty(),
            batch_upper.pretty(),
        )),
    );
}

2. MEDIUM -- workflow_idempotent_retry can drop the response to an empty commit, in which case it passes with the pre-fix duplicate-write behavior

test/iceberg/mzcompose.py:405

The drop is armed as soon as commits_ok >= 1 and before the first INSERT, so the commit whose response gets dropped is frequently one that carries no data files. When that happens the test proves nothing about duplicate writes: replaying an empty RowDeltaAction adds no rows, so idempotent-retry-verify.td's exact count of 13 still holds, messages_committed >= 13 still holds, and the sink is still running. The regression this workflow exists to catch (reverting to Transaction::commit, whose rebase-and-retry re-adds the dropped batch's data files) would slip through.

Details

Empty commits are the steady state here: the sink commits on the 2s interval whether or not data arrived, the three setup rows all land in the snapshot batch [as_of, as_of+1), and the catch-up batch [as_of+1, observed_frontier) that immediately follows it is therefore empty. So the armed drop lands on a data-carrying commit only if the first c.sql INSERT round-trip completes before the next batch boundary, which is a timing race rather than something the test establishes.

The proxy already has the request body in hand at test/iceberg/polaris_proxy.py:96, so the fire condition can require a data-carrying commit, for example by only setting should_drop when the body's add-snapshot update has a non-zero added-data-files/added-records in its summary. Asserting after the fact that the dropped commit carried data would work equally well.

@patrickwwbutler patrickwwbutler left a comment •

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I think just fix this main issue and it looks ok? Would also take a look at the test setup issue in the LLM review but that's less important to me.

Merge issues probably just from my iceberg stuff last week, should be pretty easy although you'll need to rebase your iceberg-rust changes and update the revision

Comment thread src/storage/src/sink/iceberg.rs Outdated
Comment thread Cargo.toml Outdated
@def-

def- commented Sep 2, 2026

Copy link
Copy Markdown
Contributor

QA LLM Review

1. MEDIUM -- Hand-rolled commit drops iceberg-rust's guard against committing to an encrypted table

src/storage/src/sink/iceberg.rs:823

Transaction::commit, which do_commit replaces, refuses to commit to a table whose encryption.key-id is set, because iceberg-rust does not support encrypted writes. The new path carries no such check, so a sink pointed at a table configured for encryption at rest now writes plaintext data files and plaintext manifests and commits a snapshot referencing them, with no error surfaced.

Details

The guard is in Transaction::commit at the pinned rev (crates/iceberg/src/transaction/mod.rs): table_properties()?.encryption_key_id.is_some() returns ErrorKind::FeatureUnsupported. It is the only thing standing in the way, because nothing below it is encryption-aware: SnapshotProducer::commit opens the manifest list with a plain table.file_io().new_output(...), and the sink's own DataFileWriterBuilder path writes plaintext Parquet.

Reachable shape: a table pre-created in the catalog with encryption.key-id set and no snapshots yet. load_or_create_table (src/storage/src/sink/iceberg.rs:971) loads any existing table with no property or format-version validation, and RowDeltaOperation::existing_manifest short-circuits when there is no current snapshot, so nothing needs decrypting and the commit lands. Every later commit then sees only the sink's own unencrypted snapshots and keeps succeeding.

Before this change tx.commit(catalog) returned FeatureUnsupported, which reached mz_sink_statuses as "Cannot commit to an encrypted table" and left the table untouched. The regression is from a visible refusal to a silent write that breaks the table's encryption-at-rest guarantee and leaves encryption-aware readers with a table whose snapshots disagree.

Best fixed by rejecting the table up front rather than at commit time: check encryption_key_id where the table is loaded, so the sink fails with a clear terminal error. Putting it in do_commit would classify as CommitError::Local, i.e. retryable, and just burn the five attempts before halting.

@patrickwwbutler patrickwwbutler left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM

@ublubu

ublubu commented Sep 15, 2026

Copy link
Copy Markdown
Contributor Author

added a lil formatting fix for the broken testdrive from #38732

@ublubu
ublubu merged commit f631c65 into MaterializeInc:main Sep 15, 2026
82 checks passed
@ublubu
ublubu deleted the iceberg-retry branch September 15, 2026 18:42
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants