From 005ec51d8b07d6739fed128dfaa03532fd8a81b7 Mon Sep 17 00:00:00 2001 From: Marius Bogoevici Date: Tue, 4 Apr 2017 12:09:49 -0400 Subject: [PATCH] Improve isolation of Kafka tests Fix #120 --- .../stream/binder/kafka/KafkaBinderTests.java | 62 +++++++++---------- 1 file changed, 31 insertions(+), 31 deletions(-) diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java index 03cf4ca2f..c9c40cdcf 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java @@ -374,10 +374,10 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests producerProperties = createProducerProperties(); producerProperties.getExtension().setCompressionType( KafkaProducerProperties.CompressionType.valueOf(codec.toString())); - Binding producerBinding = binder.bindProducer("foo.0", moduleOutputChannel, + Binding producerBinding = binder.bindProducer("testCompression", moduleOutputChannel, producerProperties); ExtendedConsumerProperties consumerProperties = createConsumerProperties(); - Binding consumerBinding = binder.bindConsumer("foo.0", "test", moduleInputChannel, + Binding consumerBinding = binder.bindConsumer("testCompression", "test", moduleInputChannel, consumerProperties); Message message = org.springframework.integration.support.MessageBuilder.withPayload(testPayload) .build(); @@ -732,13 +732,13 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests producerBinding = binder.bindProducer("foo.x", moduleOutputChannel, + Binding producerBinding = binder.bindProducer("testManualAckSucceedsWhenAutoCommitOffsetIsTurnedOff", moduleOutputChannel, createProducerProperties()); ExtendedConsumerProperties consumerProperties = createConsumerProperties(); consumerProperties.getExtension().setAutoCommitOffset(false); - Binding consumerBinding = binder.bindConsumer("foo.x", "test", moduleInputChannel, + Binding consumerBinding = binder.bindConsumer("testManualAckSucceedsWhenAutoCommitOffsetIsTurnedOff", "test", moduleInputChannel, consumerProperties); String testPayload1 = "foo" + UUID.randomUUID().toString(); @@ -774,12 +774,12 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests producerBinding = binder.bindProducer("foo.x", moduleOutputChannel, + Binding producerBinding = binder.bindProducer("testManualAckIsNotPossibleWhenAutoCommitOffsetIsEnabledOnTheBinder", moduleOutputChannel, createProducerProperties()); ExtendedConsumerProperties consumerProperties = createConsumerProperties(); - Binding consumerBinding = binder.bindConsumer("foo.x", "test", moduleInputChannel, + Binding consumerBinding = binder.bindConsumer("testManualAckIsNotPossibleWhenAutoCommitOffsetIsEnabledOnTheBinder", "test", moduleInputChannel, consumerProperties); String testPayload1 = "foo" + UUID.randomUUID().toString(); @@ -1420,7 +1420,7 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests outputBinding = binder.bindProducer("partJ.0", output, properties); + Binding outputBinding = binder.bindProducer("partJ.raw.0", output, properties); ExtendedConsumerProperties consumerProperties = createConsumerProperties(); consumerProperties.setConcurrency(2); @@ -1431,15 +1431,15 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests input0Binding = binder.bindConsumer("partJ.0", "test", input0, consumerProperties); + Binding input0Binding = binder.bindConsumer("partJ.raw.0", "test", input0, consumerProperties); consumerProperties.setInstanceIndex(1); QueueChannel input1 = new QueueChannel(); input1.setBeanName("test.input1J"); - Binding input1Binding = binder.bindConsumer("partJ.0", "test", input1, consumerProperties); + Binding input1Binding = binder.bindConsumer("partJ.raw.0", "test", input1, consumerProperties); consumerProperties.setInstanceIndex(2); QueueChannel input2 = new QueueChannel(); input2.setBeanName("test.input2J"); - Binding input2Binding = binder.bindConsumer("partJ.0", "test", input2, consumerProperties); + Binding input2Binding = binder.bindConsumer("partJ.raw.0", "test", input2, consumerProperties); output.send(new GenericMessage<>(new byte[] {(byte) 0})); output.send(new GenericMessage<>(new byte[] {(byte) 1})); @@ -1473,11 +1473,11 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests outputBinding = binder.bindProducer("part.0", output, properties); + Binding outputBinding = binder.bindProducer("part.raw.0", output, properties); try { Object endpoint = extractEndpoint(outputBinding); assertThat(getEndpointRouting(endpoint)) - .contains(getExpectedRoutingBaseDestination("part.0", "test") + "-' + headers['partition']"); + .contains(getExpectedRoutingBaseDestination("part.raw.0", "test") + "-' + headers['partition']"); } catch (UnsupportedOperationException ignored) { } @@ -1492,15 +1492,15 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests input0Binding = binder.bindConsumer("part.0", "test", input0, consumerProperties); + Binding input0Binding = binder.bindConsumer("part.raw.0", "test", input0, consumerProperties); consumerProperties.setInstanceIndex(1); QueueChannel input1 = new QueueChannel(); input1.setBeanName("test.input1S"); - Binding input1Binding = binder.bindConsumer("part.0", "test", input1, consumerProperties); + Binding input1Binding = binder.bindConsumer("part.raw.0", "test", input1, consumerProperties); consumerProperties.setInstanceIndex(2); QueueChannel input2 = new QueueChannel(); input2.setBeanName("test.input2S"); - Binding input2Binding = binder.bindConsumer("part.0", "test", input2, consumerProperties); + Binding input2Binding = binder.bindConsumer("part.raw.0", "test", input2, consumerProperties); Message message2 = org.springframework.integration.support.MessageBuilder.withPayload(new byte[] {2}) .setHeader(IntegrationMessageHeaderAccessor.CORRELATION_ID, "kafkaBinderTestCommonsDelegate") @@ -1531,19 +1531,19 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests producerProperties = createProducerProperties(); producerProperties.setHeaderMode(HeaderMode.raw); - Binding producerBinding = binder.bindProducer("0", moduleOutputChannel, + Binding producerBinding = binder.bindProducer("raw.0", moduleOutputChannel, producerProperties); ExtendedConsumerProperties consumerProperties = createConsumerProperties(); consumerProperties.setHeaderMode(HeaderMode.raw); - Binding consumerBinding = binder.bindConsumer("0", "test", moduleInputChannel, + Binding consumerBinding = binder.bindConsumer("raw.0", "test", moduleInputChannel, consumerProperties); - Message message = org.springframework.integration.support.MessageBuilder.withPayload("kafkaBinderTestCommonsDelegate".getBytes()).build(); + Message message = org.springframework.integration.support.MessageBuilder.withPayload("testSendAndReceiveWithRawMode".getBytes()).build(); // Let the consumer actually bind to the producer before sending a msg binderBindUnbindLatency(); moduleOutputChannel.send(message); Message inbound = receive(moduleInputChannel); assertThat(inbound).isNotNull(); - assertThat(new String((byte[]) inbound.getPayload())).isEqualTo("kafkaBinderTestCommonsDelegate"); + assertThat(new String((byte[]) inbound.getPayload())).isEqualTo("testSendAndReceiveWithRawMode"); producerBinding.unbind(); consumerBinding.unbind(); } @@ -1556,19 +1556,19 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests producerProperties = createProducerProperties(); producerProperties.setHeaderMode(HeaderMode.raw); - Binding producerBinding = binder.bindProducer("0", moduleOutputChannel, + Binding producerBinding = binder.bindProducer("raw.string.0", moduleOutputChannel, producerProperties); ExtendedConsumerProperties consumerProperties = createConsumerProperties(); consumerProperties.setHeaderMode(HeaderMode.raw); - Binding consumerBinding = binder.bindConsumer("0", "test", moduleInputChannel, + Binding consumerBinding = binder.bindConsumer("raw.string.0", "test", moduleInputChannel, consumerProperties); - Message message = org.springframework.integration.support.MessageBuilder.withPayload("kafkaBinderTestCommonsDelegate").build(); + Message message = org.springframework.integration.support.MessageBuilder.withPayload("testSendAndReceiveWithRawModeAndStringPayload").build(); // Let the consumer actually bind to the producer before sending a msg binderBindUnbindLatency(); moduleOutputChannel.send(message); Message inbound = receive(moduleInputChannel); assertThat(inbound).isNotNull(); - assertThat(new String((byte[]) inbound.getPayload())).isEqualTo("kafkaBinderTestCommonsDelegate"); + assertThat(new String((byte[]) inbound.getPayload())).isEqualTo("testSendAndReceiveWithRawModeAndStringPayload"); producerBinding.unbind(); consumerBinding.unbind(); } @@ -1584,30 +1584,30 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests producerProperties = createProducerProperties(); producerProperties.setHeaderMode(HeaderMode.raw); - Binding producerBinding = binder.bindProducer("baz.0", moduleOutputChannel, + Binding producerBinding = binder.bindProducer("baz.raw.0", moduleOutputChannel, producerProperties); ExtendedConsumerProperties consumerProperties = createConsumerProperties(); consumerProperties.setHeaderMode(HeaderMode.raw); consumerProperties.getExtension().setAutoRebalanceEnabled(false); - Binding input1Binding = binder.bindConsumer("baz.0", "test", module1InputChannel, + Binding input1Binding = binder.bindConsumer("baz.raw.0", "test", module1InputChannel, consumerProperties); // A new module is using the tap as an input channel - String fooTapName = "baz.0"; + String fooTapName = "baz.raw.0"; Binding input2Binding = binder.bindConsumer(fooTapName, "tap1", module2InputChannel, consumerProperties); // Another new module is using tap as an input channel - String barTapName = "baz.0"; + String barTapName = "baz.raw.0"; Binding input3Binding = binder.bindConsumer(barTapName, "tap2", module3InputChannel, consumerProperties); - Message message = org.springframework.integration.support.MessageBuilder.withPayload("kafkaBinderTestCommonsDelegate".getBytes()).build(); + Message message = org.springframework.integration.support.MessageBuilder.withPayload("testSendAndReceiveWithExplicitConsumerGroupWithRawMode".getBytes()).build(); boolean success = false; boolean retried = false; while (!success) { moduleOutputChannel.send(message); Message inbound = receive(module1InputChannel); assertThat(inbound).isNotNull(); - assertThat(new String((byte[]) inbound.getPayload())).isEqualTo("kafkaBinderTestCommonsDelegate"); + assertThat(new String((byte[]) inbound.getPayload())).isEqualTo("testSendAndReceiveWithExplicitConsumerGroupWithRawMode"); Message tapped1 = receive(module2InputChannel); Message tapped2 = receive(module3InputChannel); @@ -1618,8 +1618,8 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests