CHARLIE SAYS

查理如是说
DATE 2026-08-24
THEME
SERIES / JHIPSTER / P-227 · JHipster 项目开发

JHipster 开发 12:异步与消息——Kafka 集成

第 11 篇结尾留了个钩子:千条以上的批量导入不该走同步 HTTP。这类需求的答案就是本篇的主角——消息队列。JHipster 在问答里提供 Kafka 选项,生成一套可用的脚手架;但我们团队真正花时间的地方,是把它从”能跑”改造成”敢在生产发消息”:可靠性、顺序性、与数据库事务的边界,每一条都踩过坑。

先想清楚:什么场景需要异步

不是所有调用都值得消息化,先过一遍这张决策表:

动机典型场景同步做法的问题消息化收益
削峰秒杀下单、批量导入流量打穿数据库连接池队列缓冲,消费端按自己的节奏处理
解耦下单后通知库存/积分/营销调用方要知道所有下游发布事件,下游各自订阅
最终一致跨服务数据同步分布式事务复杂脆弱事件驱动,允许秒级延迟
耗时任务报表生成、PDF 导出HTTP 超时、线程占用异步处理 + 状态回写

反过来说,需要即时读到结果的操作不适合消息化:下单扣库存如果 500ms 内看不到结果,用户体验崩塌;查询类请求永远走同步。我们的红线是:凡是”用户在等响应”的链路,消息只能出现在响应之后的副作用里。

flowchart LR
  C["Client"] -->|POST /api/orders| A["OrderResource"]
  A --> B["OrderService\n同步:落库 + 返回 201"]
  B -->|事务提交后发事件| K["Kafka\norder-created"]
  K --> D["库存服务"]
  K --> E["积分服务"]
  K --> F["通知服务"]
  style B fill:#e8f0fe
  style K fill:#fce8e6

注意图里的措辞:事务提交后发事件。这个顺序是本篇后半段的伏笔。

JHipster Kafka 选项生成了什么

构建问答里选择 Asynchronous messages using Apache Kafka,或者在 JDL 里加 deployment * with kafka(微服务场景)后,生成的东西包括:

pom.xml                                        // spring-kafka / spring-cloud-stream 依赖
src/main/docker/kafka.yml                      // 本地 Kafka 的 docker-compose
src/main/resources/config/application.yml      // kafka 配置段
src/main/java/.../config/KafkaProperties.java  // 配置绑定类
src/main/java/.../service/KafkaConsumer.java   // 消费示例
src/main/java/.../service/KafkaProducer.java   // 生产示例

生成的 Producer/Consumer 是”hello world”级别的样板:一个往默认 topic 发字符串的轮询任务,一个把收到的消息打日志。它的价值是示范配置怎么接KafkaProperties 怎么绑定、bootstrap servers 怎么读环境变量),不是拿来直接写业务。

我们的决策:保留 KafkaProperties 与 docker-compose,把示例的 Producer/Consumer 删掉,业务消息统一迁到 Spring Cloud Stream 的函数式模型上。理由下一节讲。

Spring Cloud Stream 函数式模型

JHipster 微服务 + Kafka 的组合默认引入 Spring Cloud Stream,单体选 Kafka 时也可以自己加上。函数式模型的核心:不写 @KafkaListener,只声明 bean

@Configuration
public class OrderStreamConfig {

    @Bean
    public Consumer<OrderCreatedEvent> orderCreatedConsumer(OrderInventoryHandler handler) {
        return event -> handler.handle(event);
    }

    @Bean
    public Function<OrderCreatedEvent, NotificationRequest> orderToNotification() {
        return event -> new NotificationRequest(event.customerId(), "订单已创建");
    }

    @Bean
    public Supplier<Flux<OrderEvent>> orderEventSupplier(EventPublisher publisher) {
        return () -> publisher.flux();
    }
}

三种 bean 对应三种角色:

Bean 类型角色典型用法
Consumer<T>订阅者处理某个事件
Function<T, R>处理再转发事件转换、内容 enrichment
Supplier<T>发布者定时/流式产生消息

绑定与 topic 全部声明在 YAML 里,代码零 Kafka 痕迹,换 binder(比如测试换 embedded)不动代码:

spring:
  cloud:
    function:
      definition: orderCreatedConsumer;orderToNotification
    stream:
      bindings:
        orderCreatedConsumer-in-0:
          destination: order-created
          group: inventory-service
          consumer:
            concurrency: 3
        orderToNotification-in-0:
          destination: order-created
          group: notification-transformer
        orderToNotification-out-0:
          destination: notification-request
      kafka:
        binder:
          brokers: ${KAFKA_BROKERS:localhost:9092}
        bindings:
          orderCreatedConsumer-in-0:
            consumer:
              enableDlq: true
              dlqName: order-created-dlq

注意 group 的设计:同一 topic 上 inventory-servicenotification-transformer 是不同消费组,各自拿到全量消息;组内多实例分摊分区。这是事件驱动架构的基本功。

发送侧不用 Supplier 轮询的话,直接用 StreamBridge

@Service
public class EventPublisher {

    private final StreamBridge streamBridge;

    public void publishOrderCreated(Order order) {
        OrderCreatedEvent event = new OrderCreatedEvent(order.getId(), order.getCustomerId(), order.getAmount());
        Message<OrderCreatedEvent> message = MessageBuilder.withPayload(event)
            .setHeader(KafkaHeaders.KEY, String.valueOf(order.getCustomerId()))  // 分区键
            .build();
        boolean sent = streamBridge.send("order-created", message);
        if (!sent) {
            throw new EventPublishException("Failed to publish order-created for order " + order.getId());
        }
    }
}

