温馨提示×

debian kafka怎么配消费者

小樊
38
2026-08-31 23:16:05
栏目: 智能运维

下面以 Debian 系统为背景,给你一个 Kafka 消费者配置与使用示例,分为 命令行方式代码方式(Java/Python),看你用什么场景。


一、前提条件(Debian)

1. 安装 Kafka

sudo apt update
sudo apt install default-jdk -y

下载 Kafka(示例版本 3.6.1):

wget https://downloads.apache.org/kafka/3.6.1/kafka_2.13-3.6.1.tgz
tar -xzf kafka_2.13-3.6.1.tgz
cd kafka_2.13-3.6.1

2. 启动 Kafka

# 启动 Zookeeper(老版本)
bin/zookeeper-server-start.sh config/zookeeper.properties

# 新版本可只用 KRaft(略)
bin/kafka-server-start.sh config/server.properties

二、命令行消费者(最常用)

1. 创建 Topic

bin/kafka-topics.sh \
  --create \
  --topic test \
  --bootstrap-server localhost:9092 \
  --partitions 1 \
  --replication-factor 1

2. 启动消费者

bin/kafka-console-consumer.sh \
  --topic test \
  --bootstrap-server localhost:9092 \
  --from-beginning

✅ 此后只要在另一个终端生产消息,这里就能看到。

3. 生产者测试

bin/kafka-console-producer.sh \
  --topic test \
  --bootstrap-server localhost:9092

输入:

hello kafka

三、消费者核心配置说明

常用参数(consumer.properties):

bootstrap.servers=localhost:9092
group.id=test-group
auto.offset.reset=earliest
enable.auto.commit=true
参数 说明
group.id 消费者组(非常重要)
auto.offset.reset earliest / latest
enable.auto.commit 是否自动提交 offset

启动:

bin/kafka-console-consumer.sh \
  --topic test \
  --bootstrap-server localhost:9092 \
  --consumer.config config/consumer.properties

四、Java 消费者示例(Debian + Maven)

1. 依赖

<dependency>
  <groupId>org.apache.kafka</groupId>
  <artifactId>kafka-clients</artifactId>
  <version>3.6.1</version>
</dependency>

2. Java 代码

Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("group.id", "test-group");
props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
props.put("auto.offset.reset", "earliest");

KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(List.of("test"));

while (true) {
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
    for (ConsumerRecord<String, String> r : records) {
        System.out.println(r.value());
    }
}

五、Python 消费者(Debian)

pip install kafka-python
from kafka import KafkaConsumer

consumer = KafkaConsumer(
    'test',
    bootstrap_servers='localhost:9092',
    group_id='test-group',
    auto_offset_reset='earliest'
)

for msg in consumer:
    print(msg.value.decode())

六、常见问题

1. 收不到消息?

  • 检查 group.id
  • 确认 auto.offset.reset=earliest
  • 确认 topic 名字正确

2. Debian 防火墙

sudo ufw allow 9092

如果你告诉我:

  • Kafka 版本
  • 是否用 Docker
  • 用什么语言消费

我可以给你 更精确的配置方案

0