Appearance
Kafka
这篇笔记聚焦 Kafka 这个具体产品,重点放在 Topic、Partition、Offset、Consumer Group、吞吐、顺序、回放、Spring Boot 集成 这些最常和 Kafka 一起出现的主线问题上。
如果你还没建立消息队列这条知识线,可以看:
这篇笔记不会把 Kafka 写成“另一种普通队列”,而是重点强调:Kafka 更像一个高吞吐、可保留、可回放的事件日志流平台。
1. Kafka 是什么
Kafka 可以看成 一个面向分布式日志流和事件流的消息平台。
这里的重点不是“发一条消息”,而是:
- 持续承接大量事件写入
- 把同一份数据分发给多个下游系统
- 允许不同消费者按自己的进度消费
- 在一段时间内保留历史消息,支持回放
所以 Kafka 更常被用在这些语境里:
- 日志采集
- 用户行为埋点
- 订单、支付、风控等业务事件流
- 实时计算和流式处理前置链路
它不是简单拿来和 RabbitMQ 做“一换一”的替代,而是:更适合对吞吐、日志保留、分区扩展和事件回放更敏感的场景。
2. 为什么 Kafka 常和日志流、事件流放在一起讨论
很多业务消息系统更像:来了一个任务,尽快投给下游处理。
Kafka 的思路更像:把一连串持续产生的事件,按顺序不断写进日志,再让不同消费者按各自节奏读取。
这会带来几个很重要的差异:
- 消息不会因为某个消费者读过就立刻消失
- 不同消费者组可以独立消费同一份数据
- 消费进度是“消费者自己记到哪里了”
- 历史数据在保留期内还可以重放
也正因为这样,Kafka 天然适合:
- 一份数据被多个系统复用
- 需要回溯、补算、重放历史消息
- 需要持续处理大规模数据流
3. Kafka 里最核心的几个概念
3.1 Broker
Broker 可以理解成 真正负责接收、存储和分发消息的 Kafka 服务节点。
如果是本地学习环境,常见地址就是:
localhost:9092
3.2 Topic
Topic 可以理解成 同一类消息的逻辑主题。
例如:
order-eventspayment-eventsuser-behavioraudit-logs
3.3 Partition
Partition 是 Kafka 最核心的概念之一。
它可以理解成 一个 Topic 在物理上拆成的多个分区。
它主要解决两件事:
- 提升并行处理能力
- 让同一个业务 key 可以稳定落到同一个分区里,从而维持局部顺序
3.4 Offset
Offset 可以理解成 消息在某个分区中的位置编号。
它表达的是:
- 这条消息在当前分区里排第几个
- 消费者已经处理到哪里
要注意的是:Kafka 更常是“消息保留在日志里,消费者自己记录进度”,而不是“消息被消费后就立刻删掉”。
3.5 Producer
Producer 就是消息生产者。
谁真正把消息写进 Kafka,谁就是 Producer。
它可能是:
- 订单服务
- 支付服务
- 日志接入服务
- 采集日志的 Agent
3.6 Consumer
Consumer 就是消息消费者。
谁从 Kafka 拉取消息并处理,谁就是 Consumer。
它的处理动作常包括:
- 落库
- 发告警
- 做清洗转换
- 写搜索系统
- 转发给别的下游
3.7 Consumer Group
Consumer Group 可以理解成 一组协同消费同一个 Topic 的消费者实例。
它最值得记住的规则是:
- 同一个组内,同一个分区同一时刻只会分给一个消费者
- 不同组之间互不影响,可以各自独立消费同一份数据
这也是 Kafka 为什么特别适合 一份事件,多个系统,各自独立处理。
4. 一条 Kafka 消息通常是怎么流转的
看这张整体图:
mermaid
flowchart LR
P[Producer 发送消息] --> T[Topic]
T --> P0[Partition 0]
T --> P1[Partition 1]
T --> P2[Partition 2]
P0 --> G1[Consumer Group A]
P1 --> G1
P2 --> G1
P0 --> G2[Consumer Group B]
P1 --> G2
P2 --> G2这条链路真正表达的是:
- 生产者不是把消息直接发给某个消费者
- 而是先写进
Topic Topic再按分区存储消息- 不同消费者组各自维护自己的消费进度
所以同一份 order-events 数据,可以同时被:
- 订单分析服务消费
- 审计服务消费
- 风控服务消费
- 搜索索引服务消费
5. Kafka 的顺序、吞吐和回放应该怎么理解
5.1 顺序:更常保证分区内顺序
Kafka 经常被说“支持顺序消息”,但更贴近实际的是:Kafka 更常保证分区内顺序,而不是全局顺序。
也就是说:
- 同一个分区里的消息天然有序
- 不同分区之间不保证全局有序
所以如果业务要求:同一个订单的事件必须按顺序处理
更常见的做法是:让同一个 orderId 作为 message key,这样它更稳定地落到同一个分区。
5.2 吞吐:为什么 Kafka 通常更高
Kafka 的高吞吐,通常和这些设计有关:
- 顺序追加写入
- 分区带来的并行扩展能力
- 批量发送和批量拉取
- 更偏日志流的模型设计
它不是靠某一个技巧“突然变快”,而是:存储模型、分区模型和批处理模型一起服务于高吞吐。
5.3 回放:为什么 Kafka 特别适合补消费
Kafka 很重要的一点是:消息在保留期内通常还在,消费者只是记录自己读到哪里。
所以当出现下面这些场景时,它会比较自然:
- 新增一个消费者组,重新消费历史消息
- 某个消费者故障恢复后,从上次位点继续
- 因为业务 bug 需要补算历史数据
这也是它和很多“处理完就结束”的业务队列非常不一样的地方。
6. Kafka 更适合哪些场景,不适合哪些场景
更适合
- 日志采集
- 埋点数据流
- 大规模业务事件流
- 实时计算前置链路
- 一份数据被多个下游独立消费
不那么适合
- 特别依赖复杂路由规则的业务消息
- 更强调交换器模型和细粒度投递控制的场景
- 需要把“一条消息怎么分发到多个队列”表达得非常细的场景
如果业务重点是:消息怎么按复杂路由规则投给多个下游
RabbitMQ 往往更自然。
7. 企业里常见的 Kafka 使用形态
结合实际项目,Kafka 常见的几种形态可以压缩成下面三类。
7.1 业务系统直接写 Kafka
这种方式最常见于事件驱动架构:
- 订单服务写
order-events - 支付服务写
payment-events - 下游审计、风控、报表系统各自订阅
它的优点是:
- 解耦清晰
- 一份事件可以被多个系统复用
- 扩展性较好
7.2 日志文件或标准输出先被 Agent 采集,再发 Kafka
这种方式更常见于日志平台:
- 应用先写本地日志
Filebeat、Fluent Bit、Vector等 Agent 采集- 再统一写入 Kafka
它更适合:
- 老系统较多
- 多语言混合环境
- 日志以文件形态存在
7.3 Kafka 前面接入事件,后面再接处理层和落库层
这种方式更像成熟企业架构:
- 上游系统把原始事件写进 Kafka
- 中间用
Flink、Kafka Streams或自定义服务做清洗和加工 - 最后再落到
MySQL、Elasticsearch、ClickHouse等系统
它更适合:
- 实时报表
- 实时风控
- 实时告警
- 大规模日志分析
8. Spring Boot 项目里怎么整合 Kafka
如果是 Spring Boot 项目,最常见的目标通常是:
- 项目能连接 Kafka
- 能声明或初始化 Topic
- 能写生产者
- 能写消费者
- 能明确消息模型和消费组配置
8.1 依赖通常长什么样
最基础的一组依赖通常包括:
xml
<dependencies>
<dependency>
<groupId>org.springframework.kafka</groupId>
<artifactId>spring-kafka</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-validation</artifactId>
</dependency>
<dependency>
<groupId>com.fasterxml.jackson.datatype</groupId>
<artifactId>jackson-datatype-jsr310</artifactId>
</dependency>
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
<optional>true</optional>
</dependency>
</dependencies>这几项分别在解决:
spring-kafka:Spring Boot 和 Kafka 的核心集成能力validation:校验自定义配置类jackson-datatype-jsr310:更自然地处理LocalDateTime、Instant等时间类型lombok:减少样板代码
8.2 application.yml 通常怎么配
一个更贴近 Spring Boot 工程的基础配置通常会是这样:
yml
spring:
kafka:
bootstrap-servers: localhost:9092
producer:
key-serializer: org.apache.kafka.common.serialization.StringSerializer
value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
acks: all
consumer:
group-id: order-analytics-group
auto-offset-reset: earliest
key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer
properties:
spring:
json:
trusted:
packages: com.example.demo.messaging
value:
default:
type: com.example.demo.messaging.OrderPaidEvent
listener:
ack-mode: record
app:
kafka:
topic:
order-events: order.events
partitions: 3
replicas: 1这里最值得记住的几个配置是:
bootstrap-servers:Kafka Broker 地址producer:生产者序列化和确认策略consumer.group-id:当前消费者属于哪个消费组auto-offset-reset:没有已提交位点时,从哪里开始消费listener.ack-mode:监听容器按什么粒度确认消费
8.3 Topic 声明代码通常长什么样
在 Spring Boot 项目里,常见做法是用 NewTopic 在启动阶段声明 Topic 元数据:
java
@Configuration
public class KafkaTopicConfiguration {
/**
* orderEventsTopic
* 功能:声明订单事件 Topic。
* 核心职责:
* 1. 在应用启动时向 KafkaAdmin 注册 Topic 元数据
* 2. 统一管理 Topic 名称、分区数和副本数
* 参数:
* - kafkaTopicProperties: 自定义 Kafka Topic 配置
* 返回值:
* - NewTopic: 订单事件 Topic 定义
*/
@Bean
public NewTopic orderEventsTopic(KafkaTopicProperties kafkaTopicProperties) {
return TopicBuilder
.name(kafkaTopicProperties.getOrderEvents())
.partitions(kafkaTopicProperties.getPartitions())
.replicas(kafkaTopicProperties.getReplicas())
.build();
}
}这段代码真正表达的是:Topic 不只是一个字符串名字,它还包含分区数、副本数这些很关键的元数据。
8.4 消息模型通常长什么样
如果项目走 JSON 消息体,通常会有一个明确的事件模型:
java
@Data
@Builder
@NoArgsConstructor
@AllArgsConstructor
public class OrderPaidEvent {
/**
* 事件唯一标识,用于区分每一条订单支付事件。
*/
private String eventId;
/**
* 订单号。工程上常常也会把它作为 Kafka message key。
*/
private String orderId;
/**
* 用户 ID,用于标识当前事件属于哪个用户。
*/
private Long userId;
/**
* 支付金额。
*/
private BigDecimal amount;
/**
* 事件发生时间。
*/
private LocalDateTime paidAt;
}这里把 orderId 单独拎出来很重要,因为它经常承担两层角色:
- 消息体里的业务字段
message key
第二层角色会直接影响分区路由和局部顺序。
8.5 生产者侧通常长什么样
Spring Boot 项目里最常见的生产者写法,是通过 KafkaTemplate 发送事件:
java
@Service
@RequiredArgsConstructor
public class OrderEventProducer {
private final KafkaTemplate<String, OrderPaidEvent> kafkaTemplate;
private final KafkaTopicProperties kafkaTopicProperties;
/**
* sendOrderPaidEvent
* 功能:发送订单支付成功事件。
* 核心职责:
* 1. 组装业务事件对象
* 2. 指定 Topic 和 message key
* 3. 把支付成功这件事发布到 Kafka
* 参数:
* - orderId: 订单号
* - userId: 用户 ID
* - amount: 支付金额
* 返回值:
* - void
*/
public void sendOrderPaidEvent(String orderId, Long userId, BigDecimal amount) {
OrderPaidEvent event = OrderPaidEvent.builder()
.eventId(UUID.randomUUID().toString())
.orderId(orderId)
.userId(userId)
.amount(amount)
.paidAt(LocalDateTime.now())
.build();
// 让同一个订单号稳定作为 message key,有助于维持局部顺序。
kafkaTemplate.send(kafkaTopicProperties.getOrderEvents(), orderId, event);
}
}这段生产者代码真正想说明的,不只是“怎么发消息”,而是:
- 生产者通常放在应用服务层
- 业务事件要有明确模型
message key往往要和业务顺序要求绑在一起考虑
8.6 消费者侧通常长什么样
消费者最常见的写法,是用 @KafkaListener 监听 Topic:
java
@Slf4j
@Service
@RequiredArgsConstructor
public class OrderAnalyticsConsumer {
private final OrderAnalyticsService orderAnalyticsService;
/**
* handleOrderPaid
* 功能:消费订单支付成功事件,并写入分析系统或下游存储。
* 核心职责:
* 1. 监听订单事件 Topic
* 2. 获取消息体以及分区、位点等元数据
* 3. 完成业务处理后返回,由监听容器按配置确认消费
* 参数:
* - event: 反序列化后的订单支付事件
* - partition: 当前消息所在分区
* - offset: 当前消息的位点
* - messageKey: 当前消息的 key
* 返回值:
* - void
*/
@KafkaListener(
topics = "${app.kafka.topic.order-events}",
groupId = "order-analytics-group",
concurrency = "3"
)
public void handleOrderPaid(
OrderPaidEvent event,
@Header(KafkaHeaders.RECEIVED_PARTITION) int partition,
@Header(KafkaHeaders.OFFSET) long offset,
@Header(value = KafkaHeaders.RECEIVED_KEY, required = false) String messageKey) {
log.info(
"收到订单支付事件 orderId={}, key={}, partition={}, offset={}",
event.getOrderId(),
messageKey,
partition,
offset
);
orderAnalyticsService.recordPaidEvent(event);
}
}这里最值得注意的几个点是:
topics:监听哪个 TopicgroupId:当前消费者属于哪个消费组concurrency = "3":监听容器允许的并发消费者数量partition / offset / key:这些元数据对排查顺序、定位消费进度很有帮助
8.7 这些 Spring Boot 代码到底表达了什么
把前面的 Topic、Producer、Consumer 串起来看,完整链路其实是:
mermaid
flowchart LR
A[订单服务完成支付] --> B[OrderEventProducer]
B --> C[Kafka order.events Topic]
C --> D[OrderAnalyticsConsumer]
D --> E[OrderAnalyticsService]这条链路真正表达的是:
- 业务系统负责把“发生了什么事”发布成事件
- Kafka 负责承接和分发这份事件
- 消费者组按自己的节奏处理消息
- 下游服务再决定落库、统计、告警还是继续转发
9. Kafka 集成里最常见的工程关注点
9.1 分区数不是越多越好
分区数会影响:
- 并行消费上限
- 顺序粒度
- 运维和管理复杂度
所以它不是越大越高级,而是要看:业务吞吐、消费者并发和顺序要求到底是什么。
9.2 message key 设计会直接影响顺序
如果业务要求同一个订单、有同一个用户或同一个设备的消息尽量有序,通常要:
- 明确业务 key
- 让这个 key 参与分区路由
如果不认真设计 message key,很容易出现:
- 消息被均匀打散到多个分区
- 吞吐上去了
- 但同一业务主体的顺序丢了
9.3 消费成功不等于业务完全没风险
即使 @KafkaListener 收到了消息,也不等于所有问题都解决了。
还要继续考虑:
- 业务处理是否幂等
- 异常时怎么重试
- 消费位点什么时候提交更合适
- 下游存储失败时怎么补偿
Kafka 很强,但它不会自动替业务解决所有一致性问题。
10. 一句话总结
Kafka 最值得记住的是:
它更像一个高吞吐、可分区、可保留、可回放的事件日志流平台。放到 Spring Boot 项目里,最常见的落地方式就是:用 Topic 承接事件、用 KafkaTemplate 发送消息、用 @KafkaListener 消费消息,再围绕分区、位点、消费组和 message key 去理解吞吐、顺序和工程边界。