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

  1. Treat as benign if the offsets were already committed — no data loss occurs because RocketMQ commits are idempotent per queue/offset.
  2. If it happens for the FIRST checkpoint after restore, verify savepoint/restore logic registers pending offsets correctly in snapshotState after restart.
  3. Check for duplicate checkpoint-complete callbacks; if persistent, compare engine checkpoint IDs with reader state and look for notifyCheckpointAborted interleaving.
  4. 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

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


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)