温馨提示×

debian kafka消费者怎么写

小樊
49
2026-09-03 19:10:13
栏目: 智能运维

下面分别给你 Debian 上 Kafka 消费者 的常见写法,覆盖 命令行代码(Java / Python) 两种方式,你可以按需求选。


一、环境准备(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 &

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

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

✅ 适合:

  • 调试
  • 查看消息
  • 验证生产者是否发送成功

三、Java Kafka 消费者(生产常用)

1. Maven 依赖

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

2. 消费者示例

import org.apache.kafka.clients.consumer.*;
import org.apache.kafka.common.serialization.StringDeserializer;

import java.time.Duration;
import java.util.Collections;
import java.util.Properties;

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", StringDeserializer.class.getName());
        props.put("value.deserializer", StringDeserializer.class.getName());
        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> record : records) {
                System.out.printf(
                    "offset=%d, key=%s, value=%s%n",
                    record.offset(), record.key(), record.value()
                );
            }
        }
    }
}

四、Python Kafka 消费者(轻量常用)

1. 安装库

pip install kafka-python

2. 消费者示例

from kafka import KafkaConsumer

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

for msg in consumer:
    print(f"offset={msg.offset}, value={msg.value.decode()}")

五、常见注意点

group.id

  • 相同 group 的消费者会负载均衡消费
  • 不同 group 互不影响

auto.offset.reset

  • earliest:从最早消息开始
  • latest:只消费新消息

Debian 防火墙

sudo ufw allow 9092

如果你告诉我:

  • Kafka 版本
  • 使用语言(Java / Python / Go)
  • 是否集群
  • 是否要持久化 offset

我可以直接给你 可直接跑的完整示例

0