diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/inbound/AmqpMessageSource.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/inbound/AmqpMessageSource.java index 9cccdf2a92..5401652416 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/inbound/AmqpMessageSource.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/inbound/AmqpMessageSource.java @@ -169,6 +169,7 @@ public class AmqpMessageSource extends AbstractMessageSource { .createCallback(new AmqpAckInfo(connection, channel, this.transacted, resp)); MessageProperties messageProperties = this.propertiesConverter.toMessageProperties(resp.getProps(), resp.getEnvelope(), StandardCharsets.UTF_8.name()); + messageProperties.setConsumerQueue(this.queue); Map headers = this.headerMapper.toHeadersFromRequest(messageProperties); org.springframework.amqp.core.Message amqpMessage = new org.springframework.amqp.core.Message(resp.getBody(), messageProperties); Object payload = this.messageConverter.fromMessage(amqpMessage); diff --git a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/inbound/AmqpMessageSourceTests.java b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/inbound/AmqpMessageSourceTests.java index 8cfba0711c..54cdd0f98a 100644 --- a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/inbound/AmqpMessageSourceTests.java +++ b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/inbound/AmqpMessageSourceTests.java @@ -16,6 +16,7 @@ package org.springframework.integration.amqp.inbound; +import static org.hamcrest.Matchers.equalTo; import static org.hamcrest.Matchers.instanceOf; import static org.junit.Assert.assertThat; import static org.mockito.ArgumentMatchers.anyString; @@ -32,6 +33,7 @@ import java.util.concurrent.TimeoutException; import org.junit.Test; import org.springframework.amqp.rabbit.connection.CachingConnectionFactory; +import org.springframework.amqp.support.AmqpHeaders; import org.springframework.integration.amqp.support.AmqpMessageHeaderErrorMessageStrategy; import org.springframework.integration.support.AcknowledgmentCallback.Status; import org.springframework.integration.support.StaticMessageHeaderAccessor; @@ -72,6 +74,7 @@ public class AmqpMessageSourceTests { Message received = source.receive(); assertThat(received.getHeaders().get(AmqpMessageHeaderErrorMessageStrategy.AMQP_RAW_MESSAGE), instanceOf(org.springframework.amqp.core.Message.class)); + assertThat(received.getHeaders().get(AmqpHeaders.CONSUMER_QUEUE), equalTo("foo")); // make sure channel is not cached org.springframework.amqp.rabbit.connection.Connection conn = ccf.createConnection(); Channel notCached = conn.createChannel(false); // should not have been "closed"