news 2026/10/1 15:41:05

Kafka

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Kafka

一、Kafka 环境搭建(单机版)

1. 下载 Kafka

前往 Apache Kafka 官网 下载最新稳定版(例如 3.6.0)。解压后目录结构如下:

kafka_2.13-3.6.0/ ├── bin/ # 启动脚本 ├── config/ # 配置文件 ├── libs/ # 依赖库 └── ...

2. 启动 Kafka

早期 Kafka 依赖 ZooKeeper,新版本推荐使用KRaft模式(无需 ZooKeeper)。以下分别介绍两种方式。

方式一:KRaft 模式(Kafka 3.3+ 推荐)
# 1. 生成集群 ID KAFKA_CLUSTER_ID=$(bin/kafka-storage.sh random-uuid) # 2. 格式化存储目录(使用默认配置) bin/kafka-storage.sh format -t $KAFKA_CLUSTER_ID -c config/kraft/server.properties # 3. 启动 Kafka bin/kafka-server-start.sh config/kraft/server.properties
方式二:ZooKeeper 模式(传统)
# 1. 启动 ZooKeeper bin/zookeeper-server-start.sh config/zookeeper.properties # 2. 启动 Kafka bin/kafka-server-start.sh config/server.properties

默认端口:

  • Kafka:9092
  • ZooKeeper:2181

3. 创建 Topic

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

查看 Topic 列表:

bin/kafka-topics.sh --list --bootstrap-server localhost:9092

4. 命令行测试

发送消息:

bin/kafka-console-producer.sh --topic my-topic --bootstrap-server localhost:9092 > hello kafka

消费消息:

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

二、Java 原生客户端使用

在 Maven 项目中添加依赖(pom.xml):

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

1. 生产者(Producer)

