From a52cda705718fba7cf0756a46550ed64aa3ae221 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Wed, 12 Jun 2019 18:20:37 +0200 Subject: [PATCH] Binder changes related to GH-1729 addressed PR comment Resolves #250 --- pom.xml | 2 +- spring-cloud-stream-binder-rabbit/pom.xml | 8 -------- .../stream/binder/rabbit/RabbitMessageChannelBinder.java | 8 ++++---- 3 files changed, 5 insertions(+), 13 deletions(-) diff --git a/pom.xml b/pom.xml index b50d704c6..6c0e414a5 100644 --- a/pom.xml +++ b/pom.xml @@ -7,7 +7,7 @@ org.springframework.cloud spring-cloud-build - 2.2.0.M2 + 2.2.0.BUILD-SNAPSHOT diff --git a/spring-cloud-stream-binder-rabbit/pom.xml b/spring-cloud-stream-binder-rabbit/pom.xml index a62c86c15..df8f4dd8a 100644 --- a/spring-cloud-stream-binder-rabbit/pom.xml +++ b/spring-cloud-stream-binder-rabbit/pom.xml @@ -53,14 +53,6 @@ org.springframework.boot spring-boot-starter-amqp - - org.springframework.integration - spring-integration-amqp - - - org.springframework.integration - spring-integration-core - org.springframework.integration spring-integration-jmx diff --git a/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java b/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java index 40496d3b2..11a5e95f4 100644 --- a/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java +++ b/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java @@ -494,6 +494,7 @@ public class RabbitMessageChannelBinder extends AmqpInboundChannelAdapter adapter = new AmqpInboundChannelAdapter( listenerContainer); + adapter.setBindSourceMessage(true); adapter.setBeanFactory(this.getBeanFactory()); adapter.setBeanName("inbound." + destination); DefaultAmqpHeaderMapper mapper = DefaultAmqpHeaderMapper.inboundMapper(); @@ -609,8 +610,8 @@ public class RabbitMessageChannelBinder extends public void handleMessage( org.springframework.messaging.Message message) throws MessagingException { - Message amqpMessage = (Message) message.getHeaders() - .get(AmqpMessageHeaderErrorMessageStrategy.AMQP_RAW_MESSAGE); + Message amqpMessage = StaticMessageHeaderAccessor.getSourceData(message); + if (!(message instanceof ErrorMessage)) { logger.error("Expected an ErrorMessage, not a " + message.getClass().toString() + " for: " + message); @@ -695,8 +696,7 @@ public class RabbitMessageChannelBinder extends public void handleMessage( org.springframework.messaging.Message message) throws MessagingException { - Message amqpMessage = (Message) message.getHeaders() - .get(AmqpMessageHeaderErrorMessageStrategy.AMQP_RAW_MESSAGE); + Message amqpMessage = StaticMessageHeaderAccessor.getSourceData(message); /* * NOTE: The following IF and subsequent ELSE IF should never happen * under normal interaction and it should always go to the last ELSE