{"record":{"id":"6f2263eebf72bfde","repo":"apache/seatunnel","slug":"rabbitmq-05","errorCode":"RABBITMQ-05","errorMessage":"messages could not be acknowledged during checkpoint creation","messagePattern":"messages could not be acknowledged during checkpoint creation","errorType":"error_code","errorClass":"RabbitmqConnectorException","httpStatus":null,"severity":"error","filePath":"seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/source/RabbitmqSourceReader.java","lineNumber":264,"sourceCode":"        List<Long> pendingDeliveryTags = pendingDeliveryTagsToCommit.remove(checkpointId);\n        Set<String> pendingCorrelationIds = pendingCorrelationIdsToCommit.remove(checkpointId);\n\n        if (pendingDeliveryTags != null && !autoAck) {\n            acknowledgeDeliveryTags(pendingDeliveryTags);\n        }\n        if (pendingCorrelationIds != null) {\n            correlationIdsProcessedButNotAcknowledged.removeAll(pendingCorrelationIds);\n        }\n    }\n\n    protected void acknowledgeDeliveryTags(List<Long> deliveryTags) {\n        try {\n            for (long id : deliveryTags) {\n                channel.basicAck(id, false);\n            }\n            channel.txCommit();\n        } catch (IOException e) {\n            throw new RabbitmqConnectorException(MESSAGE_ACK_FAILED, e);\n        }\n    }\n\n    /**\n     * Verify message identifier.\n     *\n     * @param correlationId correlation id\n     * @param deliveryTag delivery tag\n     * @return true if valid\n     */\n    public boolean verifyMessageIdentifier(String correlationId, long deliveryTag) {\n        if (!autoAck && usesCorrelationId) {\n            if (correlationId == null) {\n                log.warn(\n                        \"CorrelationId is missing but required, rejecting message tag: {}\",\n                        deliveryTag);\n                try {\n                    channel.basicReject(deliveryTag, false);","sourceCodeStart":246,"sourceCodeEnd":282,"githubUrl":"https://github.com/apache/seatunnel/blob/cf67b549a7a6c35fa0beb12d83c62892427ea919/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/source/RabbitmqSourceReader.java#L246-L282","documentation":"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.","triggerScenarios":"notifyCheckpointComplete -> acknowledgeDeliveryTags, when the channel is closed by the broker, connection dropped, or txCommit fails on a non-transactional/closed channel.","commonSituations":"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.","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"],"exampleFix":"// before: blind ack loop throws if channel died\ntry {\n    for (long id : deliveryTags) { channel.basicAck(id, false); }\n    channel.txCommit();\n} catch (IOException e) { throw new RabbitmqConnectorException(MESSAGE_ACK_FAILED, e); }\n// after: check channel is open before acking\nif (channel.isOpen()) {\n    for (long id : deliveryTags) { channel.basicAck(id, false); }\n    channel.txCommit();\n} else {\n    throw new RabbitmqConnectorException(MESSAGE_ACK_FAILED, new IOException(\"channel closed, acks will be redelivered\"));\n}","handlingStrategy":"try-catch","validationCode":"// before acking\nif (!channel.isOpen()) { /* recreate channel/connection before ack */ }","typeGuard":null,"tryCatchPattern":"try { reader.acknowledgeDeliveryTags(...); } catch (RabbitmqConnectorException e) {\n  // rely on broker redelivery: reset delivery tags, recreate channel, retry from last checkpoint\n}","preventionTips":["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"],"tags":["rabbitmq","checkpoint","ack","network"],"backgroundTag":"network-request-failed","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"}