I set up a Spring Integration flow to process a topic having 3 partitions and set the listener container's concurrency to 3. As expected, I see three threads processing batches from all 3 partitions. However, I see that in some cases, one of the listener threads may process a single batch containing messages from multiple partitions. My data is partitioned in kafka by an id so that it may be processed concurrently with other ids, but not with the same ids on another thread (which is what I was surprised to observe is happening). I thought from reading the docs that each thread would be assigned a partition. I'm using a KafkaMessageDrivenChannelAdapter like this:
private static final Class<List<MyEvent>> payloadClass = (Class<List<MyEvent>>)(Class) List.class;
public KafkaMessageDrivenChannelAdapterSpec.KafkaMessageDrivenChannelAdapterListenerContainerSpec<String, MyEvent> myChannelAdapterSpec() {
return Kafka.messageDrivenChannelAdapter(tstatEventConsumerFactory(),
KafkaMessageDrivenChannelAdapter.ListenerMode.batch, "my-topic") //3 partitions
.configureListenerContainer(c -> {
c.ackMode(ContainerProperties.AckMode.BATCH);
c.id(_ID);
c.concurrency(3);
RecoveringBatchErrorHandler errorHandler = new RecoveringBatchErrorHandler(
(record, exception) -> log.error("failed to handle record at offset {}: {}",
record.offset(), record.value(), exception),
new FixedBackOff(FixedBackOff.DEFAULT_INTERVAL, 2)
);
c.errorHandler(errorHandler);
});
}
@Bean
public IntegrationFlow myIntegrationFlow() {
return IntegrationFlows.from(myChannelAdapterSpec())
.handle(payloadClass, (payload, headers) -> {
service.performSink(payload);
return null;
})
.get();
}
How do I set this up so that each listener container thread only processes messages from one partition?