{"record":{"id":"a9309a2a2984bdb6","repo":"apache/iceberg","slug":"received-commit-ready-when-no-commit-in-progress","errorCode":null,"errorMessage":"Received commit ready when no commit in progress, this can happen during recovery. Commit ID: {}","messagePattern":"Received commit ready when no commit in progress, this can happen during recovery\\. Commit ID: (.+?)","errorType":"console","errorClass":null,"httpStatus":null,"severity":"warning","filePath":"kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/CommitState.java","lineNumber":65,"sourceCode":"  CommitState(IcebergSinkConfig config) {\n    this.config = config;\n  }\n\n  void addResponse(Envelope envelope) {\n    commitBuffer.add(envelope);\n    if (!isCommitInProgress()) {\n      DataWritten dataWritten = (DataWritten) envelope.event().payload();\n      LOG.warn(\n          \"Received commit response when no commit in progress, this can happen during recovery. Commit ID: {}\",\n          dataWritten.commitId());\n    }\n  }\n\n  void addReady(Envelope envelope) {\n    DataComplete dataComplete = (DataComplete) envelope.event().payload();\n    readyBuffer.add(dataComplete);\n    if (!isCommitInProgress()) {\n      LOG.warn(\n          \"Received commit ready when no commit in progress, this can happen during recovery. Commit ID: {}\",\n          dataComplete.commitId());\n    } else if (Objects.equals(currentCommitId, dataComplete.commitId())) {\n      receivedPartitionCount += dataComplete.assignments().size();\n    }\n  }\n\n  UUID currentCommitId() {\n    return currentCommitId;\n  }\n\n  boolean isCommitInProgress() {\n    return currentCommitId != null;\n  }\n\n  boolean isCommitIntervalReached() {\n    if (startTime == 0) {\n      startTime = System.currentTimeMillis();","sourceCodeStart":47,"sourceCodeEnd":83,"githubUrl":"https://github.com/apache/iceberg/blob/86d9c8fc543e7c56c9f624eb725f76c9baff9570/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/CommitState.java#L47-L83","documentation":"A DataComplete (commit ready) envelope arrived while no commit was in progress, so its partition assignments cannot be counted toward a current commit. Like the commit-response case, this typically reflects stale messages replayed during recovery; the payload is still buffered but the partition count is not incremented.","triggerScenarios":"CommitState.addReady is called outside an active commit, or with a DataComplete whose commitId differs from currentCommitId (the commitId-match branch is skipped).","commonSituations":"Zombie coordinator from a previous generation still publishing DataComplete messages after a rebalance; recovery replay of coordinator-topic records; data workers responding to a commit that already timed out.","solutions":["Benign during recovery — stale DataComplete messages are logged and effectively ignored; no action needed for one-off occurrences.","Ensure old coordinator workers are fully stopped after rebalances so zombie coordinators stop publishing (check task liveness / fencing).","Check consumer offsets on the coordinator topic to avoid reprocessing old records after restarts.","If persistent, align all workers on the same connector version and verify commit timeouts are not shorter than data-worker response latency."],"exampleFix":null,"handlingStrategy":"retry","validationCode":null,"typeGuard":null,"tryCatchPattern":null,"preventionTips":["Confirm a single coordinator generation is active (no zombie coordinators after rebalance).","Keep data-worker response latency well below the commit timeout.","Verify consumer offsets on the coordinator topic to avoid replay after restarts.","Run compatible connector versions across all tasks."],"tags":["kafka","connect","coordinator","recovery"],"backgroundTag":"unexpected-response-shape","analyzedSha":"86d9c8fc543e7c56c9f624eb725f76c9baff9570","analyzedAt":"2026-09-12T00:46:39.097Z","contentChangedAt":"2026-09-12T00:46:39.097Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}