温馨提示×

debian kafka消费者怎么配

小樊
47
2026-09-09 21:25:03
栏目: 智能运维

下面按Debian 系统上配置 Kafka 消费者来给你讲,分 基础环境、消费者配置、命令行消费者、Java 消费者示例 四部分。


一、环境准备(Debian)

1. 安装 Java(Kafka 依赖)

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

2. 安装 Kafka(以官方二进制为例)

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

3. 启动 Kafka(单机)

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

# 启动 Kafka
bin/kafka-server-start.sh config/server.properties &

Kafka 3.x 也可用 KRaft(无 Zookeeper),但初学建议先用上面方式。


二、Kafka 消费者核心配置

消费者主要配置项(consumer.properties):

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

常用参数说明

参数 说明
bootstrap.servers Kafka 地址
group.id 消费组
auto.offset.reset earliest / latest
enable.auto.commit 是否自动提交 offset
session.timeout.ms 消费者超时

三、命令行消费者(最简单)

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-consumer.sh \
  --topic test \
  --bootstrap-server localhost:9092 \
  --group my-group

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

1. Maven 依赖

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

2. Java 代码

import org.apache.kafka.clients.consumer.*;
import java.time.Duration;
import java.util.*;

public class KafkaConsumerDemo {
    public static void main(String[] args) {
        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(Collections.singletonList("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());
            }
        }
    }
}

五、常见问题(Debian)

1. 连不上 Kafka

  • 检查 server.properties
listeners=PLAINTEXT://0.0.0.0:9092
advertised.listeners=PLAINTEXT://<服务器IP>:9092

2. 消费不到数据

  • auto.offset.reset=earliest
  • 确认 topic 有数据
  • 检查 group 是否重复

如果你告诉我:

  • 用的是 Kafka 几.x
  • 是否 集群
  • Shell / Java / Python

我可以直接给你一份可复制的配置文件

0