{"record":{"id":"3cd35c1a026c8f04","repo":"apache/seatunnel","slug":"checkpoint-do-not-exist-or-have-already-been-co","errorCode":null,"errorMessage":"checkpoint {} do not exist or have already been committed.","messagePattern":"checkpoint (.+?) do not exist or have already been committed\\.","errorType":"console","errorClass":null,"httpStatus":null,"severity":"warning","filePath":"seatunnel-connectors-v2/connector-rocketmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rocketmq/source/RocketMqSourceReader.java","lineNumber":246,"sourceCode":"                s -> {\n                    try {\n                        pendingPartitionsQueue.put(s);\n                    } catch (InterruptedException e) {\n                        throw new RocketMqConnectorException(\n                                RocketMqConnectorErrorCode.ADD_SPLIT_CHECKPOINT_FAILED, e);\n                    }\n                });\n    }\n\n    @Override\n    public void handleNoMoreSplits() {\n        // No-op\n    }\n\n    @Override\n    public void notifyCheckpointComplete(long checkpointId) throws Exception {\n        if (!checkpointOffsets.containsKey(checkpointId)) {\n            log.warn(\"checkpoint {} do not exist or have already been committed.\", checkpointId);\n        } else {\n            Map<MessageQueue, Long> messageQueueOffset = checkpointOffsets.remove(checkpointId);\n            for (Map.Entry<MessageQueue, Long> entry : messageQueueOffset.entrySet()) {\n                MessageQueue messageQueue = entry.getKey();\n                Long offset = entry.getValue();\n                try {\n                    if (messageQueue != null && offset != null) {\n                        RocketMqConsumerThread rocketMqConsumerThread =\n                                consumerThreads.get(messageQueue);\n                        if (rocketMqConsumerThread != null) {\n                            rocketMqConsumerThread\n                                    .getTasks()\n                                    .put(\n                                            consumer -> {\n                                                if (this.metadata.isEnabledCommitCheckpoint()) {\n                                                    consumer.getOffsetStore()\n                                                            .updateOffset(\n                                                                    messageQueue, offset, false);","sourceCodeStart":228,"sourceCodeEnd":264,"githubUrl":"https://github.com/apache/seatunnel/blob/cf67b549a7a6c35fa0beb12d83c62892427ea919/seatunnel-connectors-v2/connector-rocketmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rocketmq/source/RocketMqSourceReader.java#L228-L264","documentation":"RocketMqSourceReader.notifyCheckpointComplete() commits the offsets recorded for a completed checkpoint. If the checkpointId is not present in checkpointOffsets it logs this warning ('checkpoint {} do not exist or have already been committed.') and does nothing — the offsets were already committed (e.g. on a prior snapshotState/notify pair) or were never registered (e.g. after recovery). It is benign in most cases.","triggerScenarios":"The engine calls notifyCheckpointComplete(checkpointId) with an ID that was never stored by snapshotState, or that notifyCheckpointComplete/remove already consumed; commonly on duplicated checkpoint notifications, restored jobs replaying old checkpoint IDs, or after notifyCheckpointAborted removed the entry.","commonSituations":"Zeta engine restart/failover replays checkpoint completion callbacks for checkpoints taken before recovery; rapid checkpoint completion racing with an aborted checkpoint; job restored from a savepoint so pre-restore checkpoint IDs are unknown to the reader; duplicate callback delivery.","solutions":["Treat as benign if the offsets were already committed — no data loss occurs because RocketMQ commits are idempotent per queue/offset.","If it happens for the FIRST checkpoint after restore, verify savepoint/restore logic registers pending offsets correctly in snapshotState after restart.","Check for duplicate checkpoint-complete callbacks; if persistent, compare engine checkpoint IDs with reader state and look for notifyCheckpointAborted interleaving.","Upgrade SeaTunnel if a known bug in the reader's checkpoint bookkeeping (checkpointOffsets lifecycle) matches your version."],"exampleFix":null,"handlingStrategy":"try-catch","validationCode":null,"typeGuard":null,"tryCatchPattern":"// benign warning; only act if it fires for every checkpoint\nif (logMatchesUnknownCheckpoint && firstCheckpointAfterRestore) {\n    // verify snapshotState re-registered offsets after restore\n}","preventionTips":["Restore jobs from savepoints/checkpoints through the engine's supported path","Avoid restarting readers outside engine coordination (prevents ID mismatches)","Treat occasional occurrences as benign — commit is idempotent","Keep SeaTunnel/RocketMQ connector versions current for checkpoint bookkeeping fixes"],"tags":["rocketmq","checkpoint","offset-commit","idempotency"],"backgroundTag":"invalid-state-transition","analyzedSha":"cf67b549a7a6c35fa0beb12d83c62892427ea919","analyzedAt":"2026-09-10T21:44:55.265Z","contentChangedAt":"2026-09-10T21:44:55.265Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}