跳转到主要内容

Kafka 在交易系统的削峰与解耦

记录 Kafka 在交易系统中做削峰填谷与业务解耦的实践,以及顺序、幂等和积压处理上必须守住的细节。

交易系统引入 Kafka 通常出于两个动机:扛住流量尖峰,以及把非核心逻辑从主链路上摘下来。这两件事它都能做,但做法和注意点完全不同。

削峰:把洪峰变成队列

代发工资、批量还款、营销活动这些场景的流量特征是极不均匀——平时每秒几百笔,活动开始瞬间冲到几万笔。数据库扛不住的不是总量,是瞬时并发。

削峰的本质是用延迟换稳定:请求先落 Kafka,下游按自己的处理能力匀速消费。关键是消费端要限速,不能拿到消息就火力全开打数据库。

spring:
  kafka:
    consumer:
      group-id: txn-processor
      enable-auto-commit: false      # 必须手动提交
      max-poll-records: 100
      properties:
        max.poll.interval.ms: 300000  # 留足单批处理时间,避免误判掉线
    listener:
      ack-mode: manual
      concurrency: 8                  # 与分区数匹配,不要超过

concurrency 超过分区数没有任何收益,多出来的线程只会空转。要提升并行度得先加分区,而分区数一旦加了就减不回去,前期要按峰值容量规划好。

削峰只对"可以延后处理"的业务成立。用户在 App 里点转账等结果的场景不能这么做——那是同步链路,塞消息队列只是把等待转移到了别处。

解耦:从同步调用到事件

主链路上挂着一堆下游是很常见的技术债:转账成功后要发短信、更新积分、推送风控特征、写数仓。每加一个下游,主链路的失败面就大一分。

改成发一条 transfer.posted 事件,各下游自己订阅。主链路只负责记账和发事件,下游挂了不影响交易成功。

需要划清界限的是:哪些下游能异步。判断标准是这个下游失败后,业务上能不能接受"稍后补上"。更新积分可以,扣减额度不行——额度是交易准入条件,必须同步。

几个必须守住的细节

  • 顺序只在分区内保证。同一账户的消息必须用账户号做 key,否则先扣款后入账的顺序可能颠倒。跨账户不需要全局顺序,别为了顺序把分区数设成 1。
  • 消费端必须幂等。Kafka 是至少一次投递,重复是常态。用业务唯一键加唯一索引兜住,不要指望不重复。
  • 先处理成功再提交 offset。自动提交会在处理失败时丢消息。手动提交虽然会带来重复,但重复有幂等兜着,丢了就找不回来了。
  • 积压要能定位到分区。整体 lag 正常但单分区堆积,通常是某个 key 数据倾斜或者某条消息反复失败阻塞了分区。
# 按分区查看消费延迟,定位倾斜与阻塞
kafka-consumer-groups.sh --bootstrap-server "$BROKERS" \
  --group txn-processor --describe

# 排查反复失败的毒消息:从指定 offset 读一条看内容
kafka-console-consumer.sh --bootstrap-server "$BROKERS" \
  --topic txn-events --partition 3 --offset 88213 --max-messages 1

毒消息一定要有出路。处理失败超过阈值就转到死信 topic,让分区继续往下走,否则一条脏数据能把整个分区卡死几个小时。

Kafka 在交易系统里是很好的缓冲层和事件总线,但它不解决一致性问题。消息发出去了业务却回滚了,或者业务成功了消息没发出去,这些都要靠本地消息表和日终对账兜。把它当队列用,别当事务用。