Kafka 集成:spring-kafka 消费与事务消息概念

Kafka 是分布式日志型消息系统:消息追加写入分区(Partition),按 offset 顺序消费,天然适合海量事件、削峰与数据管道。Spring 侧用 spring-kafka 集成(版本由 Boot BOM 管理)。本章给消费/生产要点与事务消息的概念框架。

依赖与配置

<dependency>
    <groupId>org.springframework.kafka</groupId>
    <artifactId>spring-kafka</artifactId>
</dependency>
spring:
  kafka:
    bootstrap-servers: 127.0.0.1:9092
    consumer:
      group-id: order-group            # 消费组:组内一条消息只被一个实例消费
      auto-offset-reset: earliest      # 无 committed offset 时从最早读
      enable-auto-commit: false        # 关闭自动提交,由容器按 ack-mode 提交
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer
      properties:
        spring.json.trusted.packages: "com.example.*"   # JsonDeserializer 反序列化白名单
    listener:
      ack-mode: record                 # 每条记录处理成功即提交 offset(At-least-once 常用)

生产与消费

@Service
public class OrderEventProducer {
    private final KafkaTemplate<String, Object> kafkaTemplate;   // 构造注入

    public void publish(OrderCreatedEvent event) {
        kafkaTemplate.send("order-events", event.getOrderId(), event);
        // 输出:写入 order-events 主题,key=orderId 保证同单进同一分区保序
    }
}
@Component
public class OrderEventListener {

    @KafkaListener(topics = "order-events", groupId = "order-group")
    public void onEvent(OrderCreatedEvent event) {
        handle(event);
        // 输出:ack-mode=record,本方法正常返回即提交该条 offset
    }
}

要点:重复消费是常态(重平衡、重试、宕机都可能重投),消费者必须幂等(去重见第 41 章);处理抛异常时按重试策略重投,超过上限进 DLT(死信主题)。

事务消息概念

"本地库写成功但 Kafka 没发出去"是分布式一致性的经典缺口,Kafka 提供事务能力解决:

  1. 普通发送:无事务,失败靠重试,可能重复;
  2. 生产者事务:配置 spring.kafka.producer.transaction-id-prefix: tx-order- 后,用 kafkaTemplate.executeInTransaction(ops -> ...) 或在 @Transactional + KafkaTransactionManager 组合下发送——同批消息要么全进 Kafka,要么全不进;
  3. 读-处理-写(read-process-write):消费 → 处理 → 再发送的整条链路放进事务,配合幂等生产者与 isolation.level=read_committed 消费端,才能逼近端到端精确一次(Exactly-Once);
  4. Transactional Outbox 模式:不直接用 Kafka 事务时,把"发消息"变成"业务库里插一条 outbox 记录"(与业务同事务),后台任务扫表再可靠发布——用数据库事务兜底消息可靠性。

工程提醒:端到端精确一次代价高,多数业务用"至少一次 + 消费者幂等"就够;Kafka 事务与幂等的具体参数语义以官方文档为准。

小结:spring-kafka 用 KafkaTemplate 发、@KafkaListener 收,关闭自动提交配合 ack-mode 控制 offset;事务消息解决"本地操作与发消息不一致",生产常用 outbox 或至少一次 + 幂等消费。

笔记加载中…