基本配置项
配置项说明
bootstrap.serversKafka 集群地址,多个用逗号分隔
key.serializer键的序列化器(如StringSerializer)
value.serializer值的序列化器
acks确认机制:0不等待确认;1仅 leader 确认;all或-1所有副本确认
retries发送失败重试次数
batch.size批量发送大小(字节)
linger.ms等待更多消息加入批次的时间
buffer.memory生产者缓冲区大小
compression.type压缩类型:none、gzip、snappy、lz4、zstd
生产者代码示例
import org.apache.kafka.clients.producer.*; import org.apache.kafka.common.serialization.StringSerializer; import java.util.Properties; import java.util.concurrent.Future; public class MyProducer { public static void main(String[] args) { Properties props = new Properties(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); props.put(ProducerConfig.ACKS_CONFIG, "all"); props.put(ProducerConfig.RETRIES_CONFIG, 3); props.put(ProducerConfig.BATCH_SIZE_CONFIG, 16384); props.put(ProducerConfig.LINGER_MS_CONFIG, 1); props.put(ProducerConfig.BUFFER_MEMORY_CONFIG, 33554432); props.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, "snappy"); // 创建生产者 Producer<String, String> producer = new KafkaProducer<>(props); // 1. 发送消息(异步,不关心结果) producer.send(new ProducerRecord<>("my-topic", "key1", "value1")); // 2. 发送消息并获取 Future(可阻塞等待结果) Future<RecordMetadata> future = producer.send( new ProducerRecord<>("my-topic", "key2", "value2") ); try { RecordMetadata metadata = future.get(); System.out.println("发送成功,offset=" + metadata.offset() + ", partition=" + metadata.partition()); } catch (Exception e) { e.printStackTrace(); } // 3. 发送消息并带回调(异步) producer.send(new ProducerRecord<>("my-topic", "key3", "value3"), new Callback() { @Override public void onCompletion(RecordMetadata metadata, Exception exception) { if (exception == null) { System.out.println("发送成功: " + metadata.offset()); } else { exception.printStackTrace(); } } }); // 关闭生产者(会等待所有缓冲消息发送完成) producer.close(); } }

2. 消费者(Consumer)

基本配置项
配置项说明
bootstrap.serversKafka 集群地址
group.id消费者组 ID,相同组内的消费者共同消费
key.deserializer键的反序列化器
value.deserializer值的反序列化器
enable.auto.commit是否自动提交 offset
auto.offset.reset初始消费位置:earliest从头开始,latest从最新开始
max.poll.records一次 poll 最多返回的记录数
session.timeout.ms会话超时时间
消费者代码示例
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 MyConsumer { public static void main(String[] args) { Properties props = new Properties(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); props.put(ConsumerConfig.GROUP_ID_CONFIG, "my-group"); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "true"); // 自动提交 props.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG, "1000"); props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); // 从头消费 Consumer<String, String> consumer = new KafkaConsumer<>(props); consumer.subscribe(Collections.singletonList("my-topic")); try { 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()); } } } finally { consumer.close(); } } }
手动提交 offset(推荐生产环境使用)
props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false"); // 关闭自动提交 while (true) { ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100)); for (ConsumerRecord<String, String> record : records) { // 处理消息 System.out.println(record.value()); } // 处理完后手动提交 offset(同步提交) consumer.commitSync(); // 或者异步提交 // consumer.commitAsync(); }

三、Spring Boot 集成 Kafka

Spring Kafka 对 Kafka 进行了封装,大大简化了开发。下面演示完整流程。

1. 添加依赖

在pom.xml中添加:

<dependency> <groupId>org.springframework.kafka</groupId> <artifactId>spring-kafka</artifactId> </dependency>

(Spring Boot 会自动管理版本,通常不需要显式指定版本号)

2. 配置文件(application.yml)

spring: kafka: bootstrap-servers: localhost:9092 producer: key-serializer: org.apache.kafka.common.serialization.StringSerializer value-serializer: org.apache.kafka.common.serialization.StringSerializer acks: all retries: 3 consumer: group-id: my-group key-deserializer: org.apache.kafka.common.serialization.StringDeserializer value-deserializer: org.apache.kafka.common.serialization.StringDeserializer auto-offset-reset: earliest enable-auto-commit: false listener: ack-mode: manual # 手动提交 offset

3. 生产者服务

import org.springframework.kafka.core.KafkaTemplate; import org.springframework.stereotype.Service; @Service public class KafkaProducerService { private final KafkaTemplate<String, String> kafkaTemplate; public KafkaProducerService(KafkaTemplate<String, String> kafkaTemplate) { this.kafkaTemplate = kafkaTemplate; } public void sendMessage(String topic, String message) { kafkaTemplate.send(topic, message); } // 带 key 和回调 public void sendMessageWithCallback(String topic, String key, String message) { kafkaTemplate.send(topic, key, message).addCallback( result -> System.out.println("发送成功: " + result.getRecordMetadata().offset()), ex -> System.err.println("发送失败: " + ex.getMessage()) ); } }

4. 消费者监听器

import org.springframework.kafka.annotation.KafkaListener; import org.springframework.kafka.support.Acknowledgment; import org.springframework.stereotype.Component; @Component public class KafkaConsumer { // 自动确认(配置中 enable-auto-commit: true 时) @KafkaListener(topics = "my-topic", groupId = "my-group") public void listen(String message) { System.out.println("收到消息: " + message); } // 手动提交 offset(配置中 enable-auto-commit: false 时) @KafkaListener(topics = "my-topic", groupId = "my-group") public void listenManual(String message, Acknowledgment ack) { System.out.println("收到消息: " + message); // 处理完成后提交 offset ack.acknowledge(); } }

5. 发送消息测试

在 Controller 或测试类中注入KafkaProducerService调用即可。


四、进阶特性与最佳实践

1. 自定义序列化器

如果消息是 Java 对象,可以使用 JSON 或 Avro 序列化。常用的是 Spring Kafka 提供的JsonSerializer/JsonDeserializer,或者使用StringSerializer配合 Jackson 手动转换。

2. 分区策略

  • 默认分区器:如果指定了 key,则根据 key 的 hash 选择分区;如果 key 为 null,则使用轮询(round-robin)。
  • 自定义分区器:实现Partitioner接口,并在配置中指定partitioner.class。

3. 消息顺序性

Kafka 只能保证同一个分区内消息有序。因此,对于需要严格顺序的业务,可以将需要有序的消息发送到同一个分区(例如使用相同的 key)。

4. 幂等性

Kafka 生产者默认开启幂等(enable.idempotence=true),可以避免因网络重试导致的重复消息。对于消费者端,需要业务逻辑本身具备幂等性。

5. 事务

Kafka 支持事务,可以保证“发送多条消息要么全部成功,要么全部失败”。配置:

props.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "my-transactional-id"); producer.initTransactions(); producer.beginTransaction(); producer.send(record1); producer.send(record2); producer.commitTransaction();

6. 错误处理

在消费者中,如果处理消息失败,可以抛出异常,Spring Kafka 会进行重试。也可以配置ErrorHandler或SeekToCurrentErrorHandler来控制失败后的行为(例如重试一定次数后跳过或发送到死信队列)。


五、常见问题解答

Q1: 如何保证消息不丢失?

  • 生产者:设置acks=all,retries>0,开启幂等。
  • Broker:设置min.insync.replicas至少为 2,保证至少有一个副本同步。
  • 消费者:禁用自动提交,处理完消息后手动提交 offset。

Q2: 如何保证消息不重复?

  • Kafka 本身无法绝对避免重复,需要消费者端做幂等(如数据库唯一约束、Redis 去重等)。
  • 可以使用事务或幂等生产者减少重复发送。

Q3: 消费进度如何回溯?

  • 使用kafka-consumer-groups.sh工具重置 offset:
    bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 --group my-group --reset-offsets --to-earliest --topic my-topic --execute

Q4: 如何监控 Kafka?

  • 使用 Kafka 自带的 JMX 指标 + Prometheus + Grafana。
  • 使用 Kafka Manager、Kafka Eagle、Burrow 等第三方工具。

六、总结

使用 Kafka 的基本步骤:

  1. 启动 Kafka 集群(单机或集群)。
  2. 创建 Topic。
  3. 编写生产者:配置序列化器、acks 等,发送消息。
  4. 编写消费者:配置反序列化器、group.id,订阅 Topic 并处理消息。
  5. 生产环境:考虑消息可靠性(不丢失、不重复)、顺序性、分区策略、监控等。
版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/10/1 15:37:38

QEMU 模拟器学习(二)之 USB 设备模拟与协议学习

笔者学习qemu 模拟器&#xff0c;今天讲一下usb设备这块 基于三个文件进行分析 hw/usb/core.c — USB 控制传输与数据&#xff08;Bulk&#xff09;传输的通用逻辑hw/usb/dev-storage.c — BOT&#xff08;Bulk-Only Transport&#xff0c;协议号 0x50&#xff09;协议解析hw/…

作者头像 李华
网站建设 2026/10/1 15:36:54

Bun.js 全面解析:从 Claude Code 泄漏事件说起的前端新势力

1. 引言&#xff1a;从 Claude Code 泄漏事件说起 2025 年初&#xff0c;一则关于 Claude Code 的代码泄漏事件在开发者社区引发热议。事件中&#xff0c;一段疑似 Claude Code 内部实现的代码片段被意外公开&#xff0c;而细心的开发者发现&#xff0c;这段代码的运行环境并非…

作者头像 李华
网站建设 2026/10/1 15:36:51

资质文件标准化归档方法,节省大量时间

做投标、办资质年审、应对合规审查的人&#xff0c;最头疼的往往不是 "办资质"&#xff0c;而是 "找资质"。翻翻你的电脑&#xff1a;一堆 "新建文件夹 (3)"、"最终版 (1)"、"绝对不改版"&#xff0c;营业执照、资质证书、检…

作者头像 李华