apache/seatunnel · warning

CorrelationId is missing but required, rejecting message tag

Error message

CorrelationId is missing but required, rejecting message tag: {}

What it means

In manual-ack mode with correlation-id verification enabled, each delivery must carry a correlationId used to match responses to requests. If a message arrives with a null correlationId, verifyMessageIdentifier logs this warning and rejects the message (basicReject without requeue) instead of acknowledging it.

Solutions

  1. Fix the producer to set the correlationId property on every published message
  2. If correlation matching is not needed, disable usesCorrelationId in the source config
  3. Enable auto-ack if request/response correlation semantics are not required
  4. Check for messages already dead-lettered by the reject (basicReject requeue=false) and replay them after fixing the producer

Example fix

// before (producer)
channel.basicPublish("", queue, null, body);
// after
AMQP.BasicProperties props = new AMQP.BasicProperties.Builder().correlationId(correlationId).build();
channel.basicPublish("", queue, props, body);
Defensive patterns

Strategy: validation

Validate before calling

// producer-side guard before publishing
if (message.getProperties() == null || message.getProperties().getCorrelationId() == null) throw new IllegalArgumentException("correlationId required");

Type guard

static boolean hasCorrelationId(AMQP.BasicProperties props) { return props != null && props.getCorrelationId() != null && !props.getCorrelationId().isEmpty(); }

Try / catch

try { /* pollNext */ } catch (RabbitmqConnectorException e) { if (MESSAGE_ACK_REJECTED.equals(e.getErrorCode())) { /* producer sent uncorrelated message; inspect dead-letter queue */ } }

Prevention

When it happens

Trigger: pollNext -> verifyMessageIdentifier receives correlationId==null while autoAck=false and usesCorrelationId=true; the message's properties have no correlation id set by the producer.

Common situations: Producers publishing RPC-style messages without setting the correlationId property; mixed producers on the same queue, some legacy ones omitting correlation ids; enabling usesCorrelationId on queues fed by producers that never set it.

Understand the failure class

Background: "missing required argument" and "the following required arguments were not provided": what required-argument errors mean and how to fix them — this error's family across 20 libraries.

Related errors


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

                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);
                } 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;
            }
        }

View on GitHub (pinned to cf67b549a7)