{"record":{"id":"eee5c7b058185a37","repo":"risingwavelabs/risingwave","slug":"iceberg-pk-index-sink-commit-failed-for-sink-ep","errorCode":null,"errorMessage":"iceberg pk-index sink commit failed for sink {} epoch {}","messagePattern":"iceberg pk-index sink commit failed for sink (.+?) epoch (.+?)","errorType":"exception","errorClass":null,"httpStatus":null,"severity":"critical","filePath":"src/meta/src/manager/iceberg_pk_index_sink/coordinator.rs","lineNumber":203,"sourceCode":"    pub async fn commit(&mut self) -> Result<()> {\n        let Some(commit) = self.waiting_commit.take() else {\n            return Ok(());\n        };\n\n        let refreshed_table = commit_one_epoch(\n            self.catalog.clone(),\n            self.table.identifier().clone(),\n            self.target_branch.clone(),\n            self.sink_id,\n            &commit,\n            self.retry_num,\n        )\n        .await\n        .map_err(|err| {\n            let err_report = match err {\n                CommitError::Commit(e) | CommitError::ReloadTable(e) => e,\n            };\n            anyhow!(err_report).context(format!(\n                \"iceberg pk-index sink commit failed for sink {} epoch {}\",\n                self.sink_id, commit.epoch\n            ))\n        })?;\n        self.table = refreshed_table;\n\n        commit_and_prune_epoch(\n            &self.db,\n            self.sink_id,\n            commit.epoch,\n            self.prev_committed_epoch,\n        )\n        .await\n        .with_context(|| {\n            format!(\n                \"iceberg pk-index sink mark_committed failed for sink {} epoch {}\",\n                self.sink_id, commit.epoch\n            )","sourceCodeStart":185,"sourceCodeEnd":221,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/meta/src/manager/iceberg_pk_index_sink/coordinator.rs#L185-L221","documentation":"This error wraps the underlying iceberg commit failure (CommitError::Commit or CommitError::ReloadTable from commit_one_epoch) with context identifying the sink and epoch. After appending files to the Iceberg table (or reloading it afterwards) fails beyond commit_retry_num attempts, the coordinator surfaces this contextual error; the epoch stays in pending_sink_state so it can be retried.","triggerScenarios":"Calling commit() (live or during recovery drain) when commit_one_epoch fails: the Iceberg table commit itself errors (commit conflict/exhausted retries) or reloading the table after commit fails.","commonSituations":"Concurrent external writers to the same Iceberg table causing repeated commit conflicts beyond retry count, catalog connectivity/auth failures during commit or table reload, snapshot expiration racing the commit, branch missing (wrong branch config) making reload fail.","solutions":["Inspect the wrapped underlying error (context chain) for the real cause: commit conflict vs reload failure","Retry the commit; pending state persisted in pending_sink_state makes it idempotent and recoverable on restart","Increase commit_retry_num in iceberg config if conflicts with concurrent external writers are frequent","Stop/coordinate external writers to the same table/branch, or write to a dedicated branch","Verify catalog credentials and that the target branch exists and is writable"],"exampleFix":"// before\niceberg.commit_retry_num = 3 // conflicts with external writer\n// after\niceberg.commit_retry_num = 10\n// and/or: write to dedicated branch\niceberg.write_mode = ... // ensure commit_branch targets an RW-only branch","handlingStrategy":"retry","validationCode":null,"typeGuard":null,"tryCatchPattern":"match coordinator.commit().await {\n    Err(e) if e.to_string().contains(\"commit failed for sink\") => {\n        // inspect the source: Commit vs ReloadTable\n        warn!(error = ?e.source(), \"iceberg commit failed; will retry from pending state\");\n        retry_with_backoff(|| coordinator.commit()).await\n    }\n    other => other,\n}","preventionTips":["Increase commit_retry_num when external writers contend on the same table","Write to a dedicated branch so external maintenance does not conflict with sink commits","Persist-and-retry is built in: pending epochs survive restarts, so re-run commit rather than re-writing data","Monitor catalog auth/token expiry which commonly surfaces as reload failures"],"tags":["iceberg","commit","conflict","exactly-once"],"backgroundTag":"database-write-failed","analyzedSha":"6469eb736d691e8e9b8a419a57edd6429ca77417","analyzedAt":"2026-09-11T21:06:21.487Z","contentChangedAt":"2026-09-11T21:06:21.487Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}