From c7f009f544564314b34bf2bc69c6a9fb48631ab5 Mon Sep 17 00:00:00 2001 From: Marius Bogoevici Date: Fri, 17 Jul 2015 12:11:52 -0400 Subject: [PATCH] INTEXT-180: Fix ListenerContainer fetch pool size JIRI: https://jira.spring.io/browse/INTEXT-180 --- .../kafka/listener/KafkaMessageListenerContainer.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/listener/KafkaMessageListenerContainer.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/listener/KafkaMessageListenerContainer.java index 1440a3f62f..7a4d51ba33 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/listener/KafkaMessageListenerContainer.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/listener/KafkaMessageListenerContainer.java @@ -306,7 +306,7 @@ public class KafkaMessageListenerContainer implements SmartLifecycle { partitionsByBrokerMap.clear(); partitionsByBrokerMap.putAll(partitionsAsList.groupBy(getLeader)); if (fetchTaskExecutor == null) { - fetchTaskExecutor = Executors.newFixedThreadPool(partitionsByBrokerMap.size()); + fetchTaskExecutor = Executors.newFixedThreadPool(partitionsByBrokerMap.keysView().size()); } partitionsByBrokerMap.forEachKey(launchFetchTask); }