Spring Integration中ID与TIMESTAMP为何被声明为transient headers?
Hi László, let's break down your questions and tie in what you've already discovered to clarify things:
1. Why isn't MessageHeaders.ID automatically mapped to AmqpHeaders.MESSAGE_ID?
The core reason is that MessageHeaders.ID is marked as transient—this means it's intended to be an in-memory identifier only, not something that gets serialized or transmitted across message brokers. Spring Integration's AbstractHeaderMapper explicitly skips transient headers by default, so it won't automatically copy this value to the AMQP MESSAGE_ID header. This is by design, as the framework treats this ID as a local, per-message instance identifier rather than a cross-transport tracking ID.
2. Do I need a custom ID header and MessageKeyGenerator?
It depends on what you're trying to achieve. If you need a persistent, cross-transport message ID for tracking or idempotency, then yes—you should create a custom non-transient header (e.g., X-Custom-Message-ID) and populate it when sending messages.
You don't necessarily need a MessageKeyGenerator unless you're using a feature that requires a specific key for message deduplication or storage. For most cases, simply enriching the message with a custom ID before sending (using a header enricher) and configuring the AMQP header mapper to include this header will suffice.
That said, since you've already fixed the infinite loop issue with defaultRequeueRejected(false), if you don't need message ID tracking, you might not even need to tackle this part.
3. How to implement this with the latest Java DSL?
If you do want to map a custom ID (or even force the MessageHeaders.ID to be sent, though that's not recommended), you can configure a custom AmqpHeaderMapper in both your inbound and outbound adapters. Here's an example:
For Outbound Adapter (Sending Messages)
@Bean public IntegrationFlow webhookOutboundFlow(ConnectionFactory connectionFactory) { return IntegrationFlows.from("webhookOutboundChannel") // Enrich message with a custom persistent ID .enrichHeaders(h -> h.header("X-Custom-Message-ID", UUID.randomUUID().toString())) .handle(Amqp.outboundAdapter(connectionFactory) .routingKey("your-routing-key") .headerMapper(mapper -> { DefaultAmqpHeaderMapper amqpHeaderMapper = DefaultAmqpHeaderMapper.outboundMapper(); // Add your custom header to the list of mapped headers amqpHeaderMapper.setRequestHeaderNames("X-Custom-Message-ID", "other-headers-you-need"); return amqpHeaderMapper; })) .get(); }
For Inbound Adapter (Receiving Messages)
@Bean public IntegrationFlow webhookInboundFlow(ConnectionFactory connectionFactory, ObjectMapper objectMapper, HeaderValueRouter webhookInboundRouter) { return IntegrationFlows.from(Amqp.inboundAdapter(connectionFactory, FORGETME_WEBHOOK_QUEUE_NAME) .configureContainer(s -> s.defaultRequeueRejected(false)) .headerMapper(mapper -> { DefaultAmqpHeaderMapper amqpHeaderMapper = DefaultAmqpHeaderMapper.inboundMapper(); // Map the custom header from AMQP to Spring Integration headers amqpHeaderMapper.setRequestHeaderNames("X-Custom-Message-ID"); return amqpHeaderMapper; })) .log(Level.INFO) .transform(new ObjectToJsonNodeTransformer(objectMapper)) .route(webhookInboundRouter) .get(); }
A Note on Your Fix
Great call setting defaultRequeueRejected(false)—this prevents failed messages from being requeued indefinitely, which is the most common fix for that infinite loop issue. If you want to handle failed messages gracefully beyond just dropping them, consider configuring a dead-letter exchange (DLX) and dead-letter queue (DLQ) for your main queue. This way, failed messages are routed to the DLQ instead of being discarded, letting you inspect and reprocess them later.
内容的提问来源于stack exchange,提问作者craftingjava

