如何复用Spring AMQP自定义MessageConverter适配Spring Cloud Stream?
Great question! Let’s break this down step by step—since there’s no direct out-of-the-box converter between the two MessageConverter types, we can absolutely reuse your existing Spring AMQP custom converter by building an adapter that bridges org.springframework.amqp.core.Message and org.springframework.messaging.Message.
Understanding the Core Difference
First, let’s clarify why direct conversion isn’t possible:
- Spring AMQP’s
MessageConverterhandles conversions between business objects andorg.springframework.amqp.core.Message(AMQP-specific message format withMessagePropertiesfor AMQP headers likecontentType,deliveryMode, etc.). - Spring Cloud Stream’s
MessageConverterhandles conversions between business objects andorg.springframework.messaging.Message(a generic messaging abstraction withMessageHeadersthat includes both AMQP-specific and Stream-specific metadata).
They operate on different message abstractions, so we need an adapter to translate between the two while reusing your existing AMQP converter logic.
Solution: Build an Adapter Converter
The approach is to create a custom Spring Cloud Stream MessageConverter that acts as a wrapper around your existing AMQP converter. This adapter will:
- Convert incoming
org.springframework.messaging.Messagetoorg.springframework.amqp.core.Message, then delegate to your AMQP converter to get the business object. - Convert outgoing business objects to
org.springframework.amqp.core.Messagevia your AMQP converter, then map that toorg.springframework.messaging.Messagefor Stream processing.
Example Implementation
Here’s a concrete example of the adapter converter:
import org.springframework.amqp.core.Message; import org.springframework.amqp.core.MessageProperties; import org.springframework.amqp.support.converter.MessageConverter; import org.springframework.messaging.MessageHeaders; import org.springframework.messaging.converter.AbstractMessageConverter; import org.springframework.util.MimeType; import java.util.HashMap; import java.util.Map; public class AmqpAdaptingStreamConverter extends AbstractMessageConverter { private final org.springframework.amqp.support.converter.MessageConverter amqpConverter; // Inject your existing AMQP custom converter public AmqpAdaptingStreamConverter(MessageConverter amqpConverter) { super(MimeType.valueOf("application/json")); // Adjust to your actual content type this.amqpConverter = amqpConverter; } // Handle incoming messages (Stream Message -> Business Object) @Override protected Object convertFromInternal(org.springframework.messaging.Message<?> streamMessage, Class<?> targetClass, Object conversionHint) { // Map Stream Message headers to AMQP MessageProperties MessageProperties amqpProps = new MessageProperties(); streamMessage.getHeaders().forEach((key, value) -> { // Skip Stream-specific headers (prefixed with "sc-") and let AMQP handle content type separately if (!key.startsWith("sc-") && !key.equals(MessageHeaders.CONTENT_TYPE)) { amqpProps.setHeader(key, value); } }); if (streamMessage.getHeaders().containsKey(MessageHeaders.CONTENT_TYPE)) { amqpProps.setContentType(streamMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE).toString()); } // Create AMQP Message from Stream Message payload and mapped properties Message amqpMessage = new Message((byte[]) streamMessage.getPayload(), amqpProps); // Delegate to your existing AMQP converter to get the business object return amqpConverter.fromMessage(amqpMessage, targetClass); } // Handle outgoing messages (Business Object -> Stream Message) @Override protected Object convertToInternal(Object payload, org.springframework.messaging.MessageHeaders streamHeaders, Object conversionHint) { // Use your AMQP converter to turn the business object into an AMQP Message Message amqpMessage = amqpConverter.toMessage(payload, null); // Merge AMQP headers with existing Stream headers Map<String, Object> mergedHeaders = new HashMap<>(streamHeaders); mergedHeaders.putAll(amqpMessage.getMessageProperties().getHeaders()); mergedHeaders.put(MessageHeaders.CONTENT_TYPE, amqpMessage.getMessageProperties().getContentType()); // Return the AMQP message body; headers will be handled by the Stream framework return amqpMessage.getBody(); } }
Configure the Adapter in Spring Cloud Stream
Register the adapter as a Spring bean so Stream uses it instead of the default converter:
import org.springframework.amqp.support.converter.MessageConverter; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.messaging.converter.MessageConverter; @Configuration public class StreamConverterConfig { // Assume your custom AMQP converter is already a bean in the context @Bean public MessageConverter streamMessageConverter(MessageConverter customAmqpConverter) { return new AmqpAdaptingStreamConverter(customAmqpConverter); } }
Key Considerations
- Header Mapping: Adjust the header filtering logic to match your business needs. Some Stream-specific headers (like
sc-destination) don’t need to be passed to AMQP, while certain AMQP headers (likedeliveryMode) might need special handling. - Content Type: Ensure the content type is correctly mapped between AMQP’s
contentTypeand Stream’sMessageHeaders.CONTENT_TYPE. - Testing: Validate both incoming and outgoing message flows to confirm that the body and headers are converted correctly, and that your existing AMQP converter logic behaves as expected.
This approach lets you fully reuse your existing Spring AMQP converter without rewriting any core conversion logic—perfect for migrating to Spring Cloud Stream!
内容的提问来源于stack exchange,提问作者Johan

