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