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
- Check network stability between the SeaTunnel worker and the RabbitMQ broker; fix flaky connections first
- Reduce checkpoint interval so acks happen before broker heartbeat/channel timeouts close the channel
- Inspect the wrapped cause for 'channel is closed' or 'unknown delivery tag' and re-create the channel on failure
- Verify txSelect() was called on the channel if relying on txCommit for transactional acks
- 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
- Keep checkpoint interval shorter than broker heartbeat/channel timeout
- Enable automatic connection recovery in the connection factory
- Always call txSelect() when relying on txCommit
- Monitor network stability between workers and broker
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
- Both channel and connection closing failed. Logging channel…
- CorrelationId is missing but required, rejecting message tag
- RABBITMQ-02
- RABBITMQ-03
- RABBITMQ-04
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)