From e18cff39a6dfb1bbb06af455ae0015fb765c46a9 Mon Sep 17 00:00:00 2001 From: Marius Bogoevici Date: Thu, 7 Apr 2016 12:13:00 -0400 Subject: [PATCH] Make sure that raw tests exercise raw mode --- .../binder/kafka/RawModeKafkaBinderTests.java | 31 +++++++++++++------ 1 file changed, 22 insertions(+), 9 deletions(-) diff --git a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/RawModeKafkaBinderTests.java b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/RawModeKafkaBinderTests.java index 3f2910392..5f88e7009 100644 --- a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/RawModeKafkaBinderTests.java +++ b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/RawModeKafkaBinderTests.java @@ -26,13 +26,11 @@ import static org.junit.Assert.assertThat; import java.util.Arrays; -import org.junit.Ignore; -import org.junit.Test; - import org.springframework.cloud.stream.binder.BinderHeaders; import org.springframework.cloud.stream.binder.Binding; import org.springframework.cloud.stream.binder.ExtendedConsumerProperties; import org.springframework.cloud.stream.binder.ExtendedProducerProperties; +import org.springframework.cloud.stream.binder.HeaderMode; import org.springframework.integration.IntegrationMessageHeaderAccessor; import org.springframework.integration.channel.DirectChannel; import org.springframework.integration.channel.QueueChannel; @@ -42,6 +40,9 @@ import org.springframework.messaging.Message; import org.springframework.messaging.MessageChannel; import org.springframework.messaging.support.GenericMessage; +import org.junit.Ignore; +import org.junit.Test; + /** * @author Marius Bogoevici * @author David Turanski @@ -55,6 +56,7 @@ public class RawModeKafkaBinderTests extends KafkaBinderTests { public void testPartitionedModuleJava() throws Exception { KafkaTestBinder binder = getBinder(); ExtendedProducerProperties properties = createProducerProperties(); + properties.setHeaderMode(HeaderMode.raw); properties.setPartitionKeyExtractorClass(RawKafkaPartitionTestSupport.class); properties.setPartitionSelectorClass(RawKafkaPartitionTestSupport.class); properties.setPartitionCount(3); @@ -68,6 +70,7 @@ public class RawModeKafkaBinderTests extends KafkaBinderTests { consumerProperties.setInstanceCount(3); consumerProperties.setInstanceIndex(0); consumerProperties.setPartitioned(true); + consumerProperties.setHeaderMode(HeaderMode.raw); QueueChannel input0 = new QueueChannel(); input0.setBeanName("test.input0J"); Binding input0Binding = binder.bindConsumer("partJ.0", "test", input0, consumerProperties); @@ -111,6 +114,7 @@ public class RawModeKafkaBinderTests extends KafkaBinderTests { properties.setPartitionKeyExpression(spelExpressionParser.parseExpression("payload[0]")); properties.setPartitionSelectorExpression(spelExpressionParser.parseExpression("hashCode()")); properties.setPartitionCount(3); + properties.setHeaderMode(HeaderMode.raw); DirectChannel output = new DirectChannel(); output.setBeanName("test.output"); @@ -128,6 +132,7 @@ public class RawModeKafkaBinderTests extends KafkaBinderTests { consumerProperties.setInstanceIndex(0); consumerProperties.setInstanceCount(3); consumerProperties.setPartitioned(true); + consumerProperties.setHeaderMode(HeaderMode.raw); QueueChannel input0 = new QueueChannel(); input0.setBeanName("test.input0S"); Binding input0Binding = binder.bindConsumer("part.0", "test", input0, consumerProperties); @@ -175,8 +180,12 @@ public class RawModeKafkaBinderTests extends KafkaBinderTests { KafkaTestBinder binder = getBinder(); DirectChannel moduleOutputChannel = new DirectChannel(); QueueChannel moduleInputChannel = new QueueChannel(); - Binding producerBinding = binder.bindProducer("foo.0", moduleOutputChannel, createProducerProperties()); - Binding consumerBinding = binder.bindConsumer("foo.0", "test", moduleInputChannel, createConsumerProperties()); + ExtendedProducerProperties producerProperties = createProducerProperties(); + producerProperties.setHeaderMode(HeaderMode.raw); + Binding producerBinding = binder.bindProducer("foo.0", moduleOutputChannel, producerProperties); + ExtendedConsumerProperties consumerProperties = createConsumerProperties(); + consumerProperties.setHeaderMode(HeaderMode.raw); + Binding consumerBinding = binder.bindConsumer("foo.0", "test", moduleInputChannel, consumerProperties); Message message = MessageBuilder.withPayload("foo".getBytes()).build(); // Let the consumer actually bind to the producer before sending a msg binderBindUnbindLatency(); @@ -204,14 +213,18 @@ public class RawModeKafkaBinderTests extends KafkaBinderTests { QueueChannel module1InputChannel = new QueueChannel(); QueueChannel module2InputChannel = new QueueChannel(); QueueChannel module3InputChannel = new QueueChannel(); - Binding producerBinding = binder.bindProducer("baz.0", moduleOutputChannel, createProducerProperties()); - Binding input1Binding = binder.bindConsumer("baz.0", "test", module1InputChannel, createConsumerProperties()); + ExtendedProducerProperties producerProperties = createProducerProperties(); + producerProperties.setHeaderMode(HeaderMode.raw); + Binding producerBinding = binder.bindProducer("baz.0", moduleOutputChannel, producerProperties); + ExtendedConsumerProperties consumerProperties = createConsumerProperties(); + consumerProperties.setHeaderMode(HeaderMode.raw); + Binding input1Binding = binder.bindConsumer("baz.0", "test", module1InputChannel, consumerProperties); // A new module is using the tap as an input channel String fooTapName = "baz.0"; - Binding input2Binding = binder.bindConsumer(fooTapName, "tap1", module2InputChannel, createConsumerProperties()); + Binding input2Binding = binder.bindConsumer(fooTapName, "tap1", module2InputChannel, consumerProperties); // Another new module is using tap as an input channel String barTapName = "baz.0"; - Binding input3Binding = binder.bindConsumer(barTapName, "tap2", module3InputChannel, createConsumerProperties()); + Binding input3Binding = binder.bindConsumer(barTapName, "tap2", module3InputChannel, consumerProperties); Message message = MessageBuilder.withPayload("foo".getBytes()).build(); boolean success = false;