Spring Integration Kafka 3.0.1多消费者运行异常求助
Let's walk through the possible issues and fixes for your multi-consumer setup problem:
1. Make Sure Concurrency is Applied to Spring Integration's Kafka Adapter
The global spring.kafka.listener.concurrency setting works seamlessly for @KafkaListener annotations, but Spring Integration's MessageDrivenChannelAdapter (the core component handling Kafka message listening in integration flows) doesn't automatically inherit this global config.
If you're using Java DSL to define your integration flow, you need to explicitly set the concurrency property on the adapter:
@Bean public IntegrationFlow kafkaInboundFlow(ConsumerFactory<String, String> consumerFactory) { return IntegrationFlows.from( Kafka.messageDrivenChannelAdapter(consumerFactory, "your-10-partition-topic") .concurrency(4) // Explicitly set concurrency here ) .handle(message -> { // Your message processing logic }) .get(); }
For XML configuration, add the concurrency attribute to the <int-kafka:message-driven-channel-adapter> element.
2. Verify Partition Allocation via Kafka CLI
First, confirm if your partitions are actually being split across consumers. Run this command to inspect your consumer group's partition assignment:
kafka-consumer-groups.sh --bootstrap-server <your-kafka-host>:9092 --describe --group test-consumer
- If all 10 partitions are assigned to one consumer: Your concurrency config isn't being applied (go back to step 1 to fix the adapter setup).
- If partitions are spread across 4 consumers (e.g., 3 partitions for two consumers, 2 for the other two): Concurrency is working, but your message processing logic might be blocking/throttling, making it seem like only one consumer is active.
3. Check Version Compatibility & Container Setup
spring-integration-kafka 3.0.1 pairs with Spring Kafka 2.3.x. While this version supports concurrency, keep these points in mind:
- Ensure you're using
ConcurrentKafkaListenerContainerFactoryif you're wiring containers manually. A plainKafkaListenerContainerFactorywon't support concurrent consumers. - Avoid conflicting configs: If you've set
max.poll.recordsto an extremely high value, a single consumer might process all messages before others get a turn, even with concurrency enabled.
4. Rule Out Accidental Configuration Overrides
Double-check your configuration classes for any code that might override the concurrency setting. For example, if you're customizing the KafkaMessageDrivenChannelAdapter bean and explicitly set concurrency(1) somewhere, that would overwrite your YAML config.
5. Consider Version-Specific Bugs
While uncommon, there could be a bug in spring-integration-kafka 3.0.1 affecting concurrency. You can check the project's issue tracker for reports of similar problems in this version. If you find a matching bug, upgrading to a newer minor version (like 3.0.5, the latest in the 3.0.x line) or moving to 3.1.x might resolve the issue.
内容的提问来源于stack exchange,提问作者mrusinak

