一、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:90924. 命令行测试
发送消息:
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.servers | Kafka 集群地址,多个用逗号分隔 |
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.servers | Kafka 集群地址 |
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 # 手动提交 offset3. 生产者服务
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 的基本步骤:
- 启动 Kafka 集群(单机或集群)。
- 创建 Topic。
- 编写生产者:配置序列化器、acks 等,发送消息。
- 编写消费者:配置反序列化器、group.id,订阅 Topic 并处理消息。
- 生产环境:考虑消息可靠性(不丢失、不重复)、顺序性、分区策略、监控等。