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

如何让Kafka中所有Consumer接收任意Producer发送的消息?

Kafka多消费端接收全量生产者消息的配置调整

我是Kafka新手,目前已实现多个Producer向单个Consumer发送消息,但希望任意Producer发送的消息能被所有Consumer接收,以下是当前使用的配置与代码:

docker-compose.yml

version: '2'
services:
  zookeeper:
    image: wurstmeister/zookeeper
    container_name: zookeeper
    ports:
     - "2181:2181"
    networks:
     - my_network
    restart: unless-stopped

  kafka:
    image: wurstmeister/kafka
    container_name: kafka
    ports:
     - "9092:9092"
    expose:
     - "9093"
    environment:
      KAFKA_ADVERTISED_HOST_NAME: 192.168.11.30
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
      KAFKA_AUTO_CREATE_TOPICS_ENABLE: 'true'
      KAFKA_CREATE_TOPICS: "topic:5:5"
      KAFKA_LOG_RETENTION_HOURS: 1
      KAFKA_LOG_RETENTION_BYTES: 4073741824
      KAFKA_LOG_SEGMENT_BYTES: 1073741824
      KAFKA_RETENTION_CHECK_INTERVAL_MS: 300000
    volumes:
     - /var/run/docker.sock:/var/run/docker.sock
    networks:
     - my_network
    restart: unless-stopped

networks:
  my_network:
    driver: bridge

producer.py

import json
from kafka import KafkaProducer

kafka_server = "192.168.11.30:9092"

topic = "topic"

producer = KafkaProducer(
    bootstrap_servers=[kafka_server],
    value_serializer=lambda v: json.dumps(v).encode("utf-8"),
    acks='all',
)

while True:
    data = input()
    producer.send(topic, value=data)
    producer.flush()

consumer.py

import json
from kafka import KafkaConsumer

kafka_server = "192.168.11.30:9092"

topic = "topic"

consumer = KafkaConsumer(
    bootstrap_servers=[kafka_server],
    value_deserializer=json.loads,
    enable_auto_commit=False,
    auto_offset_reset="latest",
    group_id="my-group",
)

consumer.subscribe([topic])

try:
    for data in consumer:
        print(data.value)
        consumer.commit()
except KeyboardInterrupt:
    pass
finally:
    consumer.close()

问题原因与解决方法

Kafka中,同一**消费组(group_id)**内的Consumer会分摊订阅主题的分区消息,因此你当前所有Consumer共用my-group的配置,导致消息被分发到不同Consumer,无法实现全量接收。要让每个Consumer都收到所有消息,可采用以下两种方法:

方法1:为每个Consumer设置唯一的group_id

每个Consumer实例使用不同的group_id,这样每个消费组会独立消费主题的所有分区消息,确保全量接收。

修改后的consumer.py示例:

import json
from kafka import KafkaConsumer

kafka_server = "192.168.11.30:9092"
topic = "topic"

# 每个Consumer实例使用唯一的group_id
consumer = KafkaConsumer(
    bootstrap_servers=[kafka_server],
    value_deserializer=json.loads,
    enable_auto_commit=False,
    auto_offset_reset="latest",
    group_id="consumer-001",  # 每个实例修改此值,比如consumer-002、consumer-003等
)

consumer.subscribe([topic])

try:
    for data in consumer:
        print(data.value)
        consumer.commit()
except KeyboardInterrupt:
    pass
finally:
    consumer.close()

方法2:为每个Consumer创建独立主题(可选)

若业务场景允许,可为每个Consumer单独创建主题,Producer向所有主题发送消息。但这种方式会增加主题维护成本,一般优先选择方法1。

补充说明

你的主题配置topic:5:5已设置5个分区和5个副本,该配置可满足多消费组并行消费的需求,无需调整。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 14:30:02