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
This commit is contained in:
committed by
Oleg Zhurakousky
parent
81f4a861c5
commit
bf1c366d9b
@@ -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<KafkaProducerProperties> properties = createProducerProperties();
|
||||
properties.setPartitionCount(6);
|
||||
|
||||
DirectChannel output = createBindableChannel("output", createProducerBindingProperties(properties));
|
||||
output.setBeanName("test.output");
|
||||
Binding<MessageChannel> outputBinding = binder.bindProducer("partNative.raw.0", output, properties);
|
||||
|
||||
ExtendedConsumerProperties<KafkaConsumerProperties> consumerProperties = createConsumerProperties();
|
||||
QueueChannel input0 = new QueueChannel();
|
||||
input0.setBeanName("test.inputNative");
|
||||
Binding<MessageChannel> 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 {
|
||||
|
||||
Reference in New Issue
Block a user