apache/seatunnel · error · RabbitmqConnectorException

RABBITMQ-05

RABBITMQ-05

Error message

messages could not be acknowledged during checkpoint creation

What it means

Thrown by RabbitmqSourceReader.acknowledgeDeliveryTags when acking messages to RabbitMQ (channel.basicAck per tag followed by channel.txCommit) fails with an IOException. It signals that processed messages could not be confirmed to the broker during checkpoint creation, wrapped as MESSAGE_ACK_FAILED.

Solutions

  1. Check network stability between the SeaTunnel worker and the RabbitMQ broker; fix flaky connections first
  2. Reduce checkpoint interval so acks happen before broker heartbeat/channel timeouts close the channel
  3. Inspect the wrapped cause for 'channel is closed' or 'unknown delivery tag' and re-create the channel on failure
  4. Verify txSelect() was called on the channel if relying on txCommit for transactional acks
  5. Enable RabbitMQ connection recovery (automaticRecovery) in the connection factory

Example fix

// before: blind ack loop throws if channel died
try {
    for (long id : deliveryTags) { channel.basicAck(id, false); }
    channel.txCommit();
} catch (IOException e) { throw new RabbitmqConnectorException(MESSAGE_ACK_FAILED, e); }
// after: check channel is open before acking
if (channel.isOpen()) {
    for (long id : deliveryTags) { channel.basicAck(id, false); }
    channel.txCommit();
} else {
    throw new RabbitmqConnectorException(MESSAGE_ACK_FAILED, new IOException("channel closed, acks will be redelivered"));
}
Defensive patterns

Strategy: try-catch

Validate before calling

// before acking
if (!channel.isOpen()) { /* recreate channel/connection before ack */ }

Try / catch

try { reader.acknowledgeDeliveryTags(...); } catch (RabbitmqConnectorException e) {
  // rely on broker redelivery: reset delivery tags, recreate channel, retry from last checkpoint
}

Prevention

When it happens

Trigger: notifyCheckpointComplete -> acknowledgeDeliveryTags, when the channel is closed by the broker, connection dropped, or txCommit fails on a non-transactional/closed channel.

Common situations: Long checkpoint intervals letting the broker close idle channels; network interruption between checkpoint and broker; broker restart; using autoAck=false with a channel that was recovered externally, invalidating delivery tags.

Related errors


AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10). Data as JSON: /api/errors/6f2263eebf72bfde. Report an issue: GitHub.

Appendix: source

Thrown at seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/source/RabbitmqSourceReader.java:264

        List<Long> pendingDeliveryTags = pendingDeliveryTagsToCommit.remove(checkpointId);
        Set<String> pendingCorrelationIds = pendingCorrelationIdsToCommit.remove(checkpointId);

        if (pendingDeliveryTags != null && !autoAck) {
            acknowledgeDeliveryTags(pendingDeliveryTags);
        }
        if (pendingCorrelationIds != null) {
            correlationIdsProcessedButNotAcknowledged.removeAll(pendingCorrelationIds);
        }
    }

    protected void acknowledgeDeliveryTags(List<Long> deliveryTags) {
        try {
            for (long id : deliveryTags) {
                channel.basicAck(id, false);
            }
            channel.txCommit();
        } catch (IOException e) {
            throw new RabbitmqConnectorException(MESSAGE_ACK_FAILED, e);
        }
    }

    /**
     * Verify message identifier.
     *
     * @param correlationId correlation id
     * @param deliveryTag delivery tag
     * @return true if valid
     */
    public boolean verifyMessageIdentifier(String correlationId, long deliveryTag) {
        if (!autoAck && usesCorrelationId) {
            if (correlationId == null) {
                log.warn(
                        "CorrelationId is missing but required, rejecting message tag: {}",
                        deliveryTag);
                try {
                    channel.basicReject(deliveryTag, false);

View on GitHub (pinned to cf67b549a7)