From bf1c366d9bd296667c94153d952d370bbb3d3782 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Wed, 21 Nov 2018 13:14:53 -0500 Subject: [PATCH] GH-502: Add test for native partitioning Resolves https://github.com/spring-cloud/spring-cloud-stream-binder-kafka/issues/502 Before applying the fix for, https://github.com/spring-cloud/spring-cloud-stream/issues/1531 failed with: ``` org.springframework.messaging.MessageDeliveryException: failed to send Message to channel 'test.output'; nested exception is java.lang.IllegalArgumentException: Partition key cannot be null, failedMessage=GenericMessage [payload=byte[3], headers={kafka_partitionId=5, id=3350a823-c876-f7a9-f98b-fdbd2aaa4c12, timestamp=1542823925354}] ... Caused by: java.lang.IllegalArgumentException: Partition key cannot be null at org.springframework.util.Assert.notNull(Assert.java:198) at org.springframework.cloud.stream.binder.PartitionHandler.extractKey(PartitionHandler.java:112) at org.springframework.cloud.stream.binder.PartitionHandler.determinePartition(PartitionHandler.java:93) at org.springframework.cloud.stream.binding.MessageConverterConfigurer$PartitioningInterceptor.preSend(MessageConverterConfigurer.java:381) at org.springframework.integration.channel.AbstractMessageChannel$ChannelInterceptorList.preSend(AbstractMessageChannel.java:589) at org.springframework.integration.channel.AbstractMessageChannel.send(AbstractMessageChannel.java:435) ... 31 more ``` Resolves #503 --- .../stream/binder/kafka/KafkaBinderTests.java | 30 +++++++++++++++++++ 1 file changed, 30 insertions(+) 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 6b3104a8c..035c98b57 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 @@ -35,6 +35,7 @@ import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicReference; import com.fasterxml.jackson.databind.ObjectMapper; + import org.apache.kafka.clients.admin.AdminClient; import org.apache.kafka.clients.admin.AdminClientConfig; import org.apache.kafka.clients.admin.CreateTopicsResult; @@ -1591,6 +1592,35 @@ public class KafkaBinderTests extends outputBinding.unbind(); } + @Test + @SuppressWarnings({ "unchecked", "rawtypes" }) + public void testPartitionedNative() throws Exception { + Binder binder = getBinder(); + ExtendedProducerProperties properties = createProducerProperties(); + properties.setPartitionCount(6); + + DirectChannel output = createBindableChannel("output", createProducerBindingProperties(properties)); + output.setBeanName("test.output"); + Binding outputBinding = binder.bindProducer("partNative.raw.0", output, properties); + + ExtendedConsumerProperties consumerProperties = createConsumerProperties(); + QueueChannel input0 = new QueueChannel(); + input0.setBeanName("test.inputNative"); + Binding inputBinding = binder.bindConsumer("partNative.raw.0", "test", input0, consumerProperties); + + output.send( + new GenericMessage<>("foo".getBytes(), Collections.singletonMap(KafkaHeaders.PARTITION_ID, 5))); + + Message received = receive(input0); + assertThat(received).isNotNull(); + + assertThat(received.getPayload()).isEqualTo("foo".getBytes()); + assertThat(received.getHeaders().get(KafkaHeaders.RECEIVED_PARTITION_ID)).isEqualTo(5); + + inputBinding.unbind(); + outputBinding.unbind(); + } + @Test @SuppressWarnings({"unchecked", "rawtypes"}) public void testSendAndReceiveWithRawMode() throws Exception {