{"record":{"id":"c493c3c2f88a2be4","repo":"apache/seatunnel","slug":"correlationid-is-missing-but-required-rejecting-m","errorCode":null,"errorMessage":"CorrelationId is missing but required, rejecting message tag: {}","messagePattern":"CorrelationId is missing but required, rejecting message tag: (.+?)","errorType":"console","errorClass":null,"httpStatus":null,"severity":"warning","filePath":"seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/source/RabbitmqSourceReader.java","lineNumber":278,"sourceCode":"                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);\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        }","sourceCodeStart":260,"sourceCodeEnd":296,"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#L260-L296","documentation":"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.","triggerScenarios":"pollNext -> verifyMessageIdentifier receives correlationId==null while autoAck=false and usesCorrelationId=true; the message's properties have no correlation id set by the producer.","commonSituations":"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.","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"],"exampleFix":"// before (producer)\nchannel.basicPublish(\"\", queue, null, body);\n// after\nAMQP.BasicProperties props = new AMQP.BasicProperties.Builder().correlationId(correlationId).build();\nchannel.basicPublish(\"\", queue, props, body);","handlingStrategy":"validation","validationCode":"// producer-side guard before publishing\nif (message.getProperties() == null || message.getProperties().getCorrelationId() == null) throw new IllegalArgumentException(\"correlationId required\");","typeGuard":"static boolean hasCorrelationId(AMQP.BasicProperties props) { return props != null && props.getCorrelationId() != null && !props.getCorrelationId().isEmpty(); }","tryCatchPattern":"try { /* pollNext */ } catch (RabbitmqConnectorException e) { if (MESSAGE_ACK_REJECTED.equals(e.getErrorCode())) { /* producer sent uncorrelated message; inspect dead-letter queue */ } }","preventionTips":["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"],"tags":["rabbitmq","message-validation","ack"],"backgroundTag":"missing-required-argument","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"}