You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Kafka僵尸消费者被逐出消费组后仍可Poll记录?技术问询

问题描述

我的应用从MSK Kafka集群拉取记录,自行维护各分区的offset,因此禁用了自动提交(autocommit),从不将offset提交到Kafka,而是持久化到内部数据存储。

生产环境中出现重复记录问题,直到重启存在僵尸消费者(已被逐出消费组但仍在执行Poll操作)的实例,重复才停止。当消费者无法发送心跳时会变成僵尸消费者,这通常由MSK中的IAM认证问题导致,且有时会持续较长时间。

进一步排查发现:原本分配给僵尸消费者的分区已被分配给其他活跃消费者,但僵尸消费者的Poll操作仍能返回记录。

核心疑问:僵尸消费者在被逐出组后仍应能通过Poll获取记录吗?

本地复现环境
  • 单个Broker(Docker镜像),配置三个不同的外部端口
  • 创建一个含两个分区的Topic
  • 使用Toxiproxy代理这些端口
  • 两个消费者(订阅该Topic)分别连接一个代理端口,session timeout设为10秒
  • 在其中一个代理中引入延迟,将对应的消费者逐出组
  • 向每个分区生产若干消息
  • 结果:两个消费者均持续获取到消息
Kafka Broker启动命令
docker run -d --name kafka-container \
    -e KAFKA_ENABLE_KRAFT=yes \
    -e KAFKA_CFG_NODE_ID=1 \
    -e KAFKA_CFG_PROCESS_ROLES=broker,controller \
    -e KAFKA_CFG_CONTROLLER_LISTENER_NAMES=CONTROLLER \
    -e KAFKA_BROKER_ID=1 \
    -e KAFKA_CFG_CONTROLLER_QUORUM_VOTERS=1@localhost:9094 \
    -e KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT,EXTERNAL:PLAINTEXT,EXTERNAL1:PLAINTEXT,EXTERNAL2:PLAINTEXT \
    -e KAFKA_CFG_LISTENERS=PLAINTEXT://:9092,EXTERNAL://:9093,EXTERNAL1://:29093,EXTERNAL2://:39093,CONTROLLER://:9094 \
    -e KAFKA_CFG_ADVERTISED_LISTENERS=EXTERNAL://localhost:9092,EXTERNAL1://localhost:29092,EXTERNAL2://localhost:39092,PLAINTEXT://$CONTAINER_NAME:9092 \
    -e ALLOW_PLAINTEXT_LISTENER=yes \
    -e KAFKA_KRAFT_CLUSTER_ID=5e-ZWJELQjCqW2n3bwMTWg \
    -e KAFKA_AUTO_CREATE_TOPICS_ENABLE=false \
    -p "9092:9093" \
    -p "29093:29093" \
    -p "39093:39093" \
    bitnami/kafka:latest
Toxiproxy代理配置
toxiproxy-cli create -l localhost:29092 -u localhost:29093 zombie_consumer_scenario_healthy_proxy

toxiproxy-cli create -l localhost:39092 -u localhost:39093 zombie_consumer_scenario_unhealthy_proxy
Kafka控制台消费者启动命令
kafka-console-consumer --group test-cg --topic test --consumer-property enable.auto.commit=false --consumer-property session.timeout.ms=10000 --bootstrap-server localhost:29092

kafka-console-consumer --group test-cg --topic test --consumer-property enable.auto.commit=false --consumer-property session.timeout.ms=10000 --bootstrap-server localhost:39092
引入延迟的Toxic配置
toxiproxy-cli toxic add -t latency -n myToxic -a latency=10000 -a jitter=50 zombie_consumer_scenario_unhealthy_proxy

我正考虑在应用中自行控制分区分配,但认为Kafka消费者在被逐出组后不应通过Poll返回记录。


内容的提问来源于stack exchange,提问作者Ravi Gupta

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.15 21:13:15