第 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-service 和 notification-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,成功后标记已发
团队里的务实分级:
- 能容忍偶发丢失(日志类通知):直接发,接受
send返回值告警; - 不能丢但量小:outbox 表 + 定时轮询投递,两百行代码以内搞定;
- 高吞吐且不能丢:Debezium CDC 直读 outbox 表,运维成本上一个台阶,值得为它单独立项。
小结
- 先判断动机(削峰/解耦/最终一致/耗时任务),用户等结果的链路不要消息化;
- JHipster 的 Kafka 脚手架价值在配置示范,业务消息统一走 Spring Cloud Stream 函数式模型;
- 本地 docker-compose、测试 Testcontainers,配置全部外置到 YAML;
- 可靠性三件套:生产
acks=all+ 幂等、消费重试 + DLT 告警、消费端幂等; - 顺序性靠分区键,事务边界靠 outbox——“事务提交后发事件”是底线。
系列导航
- 上一篇:REST 层设计——生成的代码够用吗
- 下一篇:Angular 端架构——生成的代码怎么读
- 相关阅读:微服务 vs 单体——怎么选 · 测试全景——单元到端到端 · 可观测性——日志、指标、链路