温馨提示×

ubuntu kafka如何配置消费者

小樊
50
2026-09-17 18:37:05
栏目: 智能运维

在 Ubuntu 上配置 Kafka 消费者,通常分为 命令行方式(测试用)程序方式(Java / Python 等)。下面从安装到配置一步一步说明。


一、前提条件(Ubuntu)

1. 安装 Kafka

# 安装 Java
sudo apt update
sudo apt install openjdk-11-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

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

(Kafka 3.x 也可用 KRaft 模式,不需要 Zookeeper)


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

1. 创建 Topic

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

2. 启动消费者

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

参数说明:

  • --from-beginning:从最早消息开始消费
  • 不加该参数:只消费启动后的新消息

三、消费者核心配置(重要)

Kafka 消费者主要配置项(命令行或代码中通用):

配置项 说明
bootstrap.servers Kafka 地址
group.id 消费者组(非常重要)
key.deserializer key 反序列化
value.deserializer value 反序列化
auto.offset.reset earliest / latest
enable.auto.commit 是否自动提交 offset

四、消费者组(Consumer Group)

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

查看消费进度:

bin/kafka-consumer-groups.sh \
  --describe \
  --group my-group \
  --bootstrap-server localhost:9092

五、Java 消费者示例

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("auto.offset.reset", "earliest");
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-topic"));

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());
    }
}

六、Python 消费者(可选)

pip install kafka-python
from kafka import KafkaConsumer

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

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

七、常见问题

1. 消费不到消息?

  • 检查 bootstrap-server 是否正确
  • 检查 group.id
  • 检查 auto.offset.reset

2. 重复消费?

  • 关闭 enable.auto.commit
  • 手动 commitSync()

3. Ubuntu 防火墙

sudo ufw allow 9092

如果你告诉我:

  • Kafka 版本
  • 使用语言(Java / Python / Go)
  • 单机还是集群

我可以给你更精准的配置示例。

0