diff --git a/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/MessageChannelConfigurerTests.java b/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/MessageChannelConfigurerTests.java index 0f41ffa71..2a16910c0 100644 --- a/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/MessageChannelConfigurerTests.java +++ b/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/MessageChannelConfigurerTests.java @@ -28,10 +28,12 @@ import org.springframework.beans.DirectFieldAccessor; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; import org.springframework.boot.test.context.SpringBootTest; -import org.springframework.cloud.stream.annotation.Bindings; import org.springframework.cloud.stream.annotation.EnableBinding; +import org.springframework.cloud.stream.binder.BinderHeaders; import org.springframework.cloud.stream.converter.CompositeMessageConverterFactory; import org.springframework.cloud.stream.messaging.Sink; +import org.springframework.cloud.stream.messaging.Source; +import org.springframework.cloud.stream.test.binder.MessageCollector; import org.springframework.context.annotation.PropertySource; import org.springframework.integration.support.MessageBuilder; import org.springframework.messaging.Message; @@ -49,19 +51,24 @@ import static org.assertj.core.api.Assertions.assertThat; * @author Ilayaperumal Gopinathan */ @RunWith(SpringJUnit4ClassRunner.class) -@SpringBootTest(classes = {MessageChannelConfigurerTests.TestSink.class}) +@SpringBootTest(classes = {MessageChannelConfigurerTests.TestSink.class, MessageChannelConfigurerTests.TestSource.class}) public class MessageChannelConfigurerTests { @Autowired - @Bindings(TestSink.class) private Sink testSink; + @Autowired + private Source testSource; + @Autowired private CompositeMessageConverterFactory messageConverterFactory; @Autowired private ObjectMapper objectMapper; + @Autowired + private MessageCollector messageCollector; + @Test public void testMessageConverterConfigurer() throws Exception { final CountDownLatch latch = new CountDownLatch(1); @@ -76,7 +83,7 @@ public class MessageChannelConfigurerTests { }; testSink.input().subscribe(messageHandler); testSink.input().send(MessageBuilder.withPayload("{\"message\":\"Hi\"}").build()); - assertThat(latch.await(10, TimeUnit.SECONDS)); + assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue(); testSink.input().unsubscribe(messageHandler); } @@ -95,10 +102,31 @@ public class MessageChannelConfigurerTests { } } + @Test + public void testPartitionHeader() throws Exception { + testSource.output().send(MessageBuilder.withPayload("{\"message\":\"Hi\"}").build()); + Message message = messageCollector.forChannel(testSource.output()).poll(1, TimeUnit.SECONDS); + assertThat(message.getHeaders().get(BinderHeaders.PARTITION_HEADER).equals(0)); + } + + @Test + public void testPartitionHeaderWithExplicitHeader() throws Exception { + testSource.output().send(MessageBuilder.withPayload("{\"message\":\"Hi\"}").setHeader(BinderHeaders.PARTITION_HEADER, "customerId-123").build()); + Message message = messageCollector.forChannel(testSource.output()).poll(1, TimeUnit.SECONDS); + assertThat(message.getHeaders().get(BinderHeaders.PARTITION_HEADER).equals("customerId-123")); + } + @EnableBinding(Sink.class) @EnableAutoConfiguration @PropertySource("classpath:/org/springframework/cloud/stream/config/channel/sink-channel-configurers.properties") public static class TestSink { } + + @EnableBinding(Source.class) + @EnableAutoConfiguration + @PropertySource("classpath:/org/springframework/cloud/stream/config/channel/partitioned-configurers.properties") + public static class TestSource { + + } } diff --git a/spring-cloud-stream-integration-tests/src/test/resources/org/springframework/cloud/stream/config/channel/partitioned-configurers.properties b/spring-cloud-stream-integration-tests/src/test/resources/org/springframework/cloud/stream/config/channel/partitioned-configurers.properties new file mode 100644 index 000000000..9b7e4bd95 --- /dev/null +++ b/spring-cloud-stream-integration-tests/src/test/resources/org/springframework/cloud/stream/config/channel/partitioned-configurers.properties @@ -0,0 +1,3 @@ +spring.cloud.stream.bindings.output.destination=partOut +spring.cloud.stream.bindings.output.producer.partitionKeyExpression=payload +spring.cloud.stream.bindings.output.producer.partitionCount=3