温馨提示×

debian kafka消费者咋配

小樊
43
2026-09-03 18:56:16
栏目: 智能运维

在 Debian 上配置 Kafka 消费者,一般分几种场景:

  • 命令行消费者(最常用,快速测试)
  • Java / Python / Go 等客户端
  • systemd 服务里跑消费者程序

下面按常见情况给你说明。


一、环境准备(Debian)

1. 安装 Java(Kafka 依赖)

sudo apt update
sudo apt install -y default-jdk
java -version

2. 安装 Kafka(如未安装)

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

启动 Zookeeper / Kafka(老版本用 Zookeeper,新版本可用 KRaft):

# Zookeeper 模式
bin/zookeeper-server-start.sh config/zookeeper.properties &
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

常用参数:

  • --group my-group:指定消费者组
  • --max-messages 10:消费 10 条后退出
  • --property print.key=true:打印 key

示例(带 group):

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

三、消费者配置文件方式

Kafka 也支持 consumer.properties

1. 编辑配置

vi config/consumer.properties

示例内容:

bootstrap.servers=localhost:9092
group.id=test-group
enable.auto.commit=true
auto.commit.interval.ms=1000
key.deserializer=org.apache.kafka.common.serialization.StringDeserializer
value.deserializer=org.apache.kafka.common.serialization.StringDeserializer

2. 使用配置文件

bin/kafka-console-consumer.sh \
  --topic test \
  --consumer.config config/consumer.properties

四、Java 消费者示例(Debian 上运行)

1. Maven 依赖

<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("enable.auto.commit", "true");
props.put("key.deserializer",
    "org.apache.kafka.common.serialization.StringDeserializer");
props.put("value.deserializer",
    "org.apache.kafka.common.serialization.StringDeserializer");

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.printf("offset=%d, key=%s, value=%s%n",
            r.offset(), r.key(), r.value());
    }
}

运行:

java -cp your.jar Main

五、Debian systemd 跑消费者(生产推荐)

示例服务文件:

sudo vim /etc/systemd/system/kafka-consumer.service
[Unit]
Description=Kafka Consumer
After=network.target

[Service]
ExecStart=/opt/kafka/bin/kafka-console-consumer.sh \
  --topic test \
  --bootstrap-server localhost:9092 \
  --group test-group
Restart=always
User=kafka

[Install]
WantedBy=multi-user.target
sudo systemctl daemon-reload
sudo systemctl enable kafka-consumer
sudo systemctl start kafka-consumer

六、常见问题

1. 连不上 broker

  • 检查 listenersadvertised.listeners
  • 云服务器要配公网 IP

2. 消费不到数据

  • 是否指定了 --from-beginning
  • group 是否已消费过 offset

3. 中文乱码

  • 确保 producer / consumer 编码一致(UTF-8)

如果你告诉我:

  • Kafka 版本
  • 用命令行还是代码
  • 单机还是集群
  • 是否云服务器

我可以给你 更精确的 Debian 配置示例

0