第 63 章:Kafka 与流处理
学习目标
- 理解 Kafka 的核心架构与设计哲学
- 掌握 Spring Kafka 集成与生产者消费者模式
- 学会 Kafka 在日志、事件溯源、CDC 中的应用
一、为什么需要 Kafka?
Kafka 的本质:分布式持久化日志。
| 维度 | RabbitMQ | Kafka |
|---|---|---|
| 定位 | 消息队列(任务分发) | 分布式日志(事件流) |
| 吞吐量 | 万级 | 百万级 |
| 消费模式 | 队列(消费即删除) | 日志(可重放,保留 N 天) |
| 延迟 | 微秒级 | 毫秒级 |
| 场景 | 业务异步、可靠投递 | 日志、事件溯源、CDC、流计算 |
二、Kafka 核心概念
| 概念 | 说明 |
|---|---|
| Producer | 生产者 |
| Consumer | 消费者 |
| Topic | 主题(消息分类) |
| Partition | 分区(Topic 的并行单位) |
| Offset | 消息在分区里的位置(类似数组下标) |
| Broker | Kafka 服务器节点 |
| Consumer Group | 消费者组(一个分区只被组内一个消费者消费) |
三、Spring Boot 集成
xml
<dependency>
<groupId>org.springframework.kafka</groupId>
<artifactId>spring-kafka</artifactId>
</dependency>yaml
spring:
kafka:
bootstrap-servers: localhost:9092,localhost:9093,localhost:9094
producer:
key-serializer: org.apache.kafka.common.serialization.StringSerializer
value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
acks: all # ① 最强一致性
retries: 3
properties:
enable.idempotence: true # ② 幂等生产者(防重复)
max.in.flight.requests.per.connection: 5
linger.ms: 20 # 批量发送延迟
compression.type: snappy # 压缩
consumer:
group-id: taskflow-consumer
auto-offset-reset: earliest
enable-auto-commit: false # ③ 手动提交 offset
max-poll-records: 100
key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer
properties:
spring.json.trusted.packages: "com.taskflow.**"
listener:
ack-mode: manual_immediate
concurrency: 3 # 并发消费者数四、生产者
java
@Component
@RequiredArgsConstructor
@Slf4j
public class UserEventProducer {
private final KafkaTemplate<String, Object> kafkaTemplate;
// ① 发送事件(key 保证同一用户的消息有序)
public void publishUserEvent(UserEvent event) {
String key = String.valueOf(event.getUserId()); // ① 同用户的事件进同一分区
CompletableFuture<SendResult<String, Object>> future =
kafkaTemplate.send("user.events", key, event);
// ② 回调处理发送结果
future.whenComplete((result, ex) -> {
if (ex != null) {
log.error("发送失败", ex);
// 失败重试或入本地消息表
} else {
log.info("发送成功 partition={} offset={}",
result.getRecordMetadata().partition(),
result.getRecordMetadata().offset());
}
});
}
// ② 同步发送(业务关键场景)
public void publishSync(UserEvent event) {
try {
SendResult<String, Object> result =
kafkaTemplate.send("user.events", String.valueOf(event.getUserId()), event)
.get(5, TimeUnit.SECONDS);
log.info("同步发送成功 offset={}",
result.getRecordMetadata().offset());
} catch (Exception e) {
log.error("同步发送失败", e);
throw new BusinessException("事件发布失败");
}
}
}KafkaTemplate 的事务支持
java
@Configuration
public class KafkaConfig {
@Bean
public KafkaTemplate<String, Object> kafkaTemplate(
ProducerFactory<String, Object> factory) {
KafkaTemplate<String, Object> template = new KafkaTemplate<>(factory);
// ① 开启事务(生产消息 + DB 操作可以事务化)
template.setTransactionIdPrefix("taskflow-tx-");
return template;
}
}
// 用法
@Transactional("kafkaTransactionManager")
public void placeOrder(OrderDTO dto) {
orderMapper.insert(order);
kafkaTemplate.send("order.events", order); // 与 DB 同事务
}五、消费者
java
@Component
@Slf4j
public class UserEventConsumer {
// ① 监听主题
@KafkaListener(
topics = "user.events",
groupId = "user-service",
concurrency = "3"
)
public void onMessage(
ConsumerRecord<String, UserEvent> record,
@Header(KafkaHeaders.RECEIVED_TOPIC) String topic,
@Header(KafkaHeaders.RECEIVED_PARTITION) int partition,
@Header(KafkaHeaders.OFFSET) long offset,
Acknowledgment ack) {
try {
log.info("收到消息 topic={} partition={} offset={} key={} value={}",
topic, partition, offset, record.key(), record.value());
UserEvent event = record.value();
processEvent(event);
// ② 手动提交 offset
ack.acknowledge();
} catch (Exception e) {
log.error("处理失败", e);
// ③ 不提交 offset,下次重新消费
// 或转入死信主题
kafkaTemplate.send("user.events.dlq", record.value());
ack.acknowledge(); // 即使失败也 ack,避免死循环
}
}
private void processEvent(UserEvent event) {
switch (event.getType()) {
case REGISTER -> sendWelcomeEmail(event);
case LOGIN -> updateLastLogin(event);
case PURCHASE -> awardPoints(event);
}
}
}批量消费
yaml
spring:
kafka:
listener:
type: batch # ① 批量消费
ack-mode: manualjava
@KafkaListener(topics = "user.events", groupId = "user-batch-consumer")
public void onBatch(List<UserEvent> events, Acknowledgment ack) {
log.info("批量收到 {} 条消息", events.size());
// ② 批量处理(数据库批量插入)
userEventMapper.batchInsert(events);
// ③ 一次性 ack
ack.acknowledge();
}六、消费者组与 Rebalance
关键规则:
- 一个分区只能被组内一个消费者消费
- 消费者数 > 分区数 → 多余的消费者闲置
- 增减消费者会触发 Rebalance,期间整个组暂停消费
七、Offset 管理
java
// ① 自动提交(不推荐)
enable-auto-commit: true
// 风险:处理失败但 offset 已提交 → 消息丢失
// ② 手动提交(推荐)
enable-auto-commit: false
ack.acknowledge(); // 处理成功才提交
// ③ 提交特定 offset
ack.acknowledge();
// ④ 从指定 offset 开始消费
@KafkaListener(topicPartitions = @TopicPartition(
topic = "user.events",
partitions = {"0", "1"},
partitionOffsets = @PartitionOffset(partition = "0", initialOffset = "100")
))
public void onMessage(...) { ... }八、Kafka 应用场景
1. 日志收集
yaml
# Filebeat → Kafka → Logstash → ES
filebeat.inputs:
- type: log
paths:
- /var/log/app/*.log
output.kafka:
hosts: ["localhost:9092"]
topic: app-logs2. 事件溯源(Event Sourcing)
java
// 不存当前状态,存所有事件,从头重放就能恢复状态
public void placeOrder(OrderDTO dto) {
OrderCreatedEvent event = new OrderCreatedEvent(...);
kafkaTemplate.send("order.events", event); // 存事件
}
// 查询订单状态:消费所有相关事件,重建状态
public Order getOrderState(Long orderId) {
Order order = new Order();
// 读所有 OrderCreatedEvent, OrderPaidEvent, OrderShippedEvent...
// 按顺序应用,重建订单状态
return order;
}3. CDC(Change Data Capture)
典型场景:Canal/Debezium 监听 MySQL binlog,把变更事件发到 Kafka,下游各系统各自消费。
4. 用户行为追踪
java
// 用户点击 → Kafka → 实时分析
@KafkaListener(topics = "user.click")
public void trackClick(ClickEvent event) {
// 实时统计热点内容
// 推荐系统特征计算
// 实时大屏
}九、Kafka 高级特性
消息顺序性
java
// 保证同一用户的事件有序:partition key 用 userId
kafkaTemplate.send("user.events", String.valueOf(event.getUserId()), event);
// 同 userId → 同一 partition → 严格有序Exactly-Once 语义
yaml
producer:
acks: all # ① Leader + 所有 ISR 都确认
properties:
enable.idempotence: true # ② 幂等(防重复)
transactional.id: "taskflow-tx-1" # ③ 事务 ID(跨分区原子)
consumer:
isolation-level: read_committed # ④ 只读已提交消息压缩
yaml
producer:
properties:
compression.type: snappy # snappy/lz4/gzip,吞吐量提升 2-3 倍十、Kafka 运维要点
监控指标
bash
# Kafka 自带命令
kafka-topics.sh --list --bootstrap-server localhost:9092
kafka-consumer-groups.sh --list --bootstrap-server localhost:9092
kafka-consumer-groups.sh --describe --group taskflow --bootstrap-server localhost:9092| 关键指标 | 健康值 |
|---|---|
| 消息积压(Lag) | < 1000 |
| 生产 TPS | 取决于业务 |
| 消费 TPS | 接近生产 TPS |
| 分区数 | = 业务并行度 |
常见问题
bash
# ① 消费者启动后不消费
# 检查 auto-offset-reset 配置
# ② Lag 一直涨
# 增加消费者实例(不超过分区数)
# 优化消费逻辑
# ③ Rebalance 频繁
# 调大 session.timeout.ms
# 优化处理逻辑,避免超时十一、本章小结
| 要点 | 关键 |
|---|---|
| 本质 | 分布式持久化日志(区别于消息队列) |
| 吞吐量 | 百万级 |
| 核心概念 | Topic / Partition / Offset / Consumer Group |
| 分区 | 并行单位,同 key 进同一分区保证有序 |
| 消费者组 | 一个分区只被组内一个消费者消费 |
| 可靠投递 | acks=all + idempotence + transaction |
| 幂等消费 | enable.auto.commit=false + 手动 ack |
| 应用场景 | 日志、事件溯源、CDC、用户行为、实时分析 |
动手练习
练习 1:基础题
实现用户注册事件:服务 A 写入数据库后发 Kafka 消息,服务 B 消费消息发送欢迎邮件。
练习 2:进阶题
实现订单事件的 Exactly-Once:生产者事务 + 数据库操作,消费者手动提交 offset + 业务幂等。
练习 3:思考题
你的系统每天产生 1 亿条用户行为日志。如何设计 Kafka 集群和消费者,保证数据不丢、可重放、可被多个下游系统独立消费?