apache/seatunnel · warning
checkpoint do not exist or have already been committed.
Error message
checkpoint {} do not exist or have already been committed. What it means
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.
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.
Defensive patterns
Strategy: try-catch
Try / catch
// benign warning; only act if it fires for every checkpoint
if (logMatchesUnknownCheckpoint && firstCheckpointAfterRestore) {
// verify snapshotState re-registered offsets after restore
} Prevention
- 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
When it happens
Trigger: 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.
Common situations: 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.
Understand the failure class
Background: "Invalid state transition" errors: "status must be X, actually Y", "already rejected/charging/uninstalled", "cannot ... while running" — what they mean when a library rejects your call — this error's family across 31 libraries.
Related errors
- checkpointId is already set
- ACKNOWLEDGE_FAILED
- ACKNOWLEDGE_FAILED
- Ambiguous timeout on Couchbase write
- checkpoint do not exist or have already been committed.
AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10).
Data as JSON: /api/errors/3cd35c1a026c8f04.
Report an issue: GitHub.
Appendix: source
Thrown at seatunnel-connectors-v2/connector-rocketmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rocketmq/source/RocketMqSourceReader.java:246
s -> {
try {
pendingPartitionsQueue.put(s);
} catch (InterruptedException e) {
throw new RocketMqConnectorException(
RocketMqConnectorErrorCode.ADD_SPLIT_CHECKPOINT_FAILED, e);
}
});
}
@Override
public void handleNoMoreSplits() {
// No-op
}
@Override
public void notifyCheckpointComplete(long checkpointId) throws Exception {
if (!checkpointOffsets.containsKey(checkpointId)) {
log.warn("checkpoint {} do not exist or have already been committed.", checkpointId);
} else {
Map<MessageQueue, Long> messageQueueOffset = checkpointOffsets.remove(checkpointId);
for (Map.Entry<MessageQueue, Long> entry : messageQueueOffset.entrySet()) {
MessageQueue messageQueue = entry.getKey();
Long offset = entry.getValue();
try {
if (messageQueue != null && offset != null) {
RocketMqConsumerThread rocketMqConsumerThread =
consumerThreads.get(messageQueue);
if (rocketMqConsumerThread != null) {
rocketMqConsumerThread
.getTasks()
.put(
consumer -> {
if (this.metadata.isEnabledCommitCheckpoint()) {
consumer.getOffsetStore()
.updateOffset(
messageQueue, offset, false);View on GitHub (pinned to cf67b549a7)