SpringBoot与Kafka的集成是构建异步通信系统的常见方案。Spring官方提供的spring-kafka项目封装了Kafka客户端的底层细节,开发者只需少量配置就能实现消息的生产与消费。本文将从环境准备、依赖引入、配置文件、生产者与消费者的完整代码实现讲起,最后汇总集成过程中的常见问题与注意事项,适合刚接触Kafka的新手跟着操作。

一、准备工作与依赖引入
在开始编码之前,需要确保本地或远程环境有一个可用的Kafka服务。可以使用官方二进制包自行启动,也可以用Docker快速拉起一个单节点集群。以Docker为例,执行以下命令即可启动一个监听在9092端口的Kafka实例:
docker run -d --name kafka \ -p 9092:9092 \ apache/kafka:latest
接着创建一个SpringBoot项目,推荐版本2.7以上。在pom.xml中引入spring-kafka依赖,SpringBoot会自动管理版本,无需手动指定:
<dependency>
<groupId>org.springframework.kafka</groupId>
<artifactId>spring-kafka</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
spring-kafka的核心抽象有两个:KafkaTemplate用于发送消息,@KafkaListener注解用于声明消费者方法。SpringBoot的自动配置会根据application.yml中的属性自动创建这两个Bean,开发者不需要手动编写@Configuration类,这大大降低了上手门槛。
二、核心配置文件详解
在resources目录下创建application.yml,配置Kafka的连接地址和序列化方式。下面是一份生产环境常用的配置示例:
spring:
kafka:
bootstrap-servers: 127.0.0.1:9092
producer:
# 发送失败的重试次数
retries: 3
# 批量发送大小,单位字节
batch-size: 16384
# 重试次数大于0时需要保证消息顺序,设置max.in.flight.requests.per.connection为1
properties:
max.in.flight.requests.per.connection: 1
key-serializer: org.apache.kafka.common.serialization.StringSerializer
value-serializer: org.apache.kafka.common.serialization.StringSerializer
# 确认模式:all表示所有副本确认才算发送成功
acks: all
consumer:
# 消费者组ID,同组内的消费者分摊分区
group-id: demo-group
# 没有初始offset或offset越界时的策略:earliest从最早开始,latest从最新开始
auto-offset-reset: earliest
# 是否自动提交offset
enable-auto-commit: false
key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
value-deserializer: org.apache.kafka.common.serialization.StringDeserializer
listener:
# 手动确认模式,配合自动提交关闭使用
ack-mode: manual_immediate
# 并发数,不能超过分区数
concurrency: 3
几个关键参数需要重点理解。acks设置为all时,消息需要等待所有ISR副本确认才算写入成功,可靠性最高但延迟增大;batch-size控制生产者批量发送的缓冲区大小,批量发送能显著提升吞吐量;auto-offset-reset决定了消费者组第一次消费或offset失效时的行为,测试环境建议用earliest,避免漏掉历史消息。
ack-mode的取值也值得注意。默认的自动提交模式在消息处理失败时可能造成消息丢失,而manual_immediate模式允许开发者在业务逻辑执行成功后再手动提交offset,配合enable-auto-commit: false使用,是保证消息不丢失的推荐组合。
三、生产者与消费者代码实现
先编写一个简单的Controller,注入KafkaTemplate后就可以发送消息。KafkaTemplate的泛型分别对应key和value的序列化类型:
@RestController
@RequestMapping("/kafka")
public class KafkaProducerController {
@Autowired
private KafkaTemplate<String, String> kafkaTemplate;
@GetMapping("/send")
public String send(@RequestParam String message) {
// 发送消息到topic名为test-topic的主题,指定key保证同一key进入同一分区
kafkaTemplate.send("test-topic", "order-key", message);
return "消息发送成功: " + message;
}
@GetMapping("/sendAsync")
public String sendAsync(@RequestParam String message) {
// 异步发送并通过回调感知结果
CompletableFuture<SendResult<String, String>> future =
kafkaTemplate.send("test-topic", message);
future.whenComplete((result, ex) -> {
if (ex == null) {
System.out.println("发送成功,offset: "
+ result.getRecordMetadata().offset());
} else {
System.out.println("发送失败: " + ex.getMessage());
}
});
return "异步发送已提交";
}
}
send方法默认是异步执行的,返回CompletableFuture对象。如果不关心结果可以直接调用,但生产环境强烈建议注册回调,否则发送失败会被静默吞掉。如果需要保证同一业务ID的消息顺序,务必指定相同的key,Kafka会对key做哈希后路由到固定分区。
消费者的实现更加简单,只需要在方法上添加@KafkaListener注解。下面的示例演示了手动提交offset的写法:
@Component
public class KafkaConsumerService {
// 监听指定topic,支持SPEL表达式配置topic名
@KafkaListener(topics = "test-topic", groupId = "demo-group")
public void onMessage(ConsumerRecord<String, String> record, Acknowledgment ack) {
try {
System.out.println("收到消息: key=" + record.key()
+ ", value=" + record.value()
+ ", partition=" + record.partition()
+ ", offset=" + record.offset());
// 处理业务逻辑...
// 业务成功后手动提交offset
ack.acknowledge();
} catch (Exception e) {
// 不提交offset,消息会在重平衡后重新消费
System.out.println("消费异常: " + e.getMessage());
}
}
}
ConsumerRecord对象封装了消息的全部元数据,包括分区号、offset、时间戳、头部信息等。手动调用acknowledge()的前提是配置了manual或manual_immediate的ack-mode,否则启动时会报Acknowledgment参数无法解析的错误。
如果需要传输对象,可以定义一个实体类并使用JSON序列化。生产者端配置JsonSerializer,消费者端配置JsonDeserializer并指定信任的包名,这样就能直接收发Java对象:
// 发送对象消息
@Autowired
private KafkaTemplate<String, OrderInfo> jsonKafkaTemplate;
public void sendOrder(OrderInfo order) {
jsonKafkaTemplate.send("order-topic", order);
}
// 消费者直接接收对象
@KafkaListener(topics = "order-topic")
public void onOrder(OrderInfo order) {
System.out.println("收到订单: " + order.getOrderId());
}
四、发送JSON对象消息的完整方案
传输对象最常见的方式是借助SpringBoot统一的JSON配置。可以在配置类中手动构建KafkaTemplate,指定JsonSerializer作为值序列化器:
@Configuration
public class KafkaConfig {
@Bean
public ProducerFactory<String, Object> producerFactory(KafkaProperties properties) {
Map<String, Object> configs = properties.buildProducerProperties();
configs.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,
JsonSerializer.class);
return new DefaultKafkaProducerFactory<>(configs);
}
@Bean
public KafkaTemplate<String, Object> jsonKafkaTemplate(
ProducerFactory<String, Object> producerFactory) {
return new KafkaTemplate<>(producerFactory);
}
}
消费者端需要在yml中为JsonDeserializer配置信任包,否则会抛出ClassNotFound相关异常:
spring:
kafka:
consumer:
properties:
spring.json.trusted.packages: "com.example.demo.entity"
需要注意JSON方案的兼容性问题。如果生产者的实体类新增了字段而消费者还是旧版本,默认会反序列化失败。可以配置spring.json.value.default.type指定目标类型,并在实体类上使用@JsonIgnoreProperties(ignoreUnknown = true)忽略未知字段,提高前后版本的兼容性。
五、常见问题与注意事项
1. 连接超时或Connection refused。最常见的原因是bootstrap-servers地址配置错误。Kafka的broker默认会返回它配置的advertised.listeners地址给客户端,如果Kafka部署在Docker或云服务器上,需要把advertised.listeners设置为客户端可达的外部地址,否则客户端能连上9092端口但后续获取元数据时会失败。
2. 消息乱码或反序列化异常。生产者和消费者的序列化器必须匹配,生产端用StringSerializer、消费端却配置了JsonDeserializer就会出现乱码或异常。排查时先确认两端的serializer和deserializer类型是否一致。
3. 重复消费问题。当消费者处理完消息但还没来得及提交offset就发生重启或重平衡,这部分消息会被再次消费。因此消费逻辑要设计成幂等的,比如用消息key或业务ID做去重,配合数据库唯一索引兜底。
4. 消息丢失问题。三个环节都可能导致丢失:生产端acks设为0或1且broker宕机、broker端副本因子为1、消费端自动提交offset后处理失败。生产环境建议组合使用acks=all、retries大于0、副本因子不小于2、消费端手动提交,形成端到端的可靠保障。
5. 消费者不消费或分区分配不均。同一个groupId下的消费者数量超过分区数时,多余的消费者会处于空闲状态。concurrency配置的并发线程数也受此限制,比如topic只有3个分区,concurrency设为5就有2个线程浪费。如果想提高消费能力,应该先增加分区数。
6. 消息顺序性问题。Kafka只保证单个分区内的消息有序。如果发送时不指定key,消息会被轮询到不同分区导致乱序。对顺序敏感的业务,指定固定key,同时设置max.in.flight.requests.per.connection为1,防止重试时打乱顺序。
7. 版本兼容问题。spring-kafka的版本要与Kafka服务端版本相匹配,过新的客户端连接过旧的broker可能出现协议不兼容的异常。如果遇到UnsupportedVersionException,优先检查两端的版本差异,必要时显式指定spring-kafka的版本号。
六、小结
SpringBoot集成Kafka的整体流程可以概括为:引入spring-kafka依赖、编写application.yml配置、注入KafkaTemplate发送消息、使用@KafkaListener消费消息。入门门槛不高,但要在生产环境稳定运行,必须理解offset提交机制、分区与消费者组的关系、序列化方式匹配等底层概念。建议新手先在本地用Docker搭一个单节点环境跑通流程,再逐步实践手动提交、JSON序列化、异常重试等进阶特性,最终形成可靠的消息收发体系。
SpringBoot集成KafkaKafka生产者消费者Spring Kafka配置修改时间:2026-09-01 22:20:50