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 {