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
- Fix the producer to set the correlationId property on every published message
- If correlation matching is not needed, disable usesCorrelationId in the source config
- Enable auto-ack if request/response correlation semantics are not required
- 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
- Set correlationId on every produced message
- Only enable usesCorrelationId when all producers set it
- Monitor the dead-letter queue for rejected uncorrelated messages
- Document the correlation contract for queue producers
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)