{"record":{"id":"cbb7248e9fb839b8","repo":"apache/seatunnel","slug":"rabbitmq-06","errorCode":"RABBITMQ-06","errorMessage":"messages could not be acknowledged with basicReject","messagePattern":"messages could not be acknowledged with basicReject","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":284,"sourceCode":"    }\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);\n                } catch (IOException e) {\n                    throw new RabbitmqConnectorException(MESSAGE_ACK_REJECTED, e);\n                }\n                return false;\n            }\n            if (!correlationIdsProcessedButNotAcknowledged.add(correlationId)) {\n                try {\n                    channel.basicReject(deliveryTag, false);\n                } catch (IOException e) {\n                    throw new RabbitmqConnectorException(MESSAGE_ACK_REJECTED, e);\n                }\n                return false;\n            }\n        }\n        return true;\n    }\n\n    @Override\n    public void close() throws IOException {\n        for (Map.Entry<String, DefaultConsumer> entry : activeConsumers.entrySet()) {","sourceCodeStart":266,"sourceCodeEnd":302,"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#L266-L302","documentation":"Thrown by RabbitmqSourceReader.verifyMessageIdentifier when a message lacking the required CorrelationId must be rejected and channel.basicReject(deliveryTag, false) itself fails with an IOException. Error code MESSAGE_ACK_REJECTED indicates the broker could not be told to drop/requeue the message.","triggerScenarios":"pollNext -> verifyMessageIdentifier; a delivered message has a missing CorrelationId, and the basicReject call hits a closed channel, broken connection, or invalid delivery tag.","commonSituations":"Producer application not setting correlationId on outgoing messages; broker restarted or channel closed before reject; stale delivery tag after connection recovery; network blip during message validation.","solutions":["Fix the upstream producer so every message carries a CorrelationId, avoiding the reject path entirely","Check the wrapped cause for 'unknown delivery tag' or 'channel closed' and reconnect/recreate the channel","Enable connection auto-recovery in the connection factory so channels survive broker restarts","Verify broker reachability (port 5672/5671, vhost, credentials) if the failure is a dropped connection"],"exampleFix":"// before: reject without checking channel state\nchannel.basicReject(deliveryTag, false);\n// after: guard against closed channel\nif (!channel.isOpen()) {\n    throw new RabbitmqConnectorException(MESSAGE_ACK_REJECTED, new IOException(\"channel closed; message will be redelivered and rejected again\"));\n}\nchannel.basicReject(deliveryTag, false);","handlingStrategy":"validation","validationCode":"// producer side: ensure correlationId set before publish\nif (properties.getCorrelationId() == null) { throw new IllegalArgumentException(\"correlationId required\"); }","typeGuard":null,"tryCatchPattern":"try { reader.verifyMessageIdentifier(delivery); } catch (RabbitmqConnectorException e) {\n  if (!channel.isOpen()) { /* recreate channel; message will be redelivered */ }\n}","preventionTips":["Enforce correlationId on all upstream producers","Enable AMQP auto-recovery for channels","Keep channels alive; avoid long idle periods","Log and alert on missing-correlationId messages to fix producers early"],"tags":["rabbitmq","reject","io","message-validation"],"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"}