财务对账少了 3 笔积分
2 月中旬,财务同学找过来:2 月 10 号这天的积分发放和订单数据对不上,少了 3 笔。
我们链路是:订单服务发 ORDER_PAID 事件到 Kafka,积分服务消费后发积分。查了一圈,积分服务没报错,是消息压根没到。
先看生产者的配置:
spring:
kafka:
producer:
acks: 1
retries: 0
然后翻 2 月 10 号 broker 的日志,找到了:
[2021-02-10 03:12:41,882] INFO [ReplicaFetcherManager on broker 2] ...
[2021-02-10 03:12:44,118] WARN [Controller-2-to-broker-1-send-thread] ...
[2021-02-10 03:12:47,301] INFO [KafkaServer id=2] started (kafka.server.KafkaServer)
[2021-02-10 03:12:48,004] INFO [GroupCoordinator 2]: Elected leader for group ...
broker 1 在 03:12 挂了(后来查是宿主机内存 OOM 被 kill),broker 2 接管了它上面的 leader 分区。
丢消息的机制是这样的:acks=1 表示只要 leader 写入本地日志就返回成功,不等 follower 同步。如果 leader 写完就宕机,follower 还没来得及拉到这条消息,新 leader 上位后这条消息就永久消失了——因为新 leader 的日志里根本没有它。
先搞清楚三个概念
ISR(In-Sync Replicas)
每个分区的副本里,leader 维护一个"同步副本集合"。判断标准是 replica.lag.time.max.ms(默认 10 秒,Kafka 2.4 里已经改成 30 秒),follower 超过这个时间没向 leader 发 fetch 请求,就被踢出 ISR。
注意 Kafka 0.9 之后不再用"落后多少条消息"判断,改用时间,因为突发流量下消息数滞后是正常的,用条数会误杀。
$ bin/kafka-topics.sh --bootstrap-server 10.0.0.41:9092 \
--describe --topic ORDER_EVENT_TOPIC
Topic: ORDER_EVENT_TOPIC Partition: 0 Leader: 2 Replicas: 2,1,3 Isr: 2,3
这行输出里 Replicas: 2,1,3 是全部副本,Isr: 2,3 是同步副本——broker 1 已经掉队了。看到 Isr 数量少于 Replicas 数量就说明有问题。
High Watermark
分区里有个"高水位"的概念:消费者只能看到 high watermark 之前的消息。HW 取的是所有 ISR 副本中最小的 LEO(日志末端位移)。
leader LEO=100 已写 100 条
follower1 LEO=98
follower2 LEO=100
=> HW = 98,消费者最多看到第 98 条
acks 的三个取值
| acks | 语义 | 丢消息场景 | 我们的实测吞吐 |
|---|---|---|---|
0 | 发完不等任何确认 | 网络抖动、broker 没收到 | 238 MB/s |
1 | leader 写入即返回 | leader 宕机且 follower 未同步 | 210 MB/s |
all(-1) | 所有 ISR 副本都写入才返回 | 仅当 ISR 全部丢失 | 96 MB/s |
吞吐数据是在我们压测环境跑的,3 broker,1 KB 消息,lz4 压缩,batch 64 KB。
改 acks=all 就够了吗
不够。这是我们踩的第二个坑,也是很多文章没讲清楚的地方。
# broker 端默认配置
min.insync.replicas=1
min.insync.replicas 的含义是:当 ISR 里的副本数少于这个值时,producer 写入会直接失败。
如果只配 acks=all 但 min.insync.replicas=1,会发生什么?假设 3 副本的分区,2 个 follower 全挂了,ISR 只剩 leader 自己。此时 acks=all 依然返回成功(因为"所有 ISR 副本"就是 leader 一个)。leader 再一挂,消息照样丢。
这两个参数必须配对使用:
# broker 端(server.properties)
min.insync.replicas=2
replication.factor=3
unclean.leader.election.enable=false
# producer 端
acks=all
含义是:至少要有 2 个副本(leader + 1 个 follower)确认写入才算成功,且当 ISR 不足 2 个时拒绝写入。这样即使 leader 宕机,剩下那个 ISR 副本里一定有新 leader 的完整数据。
第三个参数 unclean.leader.election.enable=false 是 0.11.0 之后的默认值,它禁止"从非 ISR 副本里选 leader"。如果设成 true,数据落后的副本也能当选,那前面配的就全白费了。
代价也很明显:ISR 不足 2 个时分区不可写。这是主动选择"不可用"来保"不丢失",CAP 的经典取舍。我们的做法是把 UnavailablePartitions 纳入监控并告警。
开启幂等生产者
把 acks 改成 all、retries 从 0 改大之后,会引入新问题:重试导致的消息重复和乱序。
场景:producer 发了一批消息,broker 写入成功但 ack 网络丢了,producer 重试,同一批消息写两遍。更糟的是,如果 max.in.flight.requests.per.connection 大于 1,第二批可能比第一批先到,顺序就乱了。
解决办法是开启幂等生产者(Kafka 0.11 引入):
spring:
kafka:
producer:
acks: all
properties:
enable.idempotence: true
max.in.flight.requests.per.connection: 5
retries: 3
compression.type: lz4
linger.ms: 20
batch.size: 65536
开启后 broker 会给每个 producer 分配一个 PID,并为每个分区维护序列号,重复的消息会被丢弃,乱序会被拒绝。关键是 max.in.flight.requests.per.connection 最多只能是 5(这是幂等生产者能保证顺序的上限),超过会报错。
开启幂等的额外开销很小,我们实测吞吐从 96 MB/s 降到 92 MB/s,约 4%。这个代价完全可以接受。
一个必须知道的坑
如果你用的是 Spring Kafka 且手动创建了 ProducerFactory,enable.idempotence 必须显式设置。Spring Boot 在检测到 acks=all 且 retries > 0 时不会自动帮你开幂等(这个行为各版本还不一样),别指望默认值。
消费端的一致性:事务
我们的积分服务是"消费 ORDER_PAID → 发 POINT_GRANT"。这里有个经典问题:处理成功、发消息失败,或者发消息成功、offset 提交失败,都会导致不一致。
用 Kafka 事务解决:
@Bean
public KafkaTransactionManager<String, String> kafkaTransactionManager(
ProducerFactory<String, String> pf) {
return new KafkaTransactionManager<>(pf);
}
@KafkaListener(topics = "ORDER_PAID", groupId = "point-group")
@Transactional("kafkaTransactionManager")
public void consume(ConsumerRecord<String, String> record) {
OrderPaidEvent event = parse(record.value());
pointService.grant(event.getUserId(), event.getAmount()); // 数据库操作
kafkaTemplate.send("POINT_GRANT", buildGrantEvent(event));
// offset 提交和上面的 send 在同一个事务里
}
producer 端要配:
spring.kafka.producer.transaction-id-prefix=point-tx-
spring.kafka.producer.properties.enable.idempotence=true
transaction-id-prefix 必须配,Spring 会用它生成唯一的 transactional.id。这个 ID 的作用是"fencing"——同名的旧 producer(比如刚重启的旧实例)会被隔离掉,防止僵尸实例写数据。
消费端要把隔离级别设成读已提交,否则会读到回滚的消息:
spring.kafka.consumer.properties.isolation.level=read_committed
事务的代价不小,我们实测吞吐掉 21%,端到端延迟增加约 30 ms。所以只在积分、账务这类真的不能重复的地方用,普通的埋点日志不需要。
上线后的数据和监控
改完之后重新压测:
| 配置 | 吞吐 | P99 延迟 |
|---|---|---|
| acks=1, retries=0(原) | 210 MB/s | 46 ms |
| acks=all + min.insync=2 | 96 MB/s | 92 ms |
| 上面 + 幂等 | 92 MB/s | 95 ms |
| 上面 + 事务(仅积分链路) | 73 MB/s | 126 ms |
吞吐掉了 56%,但这是"用性能换可靠性"该付的价。为了补回吞吐,我们把分区数从 12 加到了 24(新建 topic 双写切流,不是原地扩),最终整体吞吐回到 178 MB/s。
补上的监控:
UnderReplicatedPartitions > 0 持续 2 分钟 → 警告
IsrShrinksPerSec > 0 持续 5 分钟 → 警告
UnavailablePartitions > 0 → 严重告警
ISR 数量 < replication.factor 持续 10 分钟 → 警告
UnavailablePartitions 是 min.insync.replicas=2 的直接后果,一定要告警,否则你会发现"写入全部失败"却没人知道。
小结
acks=1只保证 leader 写入,leader 宕机且 follower 未同步时就丢消息。我们 2 月 10 号那次就是这么丢的。acks=all必须配min.insync.replicas=2,否则 ISR 只剩 leader 一个时照样丢。unclean.leader.election.enable保持 false。- 重试会引入重复和乱序,开
enable.idempotence=true,同时max.in.flight.requests.per.connection不超过 5。开销约 4%。 - consume-transform-produce 场景用事务,配
transaction-id-prefix和isolation.level=read_committed。开销约 21%。 - 吞吐从 210 掉到 92 MB/s 不是终点,加分区能补回来。可靠性和性能不是二选一,是要多花钱(机器)。
- 把
UnderReplicatedPartitions和UnavailablePartitions纳入告警。
那 3 笔积分最后人工补发的。钱不多,但事后复盘时我说:如果一个系统的消息可靠性靠"broker 一般不会同时挂两台"来保障,那它迟早会出事,只是概率问题。这次概率恰好落在了财务能对出来的那 3 笔上。