apache/seatunnel · error · RabbitmqConnectorException

RABBITMQ-06

RABBITMQ-06

Error message

messages could not be acknowledged with basicReject

What it means

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.

Solutions

  1. Fix the upstream producer so every message carries a CorrelationId, avoiding the reject path entirely
  2. Check the wrapped cause for 'unknown delivery tag' or 'channel closed' and reconnect/recreate the channel
  3. Enable connection auto-recovery in the connection factory so channels survive broker restarts
  4. Verify broker reachability (port 5672/5671, vhost, credentials) if the failure is a dropped connection

Example fix

// before: reject without checking channel state
channel.basicReject(deliveryTag, false);
// after: guard against closed channel
if (!channel.isOpen()) {
    throw new RabbitmqConnectorException(MESSAGE_ACK_REJECTED, new IOException("channel closed; message will be redelivered and rejected again"));
}
channel.basicReject(deliveryTag, false);
Defensive patterns

Strategy: validation

Validate before calling

// producer side: ensure correlationId set before publish
if (properties.getCorrelationId() == null) { throw new IllegalArgumentException("correlationId required"); }

Try / catch

try { reader.verifyMessageIdentifier(delivery); } catch (RabbitmqConnectorException e) {
  if (!channel.isOpen()) { /* recreate channel; message will be redelivered */ }
}

Prevention

When it happens

Trigger: pollNext -> verifyMessageIdentifier; a delivered message has a missing CorrelationId, and the basicReject call hits a closed channel, broken connection, or invalid delivery tag.

Common situations: 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.

Related errors


AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10). Data as JSON: /api/errors/cbb7248e9fb839b8. 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:284

    }

    /**
     * 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);
                } catch (IOException e) {
                    throw new RabbitmqConnectorException(MESSAGE_ACK_REJECTED, e);
                }
                return false;
            }
            if (!correlationIdsProcessedButNotAcknowledged.add(correlationId)) {
                try {
                    channel.basicReject(deliveryTag, false);
                } catch (IOException e) {
                    throw new RabbitmqConnectorException(MESSAGE_ACK_REJECTED, e);
                }
                return false;
            }
        }
        return true;
    }

    @Override
    public void close() throws IOException {
        for (Map.Entry<String, DefaultConsumer> entry : activeConsumers.entrySet()) {

View on GitHub (pinned to cf67b549a7)