消息队列:Kafka
本教程共 48 篇 · 第 48 篇 · 更新于 2026-08-13 · 约 5 分钟阅读
本节目标:理解 Kafka 的基本概念,学会在 Spring Boot 里集成 Kafka,写生产者发消息、消费者收消息,并处理 JSON 消息。
Kafka 是什么
Kafka 是一个分布式消息系统,核心是「发布-订阅」。可以把 Kafka 想成邮局:发件人(生产者)把信件投进邮局,收件人(消费者)从邮局取信。两边互不认识,也不用同时在线。
几个必须懂的概念:
- Broker:Kafka 服务器,负责存消息。
- Topic:消息的分类,类似信件上的主题标签。
- Partition:Topic 拆成多个分区,消息按分区存储,支持并行读写。
- Offset:消息在分区里的编号,消费者靠它记录读到哪了。
- Consumer Group:一组消费者共同消费一个 Topic。一条消息只会发给组里的一个成员;不同组各拿一份。
Kafka 和 RabbitMQ 的定位不同。Kafka 吞吐量极高,擅长海量日志、事件流、数据管道;RabbitMQ 更轻,适合普通业务消息。选型看场景,没有谁全面胜出。
添加依赖与配置
Spring Boot 通过 spring-kafka 项目提供自动配置,官方文档里加的就是这个依赖:
<dependency>
<groupId>org.springframework.kafka</groupId>
<artifactId>spring-kafka</artifactId>
</dependency>
NoteKafka 没有专门的
spring-boot-starter-kafka,spring-kafka本身就是 starter 的角色,版本由 Spring Boot 统一管理。
配置集中在 spring.kafka.*:
spring:
kafka:
bootstrap-servers: "localhost:9092" # Broker 地址
consumer:
group-id: "myGroup" # 消费者组
Warning网上很多老教程写
localhost:2181,那是 ZooKeeper 的端口,早就过时了。Kafka 3.x 起用 KRaft 模式,不再需要 ZooKeeper,直连 9092 即可。
生产者:KafkaTemplate
KafkaTemplate 由自动配置创建,直接注入使用。它的 send 方法把消息发到指定 Topic:
package com.example.demo.service;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.stereotype.Service;
@Service
public class MessageProducer {
private final KafkaTemplate<String, String> kafkaTemplate;
public MessageProducer(KafkaTemplate<String, String> kafkaTemplate) {
this.kafkaTemplate = kafkaTemplate;
}
public void send(String topic, String message) {
// 第一个参数是 Topic,第二个是消息内容
kafkaTemplate.send(topic, message);
}
public void sendWithKey(String topic, String key, String message) {
// 带 key 发送:相同 key 的消息会进同一个分区,保证顺序
kafkaTemplate.send(topic, key, message);
}
}
send 返回 CompletableFuture,需要确认发送结果时可以异步等待。
Topic 不会自动创建。想在启动时建好,声明一个 NewTopic Bean:
package com.example.demo.config;
import org.apache.kafka.clients.admin.NewTopic;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@Configuration(proxyBeanMethods = false)
public class KafkaTopicConfig {
@Bean
public NewTopic orderTopic() {
// 名字、分区数、副本数
return new NewTopic("orders", 3, (short) 1);
}
}
Topic 已存在时这个 Bean 会被忽略,重复声明不会报错。
消费者:@KafkaListener
消费消息比发送更简单。在任意 Bean 的方法上标 @KafkaListener,指定 Topic 即可:
package com.example.demo.listener;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Component;
@Component
public class OrderListener {
// 直接收消息内容
@KafkaListener(topics = "orders")
public void onMessage(String message) {
System.out.println("收到订单消息:" + message);
}
// 也可以收完整的 ConsumerRecord,拿 key、分区、偏移量
@KafkaListener(topics = "orders")
public void onRecord(ConsumerRecord<String, String> record) {
System.out.println("key=" + record.key()
+ ",分区=" + record.partition()
+ ",offset=" + record.offset()
+ ",内容=" + record.value());
}
}
消费组从配置里读:spring.kafka.consumer.group-id。同一个组内,一条消息只被消费一次;想要多个系统各拿一份,用不同的组名。
监听容器(负责拉取消息、调用监听方法)由自动配置创建,参数在 spring.kafka.listener.* 下调整,比如并发数和提交方式:
spring:
kafka:
listener:
concurrency: 3 # 并发消费线程数,别超过分区数
Tip消费者默认从最新消息开始读(
auto.offset.reset=latest)。想消费历史消息,改成earliest;想手动指定位置,在监听方法里用@Header(KafkaHeaders.OFFSET)拿偏移量。
默认情况下,消息处理成功后偏移量自动提交,这叫自动确认。业务上要求「处理成功才提交」时,改成手动确认:
spring:
kafka:
listener:
ack-mode: manual # 手动提交偏移量
@Component
public class OrderListener {
@KafkaListener(topics = "orders")
public void onOrder(Order order, Acknowledgment ack) {
try {
System.out.println("处理订单:" + order.id());
ack.acknowledge(); // 处理成功后才提交偏移量
} catch (Exception e) {
// 不确认,消息会重新投递
}
}
}
手动确认的代价是代码变多,换来的是「不丢消息」的承诺:处理中途崩溃,偏移量没提交,重启后消息还会再消费一遍。
JSON 消息
业务消息很少是纯字符串,更多是对象。spring-kafka 内置了 Jackson 序列化器,配置一下就能收发 JSON:
spring:
kafka:
producer:
value-serializer: org.springframework.kafka.support.serializer.JacksonJsonSerializer
consumer:
value-deserializer: org.springframework.kafka.support.serializer.JacksonJsonDeserializer
properties:
# 反序列化成哪个类型
"[spring.json.value.default.type]": com.example.demo.Order
# 允许反序列化的包,防止反序列化攻击
"[spring.json.trusted.packages]": com.example.demo
定义消息对象:
package com.example.demo;
import java.math.BigDecimal;
public record Order(Long id, String item, BigDecimal amount) {
}
生产者直接发对象:
@Service
public class OrderProducer {
private final KafkaTemplate<String, Order> kafkaTemplate;
public OrderProducer(KafkaTemplate<String, Order> kafkaTemplate) {
this.kafkaTemplate = kafkaTemplate;
}
public void createOrder(Order order) {
kafkaTemplate.send("orders", order);
}
}
消费者方法参数直接写对象类型:
@Component
public class OrderListener {
@KafkaListener(topics = "orders")
public void onOrder(Order order) {
System.out.println("处理订单:" + order.id() + "," + order.item());
}
}
Warning
spring.json.trusted.packages一定要配。Jackson 反序列化不认识的包默认拒绝,这是防止恶意消息的安全机制,别图省事配成*。
幂等与顺序
Kafka 的投递语义是「至少一次」:消息可能被重复消费。比如消费者处理完还没来得及提交偏移量就崩溃,重启后同一条消息会再来一遍。所以消费端代码要设计成幂等的——同一个订单处理两次,结果和一次一样。常见做法是用业务主键去重:处理前查一下记录,存在就跳过。
顺序性也有讲究:同一个 key 的消息保证进同一个分区,分区内严格有序。想保证「同一个用户的订单消息按序处理」,发送时用用户 ID 做 key;消费端并发数不能超过分区数,否则一个分区被多个线程抢着消费,顺序就乱了。
事务与高级话题
消息发送和数据库操作要一起成功、一起失败时,用 Kafka 事务。配置事务前缀后,Spring Boot 自动创建 KafkaTransactionManager:
spring:
kafka:
producer:
transaction-id-prefix: "tx-" # 声明使用事务
配合 @Transactional 注解,方法内数据库写和 Kafka 发送就绑进同一个事务。
Kafka Streams 是另一座大山:用 @EnableKafkaStreams 开启,通过 StreamsBuilder 对消息做流式处理(过滤、聚合、Join)。它处理的是「持续流动的数据」,适合实时统计场景。
本地跑起来
本地开发需要先启动一个 Kafka。用 Docker Compose 最省事:
services:
kafka:
image: apache/kafka:latest
ports:
- "9092:9092"
environment:
KAFKA_NODE_ID: 1
KAFKA_PROCESS_ROLES: "broker,controller"
KAFKA_LISTENERS: "PLAINTEXT://:9092,CONTROLLER://:9093"
KAFKA_ADVERTISED_LISTENERS: "PLAINTEXT://localhost:9092"
KAFKA_CONTROLLER_LISTENER_NAMES: "CONTROLLER"
KAFKA_CONTROLLER_QUORUM_VOTERS: "1@localhost:9093"
启动后运行应用,用生产者接口发消息,观察消费者日志,一套本地消息链路就通了。
写测试时可以用 @EmbeddedKafka 注解,它会在测试里启动一个内嵌 Kafka,不依赖外部环境:
package com.example.demo;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.kafka.test.context.EmbeddedKafka;
@SpringBootTest
@EmbeddedKafka(topics = "orders", bootstrapServersProperty = "spring.kafka.bootstrap-servers")
class OrderListenerTest {
// 测试代码
}
bootstrapServersProperty 把内嵌 Broker 的地址映射进 Spring Boot 配置,自动配置无需改动。
小结
Kafka 集成是 Spring Boot 里「配置最少、威力最大」的模块之一:加 spring-kafka 依赖,写两行配置,注入 KafkaTemplate 发消息、加 @KafkaListener 收消息。进阶方向记住三条:JSON 序列化要配 trusted packages、事务靠 transaction-id-prefix、本地调试用 Docker Compose 或 @EmbeddedKafka。