多Kafka消费者无法接收消息排查(基于EmbeddedKafka与Stream Binder)
问题:同一消费组下Spring Cloud Stream Kafka Streams消费者仅一个能接收消息
现象描述
使用EmbeddedKafkaBroker和Spring Cloud Stream Kafka Streams时,两个消费者Bean绑定到同一个主题demoTopic,且通过applicationId配置为同一消费组group_id,但只有在spring.cloud.stream.function.definition中排在前面的消费者能收到消息,第二个消费者无消息。日志显示demoTopic仅存在1个分区,尽管在构造EmbeddedKafkaBroker时指定了2个分区,但配置似乎未生效。
相关代码配置
EmbeddedKafkaBroker配置类
@Configuration @Profile({"dev", "test"}) @Slf4j public class EmbeddedKafkaBrokerConfig { private static final String TMP_EMBEDDED_KAFKA_LOGS = String.format("/tmp/embedded-kafka-logs-%1$s/", UUID.randomUUID()); private static final String PORT = "port"; private static final String LOG_DIRS = "log.dirs"; private static final String LISTENERS = "listeners"; private static final Integer KAFKA_PORT = 9092; private static final String LISTENERS_VALUE = "PLAINTEXT://localhost:" + KAFKA_PORT; private static final Integer ZOOKEEPER_PORT = 2181; private EmbeddedKafkaBroker embeddedKafkaBroker; /** * bean for the embeddedKafkaBroker. * * @return local embeddedKafkaBroker */ @Bean @Qualifier("embeddedKafkaBroker") public EmbeddedKafkaBroker embeddedKafkaBroker() { Map<String, String> brokerProperties = new HashMap<>(); brokerProperties.put(LISTENERS, LISTENERS_VALUE); brokerProperties.put(PORT, KAFKA_PORT.toString()); brokerProperties.put(LOG_DIRS, TMP_EMBEDDED_KAFKA_LOGS); this.embeddedKafkaBroker = new EmbeddedKafkaBroker(1, true, 2) .kafkaPorts(KAFKA_PORT) .zkPort(ZOOKEEPER_PORT) .brokerProperties(brokerProperties); return embeddedKafkaBroker; } /** close the embeddedKafkaBroker on destroy. */ @PreDestroy public void preDestroy() { if (embeddedKafkaBroker != null) { log.warn("[EmbeddedKafkaBrokerConfig] destroying kafka broker {}", embeddedKafkaBroker); embeddedKafkaBroker.destroy(); } } }
消息发送模块
DemoController
@RestController @RequestMapping("/v1/demo/") public class DemoController { @Autowired DemoSupplier demoSupplier; @GetMapping("hello") public String helloController(){ demoSupplier.supply(); return "Hello World!"; } }
DemoSupplier
@Component public class DemoSupplier { @Autowired @Qualifier("embeddedKafkaBroker") public EmbeddedKafkaBroker kafkaBroker; @Autowired private KafkaTemplate<String,String> kafkaTemplate; @Value("${demo.topic}") private String topicName; @Bean public KafkaTemplate<String, String> stringKafkaTemplate(){ Map<String, Object> producerConfigs =new HashMap<>(); producerConfigs.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,"localhost:9092"); producerConfigs.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); producerConfigs.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,StringSerializer.class); return new KafkaTemplate<>(new DefaultKafkaProducerFactory<>(producerConfigs)); } public void supply(){ for(int i =0 ;i<100;i++){ kafkaTemplate.send(topicName, "Message:"+i*2); } } }
消费者模块
@Component public class DemoConsumer { @Bean @Qualifier("demoConsumerProcessor") public Consumer<KStream<String, String>> demoConsumerProcessor(){ return input -> input.foreach(((key, value) -> System.out.println(value))); } @Bean @Qualifier("demoConsumerProcessor2") public Consumer<KStream<String, String>> demoConsumerProcessor2(){ return input -> input.foreach(((key, value) -> System.out.println("This is second consumer 2: "+value))); } }
application.properties配置
# =============================== # = Profiles # =============================== spring.profiles.active=dev server.port=8181 # =============================== # = Kafka Topics # =============================== demo.topic=demoTopic object.demo.topic=objectDemoTopic # =============================== # = SPRING CLOUD STREAM # =============================== spring.cloud.stream.bindings.demoConsumerProcessor-in-0.destination=demoTopic spring.cloud.stream.bindings.demoConsumerProcessor2-in-0.destination=demoTopic spring.cloud.stream.function.definition=demoConsumerProcessor,demoConsumerProcessor2 spring.cloud.stream.kafka.streams.binder.functions.demoConsumerProcessor.applicationId=group_id spring.cloud.stream.kafka.streams.binder.functions.demoConsumerProcessor2.applicationId=group_id spring.kafka.producer.key-serializer=org.apache.kafka.common.serialization.StringSerializer spring.kafka.producer.value-serializer=org.apache.kafka.common.serialization.StringSerializer spring.kafka.consumer.value-deserializer=org.apache.kafka.common.serialization.StringDeserializer
关键日志
第一个消费者分区分配日志:
[Consumer clientId=group_id-359878ed-1b41-4cf0-b9b8-6e21e5e1f0fe-StreamThread-1-consumer, groupId=group_id] Updating assignment with Assigned partitions: [demoTopic-0] Current owned partitions: [] Added partitions (assigned - owned): [demoTopic-0] Revoked partitions (owned - assigned): []
第二个消费者分区分配日志:
Consumer clientId=group_id-4dce1ba5-7d97-4c18-92c3-cb79dab271b5-StreamThread-1-consumer, groupId=group_id] Updating assignment with Assigned partitions: [] Current owned partitions: [] Added partitions (assigned - owned): [] Revoked partitions (owned - assigned): []
疑问
- 构造EmbeddedKafkaBroker时已指定2个分区(
new EmbeddedKafkaBroker(1, true, 2)),但日志显示主题仅1个分区,该配置为何未生效? - 同一消费组下的两个消费者,除了调整分区数外,还有哪些配置或排查方向能让它们正常消费消息?
内容的提问来源于stack exchange,提问作者DEEPAK MALIK
相关产品推荐
相关产品推荐