本地开发:docker-compose 与 Testcontainers

src/main/docker/kafka.yml 一条命令起本地环境:

docker compose -f src/main/docker/kafka.yml up -d
# 应用起来后,生产的消息可以在容器里直接验证:
docker compose -f src/main/docker/kafka.yml exec kafka kafka-console-consumer.sh \
  --bootstrap-server localhost:9092 --topic order-created --from-beginning

集成测试不起 docker-compose,用 Testcontainers 更干净(JHipster 生成的测试基建本来就带它):

@Testcontainers
@SpringBootTest
class OrderEventIT {

    @Container
    static KafkaContainer kafka = new KafkaContainer(DockerImageName.parse("apache/kafka:3.8.0"));

    @DynamicPropertySource
    static void kafkaProps(DynamicPropertyRegistry registry) {
        registry.add("spring.cloud.stream.kafka.binder.brokers", kafka::getBootstrapServers);
    }

    @Test
    void shouldConsumeOrderCreatedEvent() {
        // 发送事件 → Awaitility 等待消费者副作用 → 断言
    }
}

@DynamicPropertySource 把容器地址注进 binder 配置,测试代码与本地/生产配置互不干扰。断言异步结果用 Awaitility 轮询而不是 Thread.sleep——sleep 是测试抖动的头号来源。

可靠性:ack、重试与死信

生产环境必答的三连问:消息会丢吗?失败了怎么办?毒消息怎么处理?

生产侧防丢:开启幂等生产者与确认。acks: all 加上 retries,配合 broker 的 min.insync.replicas >= 2,能在 broker 故障切换时不丢已确认消息:

spring:
  cloud:
    stream:
      kafka:
        binder:
          producer-properties:
            acks: all
            enable.idempotence: true

消费侧重试:消费者的 ack 默认在处理成功后提交 offset。业务处理抛异常时消息会重投,但要配重试上限,否则一条毒消息能把整个分区卡死:

spring:
  cloud:
    stream:
      bindings:
        orderCreatedConsumer-in-0:
          consumer:
            max-attempts: 4
            back-off-initial-interval: 1000
            back-off-max-interval: 10000

死信队列(DLT):重试耗尽后进 order-created-dlq(上文配置已开)。团队纪律:DLT 必须有监控告警,必须有补投工具。我们踩过的坑:某次消费端反序列化失败,消息静默进 DLT 三天没人发现,等业务方来问”积分怎么没加”才翻出来。从那以后 DLT 深度成了比 CPU 还重要的告警指标。

幂等消费:Kafka 的语义是 at-least-once,重试意味着消费者可能收到重复消息。处理端按业务键去重(唯一约束、Redis setnx、状态机判断”已处理则跳过”)是标配,不是可选项。

顺序性:分区键说了算

Kafka 只保证分区内有序。想让”同一订单的事件按发生顺序处理”,就必须让它们进同一个分区——也就是发消息时带 key(上面的 KafkaHeaders.KEY 设了 customerId):

策略效果代价
不带 key 轮询吞吐最高,完全无序业务乱序风险
业务键做 key同键有序,跨键并行热点键倾斜(大客户)
单分区全局有序吞吐退化成串行,基本不选

另一个隐形杀手:max.in.flight.requests.per.connection > 1 且开启重试时,生产者内部可能乱序。开启 enable.idempotence 后这个隐患自动消除(它强制 in-flight 不乱序),这是我们又一次”白捡”的安全。

与业务事务的边界:outbox 思想

最容易翻车的组合:数据库事务里发消息。

@Transactional
public OrderDTO createOrder(OrderDTO dto) {
    Order order = orderRepository.save(toEntity(dto));
    streamBridge.send("order-created", toEvent(order));  // 危险!
    return toDto(order);
}

两个失败模式:事务回滚但消息已发出(下游凭空创建库存扣减);事务提交但 broker 不可用导致 send 抛异常(下单失败)。两个方向都不对,本质是数据库和 Kafka 没有 shared transaction

outbox 模式的思路:事务里只写本地表,消息的发送移出事务:

@Transactional
public OrderDTO createOrder(OrderDTO dto) {
    Order order = orderRepository.save(toEntity(dto));
    outboxRepository.save(new OutboxMessage(
        "order-created",
        String.valueOf(order.getId()),
        toJson(toEvent(order))     // 事件先落库,与业务同事务
    ));
    return toDto(order);
}
// 独立的轮询任务/CDC(Debezium)把 outbox 表的记录投递到 Kafka,成功后标记已发

团队里的务实分级:

  1. 能容忍偶发丢失(日志类通知):直接发,接受 send 返回值告警;
  2. 不能丢但量小:outbox 表 + 定时轮询投递,两百行代码以内搞定;
  3. 高吞吐且不能丢:Debezium CDC 直读 outbox 表,运维成本上一个台阶,值得为它单独立项。

小结

  • 先判断动机(削峰/解耦/最终一致/耗时任务),用户等结果的链路不要消息化;
  • JHipster 的 Kafka 脚手架价值在配置示范,业务消息统一走 Spring Cloud Stream 函数式模型;
  • 本地 docker-compose、测试 Testcontainers,配置全部外置到 YAML;
  • 可靠性三件套:生产 acks=all + 幂等、消费重试 + DLT 告警、消费端幂等;
  • 顺序性靠分区键,事务边界靠 outbox——“事务提交后发事件”是底线。

系列导航

← 软件工程 012:面向过程管理:Kanban方式 目录 网络协议 012:工具:netstat查看服务及监听端口详解 →
← 返回文章列表