Spring Boot 配置 Kafka 完整指南
一、项目准备
1.1 创建 Spring Boot 项目
确保你的项目使用 Spring Boot 2.x 或 3.x 版本,推荐使用 Spring Initializr 创建项目,或者手动在 pom.xml 或 build.gradle 中添加依赖。
1.2 添加 Maven 依赖
<!– pom.xml –>
<dependencies>
<!– Spring Boot Starter Web –>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
<!– Spring Kafka –>
<dependency>
<groupId>org.springframework.kafka</groupId>
<artifactId>spring-kafka</artifactId>
</dependency>
<!– 可选:JSON 序列化 –>
<dependency>
<groupId>com.fasterxml.jackson.core</groupId>
<artifactId>jackson-databind</artifactId>
</dependency>
</dependencies>
1.3 Gradle 依赖
// build.gradle
dependencies {
implementation 'org.springframework.boot:spring-boot-starter-web'
implementation 'org.springframework.kafka:spring-kafka'
implementation 'com.fasterxml.jackson.core:jackson-databind'
}
二、基础配置
2.1 application.yml 配置
# application.yml
spring:
kafka:
# Kafka 服务器地址
bootstrap-servers: localhost:9092
# 生产者配置
producer:
# 键的序列化器
key-serializer: org.apache.kafka.common.serialization.StringSerializer
# 值的序列化器
value-serializer: org.apache.kafka.common.serialization.StringSerializer
# 可选:自定义分区器
# partitioner-class: com.example.demo.CustomPartitioner
# 可选:消息确认机制
acks: all
# 可选:重试次数
retries: 3
# 可选:批量发送大小
batch-size: 16384
# 可选:缓冲区大小
buffer-memory: 33554432
# 消费者配置
consumer:
# 键的反序列化器
key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
# 值的反序列化器
value-deserializer: org.apache.kafka.common.serialization.StringDeserializer
# 消费者组 ID(必填)
group-id: demo–group
# 自动提交偏移量
enable-auto-commit: true
# 自动提交的时间间隔
auto-commit-interval: 1000
# 消费者重置策略
auto-offset-reset: latest
# 一次拉取的最大条数
max-poll-records: 100
# 监听器配置
listener:
# 消费者并发数(消费者实例数)
concurrency: 3
# 消费者监听器类型(单线程或批处理)
type: single
2.2 配置参数说明
| bootstrap-servers | Kafka 集群地址,多个用逗号分隔 | localhost:9092 |
| key-serializer/value-serializer | 消息键和值的序列化方式 | StringSerializer 或 JsonSerializer |
| key-deserializer/value-deserializer | 消息键和值的反序列化方式 | StringDeserializer 或 JsonDeserializer |
| acks | 生产者确认机制:0不等待确认,1等待leader确认,all等待所有副本确认 | all(确保可靠性) |
| retries | 发送失败重试次数 | 3 |
| auto-offset-reset | 消费者偏移量重置策略:earliest从最早开始,latest从最新开始,none报错 | latest |
| concurrency | 监听器并发数(消费者线程数) | 根据分区数设置 |
三、生产者实现
3.1 发送简单字符串消息
package com.example.demo.producer;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.stereotype.Component;
@Component
public class MessageProducer {
@Autowired
private KafkaTemplate<String, String> kafkaTemplate;
/**
* 发送简单字符串消息
* @param topic 主题名称
* @param message 消息内容
*/
public void sendMessage(String topic, String message) {
kafkaTemplate.send(topic, message);
}
/**
* 发送带键的消息
* @param topic 主题名称
* @param key 消息键
* @param message 消息内容
*/
public void sendMessageWithKey(String topic, String key, String message) {
kafkaTemplate.send(topic, key, message);
}
}
3.2 发送 JSON 对象消息
3.2.1 创建实体类
package com.example.demo.model;
import java.io.Serializable;
public class User implements Serializable {
private Long id;
private String name;
private String email;
// 构造函数
public User() {}
public User(Long id, String name, String email) {
this.id = id;
this.name = name;
this.email = email;
}
// Getter 和 Setter 方法
public Long getId() {
return id;
}
public void setId(Long id) {
this.id = id;
}
public String getName() {
return name;
}
public void setName(String name) {
this.name = name;
}
public String getEmail() {
return email;
}
public void setEmail(String email) {
this.email = email;
}
@Override
public String toString() {
return "User{" +
"id=" + id +
", name='" + name + '\\'' +
", email='" + email + '\\'' +
'}';
}
}
3.2.2 JSON 序列化器配置
修改 application.yml,使用 JSON 序列化器:
spring:
kafka:
producer:
key-serializer: org.apache.kafka.common.serialization.StringSerializer
value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
consumer:
key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer
properties:
spring.json.trusted.packages: "*"
3.2.3 发送 JSON 消息
package com.example.demo.producer;
import com.example.demo.model.User;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.stereotype.Component;
@Component
public class UserProducer {
@Autowired
private KafkaTemplate<String, User> kafkaTemplate;
/**
* 发送 User 对象
* @param topic 主题名称
* @param user 用户对象
*/
public void sendUser(String topic, User user) {
kafkaTemplate.send(topic, user);
}
}
3.3 回调处理
发送消息时可以添加回调函数,处理发送成功或失败的情况:
@Component
public class MessageProducer {
@Autowired
private KafkaTemplate<String, String> kafkaTemplate;
public void sendMessageWithCallback(String topic, String message) {
kafkaTemplate.send(topic, message)
.addCallback(
success -> {
System.out.println("消息发送成功: " + message);
},
failure -> {
System.err.println("消息发送失败: " + failure.getMessage());
}
);
}
}
四、消费者实现
4.1 监听简单字符串消息
package com.example.demo.consumer;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Component;
@Component
public class MessageConsumer {
/**
* 监听主题消息
* @param message 消息内容
*/
@KafkaListener(topics = "demo-topic", groupId = "demo-group")
public void consumeMessage(String message) {
System.out.println("接收到消息: " + message);
}
}
4.2 监听 JSON 对象消息
package com.example.demo.consumer;
import com.example.demo.model.User;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Component;
@Component
public class UserConsumer {
/**
* 监听 User 对象消息
* @param user 用户对象
*/
@KafkaListener(topics = "user-topic", groupId = "demo-group")
public void consumeUser(User user) {
System.out.println("接收到用户: " + user);
}
}
4.3 手动提交偏移量
为了更精确地控制消息消费,可以配置手动提交偏移量:
4.3.1 修改 application.yml
spring:
kafka:
consumer:
enable-auto-commit: false
listener:
ack-mode: manual
4.3.2 手动提交偏移量
@Component
public class ManualCommitConsumer {
@KafkaListener(topics = "demo-topic", groupId = "demo-group")
public void consumeMessage(String message, Acknowledgment acknowledgment) {
System.out.println("接收到消息: " + message);
try {
// 处理消息业务逻辑
processMessage(message);
// 手动提交偏移量
acknowledgment.acknowledge();
} catch (Exception e) {
System.err.println("消息处理失败: " + e.getMessage());
// 不提交偏移量,消息会重新消费
}
}
private void processMessage(String message) {
// 模拟消息处理
System.out.println("处理消息: " + message);
}
}
4.4 批量消费消息
当需要提高吞吐量时,可以配置批量消费:
4.4.1 修改 application.yml
spring:
kafka:
consumer:
max-poll-records: 100 # 一次拉取的最大条数
listener:
type: batch # 设置为批处理模式
4.4.2 批量消费实现
@Component
public class BatchConsumer {
@KafkaListener(topics = "demo-topic", groupId = "demo-group")
public void consumeBatch(List<String> messages) {
System.out.println("批量接收到 " + messages.size() + " 条消息");
for (String message : messages) {
System.out.println("消息: " + message);
}
}
}
五、高级配置与最佳实践
5.1 消息拦截器
消息拦截器可以在消息发送和消费前后进行处理:
5.1.1 生产者拦截器
package com.example.demo.interceptor;
import org.apache.kafka.clients.producer.ProducerInterceptor;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.clients.producer.RecordMetadata;
import java.util.Map;
public class CustomProducerInterceptor implements ProducerInterceptor<String, String> {
@Override
public ProducerRecord<String, String> onSend(ProducerRecord<String, String> record) {
// 发送前处理
System.out.println("生产者拦截器:准备发送消息: " + record.value());
return record;
}
@Override
public void onAcknowledgement(RecordMetadata metadata, Exception exception) {
// 发送后处理
if (exception == null) {
System.out.println("生产者拦截器:消息发送成功,分区: " + metadata.partition());
} else {
System.err.println("生产者拦截器:消息发送失败: " + exception.getMessage());
}
}
@Override
public void close() {
System.out.println("生产者拦截器关闭");
}
@Override
public void configure(Map<String, ?> configs) {
System.out.println("生产者拦截器配置");
}
}
5.1.2 配置拦截器
在 application.yml 中添加拦截器配置:
spring:
kafka:
producer:
properties:
interceptor.classes: com.example.demo.interceptor.CustomProducerInterceptor
5.2 自定义分区策略
默认情况下,Kafka 使用 DefaultPartitioner 进行分区分配。你可以自定义分区策略:
package com.example.demo.partitioner;
import org.apache.kafka.clients.producer.Partitioner;
import org.apache.kafka.common.Cluster;
import org.apache.kafka.common.PartitionInfo;
import org.apache.kafka.common.utils.Utils;
import java.util.List;
import java.util.Map;
import java.util.concurrent.atomic.AtomicInteger;
public class CustomPartitioner implements Partitioner {
private final AtomicInteger counter = new AtomicInteger(0);
@Override
public int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster) {
List<PartitionInfo> partitions = cluster.partitionsForTopic(topic);
int numPartitions = partitions.size();
// 如果没有 key,使用轮询策略
if (keyBytes == null) {
return counter.getAndIncrement() % numPartitions;
} else {
// 如果有 key,使用 hash 策略
return Utils.toPositive(Utils.murmur2(keyBytes)) % numPartitions;
}
}
@Override
public void close() {
System.out.println("自定义分区器关闭");
}
@Override
public void configure(Map<String, ?> configs) {
System.out.println("自定义分区器配置");
}
}
在 application.yml 中配置自定义分区器:
spring:
kafka:
producer:
partitioner-class: com.example.demo.partitioner.CustomPartitioner
5.3 消息重试与死信队列
5.3.1 配置重试策略
spring:
kafka:
listener:
ack-mode: manual_immediate
consumer:
enable-auto-commit: false
template:
retry:
max-attempts: 3 # 重试次数
backoff:
delay: 1000 # 重试间隔(毫秒)
multiplier: 2 # 延迟倍数
5.3.2 死信队列配置
package com.example.demo.consumer;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.support.Acknowledgment;
import org.springframework.stereotype.Component;
@Component
public class DlqConsumer {
@Autowired
private KafkaTemplate<String, String> kafkaTemplate;
private static final String DEAD_LETTER_TOPIC = "demo-topic-dlt";
@KafkaListener(topics = "demo-topic", groupId = "demo-group")
public void consumeMessage(ConsumerRecord<String, String> record, Acknowledgment acknowledgment) {
try {
System.out.println("接收到消息: " + record.value());
// 模拟处理失败
if (record.value().contains("error")) {
throw new RuntimeException("模拟处理失败");
}
// 处理成功,提交偏移量
acknowledgment.acknowledge();
} catch (Exception e) {
System.err.println("消息处理失败: " + e.getMessage());
// 发送到死信队列
kafkaTemplate.send(DEAD_LETTER_TOPIC, record.value());
// 提交偏移量(避免重复消费)
acknowledgment.acknowledge();
}
}
}
5.4 消息幂等性
为了避免消息重复消费导致数据不一致,需要实现幂等性:
@Component
public class IdempotentConsumer {
// 使用 Redis 或数据库记录已处理的消息 ID
private final Set<String> processedMessageIds = new HashSet<>();
@KafkaListener(topics = "demo-topic", groupId = "demo-group")
public void consumeMessage(String message) {
String messageId = extractMessageId(message);
// 检查是否已处理
if (processedMessageIds.contains(messageId)) {
System.out.println("消息已处理,跳过: " + messageId);
return;
}
// 处理消息
processMessage(message);
// 记录已处理的消息 ID
processedMessageIds.add(messageId);
}
private String extractMessageId(String message) {
// 从消息中提取 ID(假设消息格式为 "id:content")
return message.split(":")[0];
}
private void processMessage(String message) {
System.out.println("处理消息: " + message);
}
}
5.5 消息事务支持
Spring Kafka 支持事务消息,确保消息的原子性:
package com.example.demo.producer;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.support.SendResult;
import org.springframework.stereotype.Component;
import org.springframework.transaction.annotation.Transactional;
import java.util.concurrent.CompletableFuture;
@Component
public class TransactionalProducer {
@Autowired
private KafkaTemplate<String, String> kafkaTemplate;
/**
* 事务发送消息
* @param topic 主题名称
* @param message 消息内容
*/
@Transactional
public void sendTransactionalMessage(String topic, String message) {
CompletableFuture<SendResult<String, String>> future = kafkaTemplate.send(topic, message);
future.whenComplete((result, ex) -> {
if (ex == null) {
System.out.println("事务消息发送成功: " + message);
} else {
System.err.println("事务消息发送失败: " + ex.getMessage());
}
});
}
}
在 application.yml 中启用事务支持:
spring:
kafka:
producer:
transaction-id-prefix: tx–
六、测试与验证
6.1 创建测试 Controller
package com.example.demo.controller;
import com.example.demo.model.User;
import com.example.demo.producer.MessageProducer;
import com.example.demo.producer.UserProducer;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.*;
@RestController
@RequestMapping("/api/kafka")
public class KafkaTestController {
@Autowired
private MessageProducer messageProducer;
@Autowired
private UserProducer userProducer;
/**
* 发送字符串消息
*/
@PostMapping("/send")
public String sendMessage(@RequestParam String message) {
messageProducer.sendMessage("demo-topic", message);
return "消息已发送: " + message;
}
/**
* 发送用户对象
*/
@PostMapping("/send-user")
public String sendUser(@RequestBody User user) {
userProducer.sendUser("user-topic", user);
return "用户信息已发送: " + user;
}
}
6.2 测试步骤
bin/kafka-topics.sh –create –topic user-topic –bootstrap-server localhost:9092 –partitions 3 –replication-factor 1
curl -X POST -H "Content-Type: application/json" \\
-d '{"id":1,"name":"张三","email":"zhangsan@example.com"}' \\
http://localhost:8080/api/kafka/send-user
七、常见问题与解决方案
7.1 消息乱序
原因:多个分区或消费者线程导致消息乱序
解决方案:
- 使用单分区主题
- 在消息中添加时间戳或序列号,消费时排序
- 使用顺序消息方案(如 RocketMQ 的有序消息)
7.2 消息丢失
原因:
- 生产者未正确配置 acks
- 消费者自动提交偏移量,但消息处理失败
- Kafka 服务端故障
解决方案:
- 生产者设置 acks: all
- 使用手动提交偏移量
- 确保 Kafka 集群有足够的副本因子
7.3 消息重复消费
原因:
- 消费者处理消息后未及时提交偏移量
- Kafka 重平衡导致偏移量重置
解决方案:
- 实现消息幂等性(使用 Redis 或数据库记录已处理消息)
- 使用手动提交偏移量
7.4 消费者积压
原因:
- 消费速度低于生产速度
- 消费者线程数不足
- 消息处理逻辑耗时过长
解决方案:
- 增加消费者并发数(concurrency)
- 增加分区数量
- 优化消息处理逻辑
7.5 连接超时
原因:
- Kafka 服务未启动或网络不通
- 配置的服务地址错误
解决方案:
- 检查 Kafka 服务是否启动
- 确认 bootstrap-servers 配置正确
- 检查防火墙设置
八、总结
本文档详细介绍了在 Spring Boot 中配置和使用 Kafka 的完整流程,涵盖了:
在实际应用中,根据业务需求合理配置 Kafka 参数,并遵循最佳实践,可以构建高性能、高可靠的分布式消息系统。
注意:本文档基于 Spring Boot 3.x 和 Kafka 3.x 版本编写,不同版本可能存在细微差异。






