Administrator
发布于 2021-02-17 / 8071 阅读
69

Kafka 消息可靠性配置:acks、ISR 与最少同步副本

财务对账少了 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
1leader 写入即返回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=allmin.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 改成 allretries 从 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 且手动创建了 ProducerFactoryenable.idempotence 必须显式设置。Spring Boot 在检测到 acks=allretries > 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/s46 ms
acks=all + min.insync=296 MB/s92 ms
上面 + 幂等92 MB/s95 ms
上面 + 事务(仅积分链路)73 MB/s126 ms

吞吐掉了 56%,但这是"用性能换可靠性"该付的价。为了补回吞吐,我们把分区数从 12 加到了 24(新建 topic 双写切流,不是原地扩),最终整体吞吐回到 178 MB/s。

补上的监控:

UnderReplicatedPartitions     > 0 持续 2 分钟 → 警告
IsrShrinksPerSec              > 0 持续 5 分钟 → 警告
UnavailablePartitions         > 0            → 严重告警
ISR 数量 < replication.factor 持续 10 分钟  → 警告

UnavailablePartitionsmin.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-prefixisolation.level=read_committed。开销约 21%。
  • 吞吐从 210 掉到 92 MB/s 不是终点,加分区能补回来。可靠性和性能不是二选一,是要多花钱(机器)。
  • UnderReplicatedPartitionsUnavailablePartitions 纳入告警。

那 3 笔积分最后人工补发的。钱不多,但事后复盘时我说:如果一个系统的消息可靠性靠"broker 一般不会同时挂两台"来保障,那它迟早会出事,只是概率问题。这次概率恰好落在了财务能对出来的那 3 笔上。

参考