Kafka消费者无法获取消息,Docker部署+Python客户端问题排查求助
Kafka连接后无法收发消息问题解决方案
1. 核心配置错误:Kafka对外监听地址配置失效
你当前配置的KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://kafka:9092是Docker内部网络的域名,仅能被同Docker网络下的服务解析访问。外部客户端完成9092端口的初始握手后,Kafka会返回该内部域名给客户端,客户端无法解析kafka域名,自然无法完成后续消息收发。
修正docker-compose.yaml配置
将kafka的环境变量调整为支持内外网同时访问的模式:
version: "3.9" services: zookeeper: image: "bitnami/zookeeper:latest" ports: - "2181:2181" environment: - ALLOW_ANONYMOUS_LOGIN=yes networks: - "pocnetwork" kafka: image: bitnami/kafka:2.8.1 ports: - "9092:9092" environment: - KAFKA_BROKER_ID=1 - KAFKA_LISTENERS=PLAINTEXT://0.0.0.0:9092,PLAINTEXT_INTERNAL://0.0.0.0:9093 - KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://192.168.10.140:9092,PLAINTEXT_INTERNAL://kafka:9093 - KAFKA_LISTENER_SECURITY_PROTOCOL_MAP=PLAINTEXT:PLAINTEXT,PLAINTEXT_INTERNAL:PLAINTEXT - KAFKA_INTER_BROKER_LISTENER_NAME=PLAINTEXT_INTERNAL - KAFKA_ZOOKEEPER_CONNECT=zookeeper:2181 - ALLOW_PLAINTEXT_LISTENER=yes - KAFKA_CFG_AUTO_CREATE_TOPICS_ENABLE=true depends_on: - zookeeper networks: - "pocnetwork" networks: pocnetwork:
修改完成后执行docker-compose down && docker-compose up -d重启容器生效。
2. 生产者代码优化
KafkaProducer.send()是异步方法,直接调用后消息可能还在本地缓冲区未实际发出,可添加等待逻辑确认发送结果、捕获异常:
from time import sleep from json import dumps from kafka import KafkaProducer producer = KafkaProducer(bootstrap_servers=['192.168.10.140:9092'], api_version=(0,11,5), value_serializer=lambda x: dumps(x).encode('utf-8')) for e in range(1000): data = {'number' : e} print(f"Sending {e}") future = producer.send('numtest', value=data) # 等待发送完成,出错会抛出异常方便排查 result = future.get(timeout=10) print(f"Sent {e}, partition: {result.partition}, offset: {result.offset}") sleep(1)
3. 消费者代码优化
consumer_timeout_ms=1000会让消费者1秒没收到新消息就自动退出,测试阶段可以先关闭该配置- 多余的
consumer.poll()会提前拉走第一批消息,导致后续for循环无法读取到已拉取内容
修正后代码如下:
from kafka import KafkaConsumer from json import loads consumer = KafkaConsumer( 'numtest', bootstrap_servers=['192.168.10.140:9092'], auto_offset_reset='earliest', enable_auto_commit=True, auto_commit_interval_ms=100, group_id='my-group', value_deserializer=lambda x: loads(x.decode('utf-8'))) print("Start consuming") for message in consumer: print(f'Received message {message.value}, offset: {message.offset}')
内容的提问来源于stack exchange,提问作者Oleksandr Masliuchenko
相关产品推荐
相关产品推荐

