这是《Spring Boot 与 Spring Cloud 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 消息用于服务解耦、削峰填谷和事件传播;调度用于定时任务、补偿任务和周期报表。两者都涉及重试、幂等、顺序、并发控制和死信处理。
13.1 事件模型
public record OrderCreatedEvent(
String eventId,
Long orderId,
String userId,
Instant occurredAt) {
}
事件字段建议:
| 字段 | 用途 |
|---|---|
| eventId | 全局唯一,幂等 |
| eventType | 类型 |
| occurredAt | 业务发生时间 |
| traceId | 链路追踪 |
| version | 契约版本 |
| payload | 业务数据 |
事件名应表达事实:
OrderCreated
PaymentSucceeded
ShipmentDelivered
避免 OrderChanged 这种信息量不足的名称。
13.2 Spring 内部事件
发布:
@Service
public class OrderService {
private final ApplicationEventPublisher publisher;
@Transactional
public void create(Order order) {
repository.save(order);
publisher.publishEvent(new OrderCreatedEvent(...));
}
}
监听:
@Component
public class InventoryListener {
@TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT)
public void on(OrderCreatedEvent event) {
inventoryService.reserve(event.orderId());
}
}
区别:
| 注解 | 执行 |
|---|---|
@EventListener |
发布线程立即执行 |
@Async |
提交线程池异步执行 |
@TransactionalEventListener |
绑定事务阶段 |
默认 AFTER_COMMIT 只在同一线程的当前事务提交后执行。事务回滚时默认不调用。跨线程和异步事件要显式设计。
13.3 Kafka 生产者
配置:
spring:
kafka:
producer:
bootstrap-servers: localhost:9092
key-serializer: org.apache.kafka.common.serialization.StringSerializer
value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
acks: all
properties:
enable.idempotence: true
max.in.flight.requests.per.connection: 5
发送:
@Service
public class OrderEventProducer {
private final KafkaTemplate<String, OrderCreatedEvent> kafkaTemplate;
public CompletableFuture<SendResult<String, OrderCreatedEvent>> send(OrderCreatedEvent event) {
return kafkaTemplate.send("order-events", event.orderId().toString(), event);
}
}
顺序通常依赖:
- 相同 key 写同一分区;
- 生产端幂等和 max in flight 配置;
- 不无限重试乱序;
- 消费端按分区顺序处理。
13.4 Kafka 消费者
@Component
public class OrderEventListener {
@KafkaListener(
topics = "order-events",
groupId = "inventory-service",
concurrency = "3")
public void on(ConsumerRecord<String, OrderCreatedEvent> record,
Acknowledgment acknowledgment) {
handle(record.value());
acknowledgment.acknowledge();
}
}
提交策略:
| 模式 | 风险 |
|---|---|
| 自动提交 | 可能丢消息或重复消费 |
| 先提交后处理 | 可能丢消息 |
| 先处理后提交 | 至少一次,可能重复 |
推荐至少一次 + 幂等。幂等键用 eventId 或业务唯一键,而不是 offset。
13.5 重试与死信
处理失败:
业务消息
-> 本地重试
-> 延迟重试 topic
-> 死信 topic
-> 人工或自动补偿
异常分类:
| 异常 | 处理 |
|---|---|
| 参数缺失 | 记录死信,不重试 |
| 数据暂时不存在 | 延迟重试 |
| 下游超时 | 指数退避 |
| 序列化失败 | 死信并告警 |
| 未知异常 | 有限重试后死信 |
死信消息必须保留:
原消息
异常栈
topic / partition / offset
重试次数
traceId
时间
13.6 幂等消费
表结构:
create table consumed_events (
event_id varchar(64) primary key,
consumer_group varchar(64) not null,
processed_at timestamp not null
);
create index idx_consumed_group_time
on consumed_events(consumer_group, processed_at);
处理:
@Transactional
public void handle(OrderCreatedEvent event) {
if (eventStore.exists(event.eventId(), "inventory-service")) {
return;
}
inventoryService.reserve(event.orderId());
eventStore.save(event.eventId(), "inventory-service");
}
注意:
- 业务写和幂等记录要在一个本地事务;
- 幂等表需要清理策略;
- 并发同一 key 可能需要唯一约束兜底;
- 重试不会破坏结果;
- 消费位移提交要晚于业务事务。
13.7 Outbox 模式
create table outbox_events (
event_id varchar(64) primary key,
event_type varchar(100) not null,
aggregate_id varchar(64) not null,
payload json not null,
status varchar(20) not null,
retry_count int default 0,
created_at timestamp not null,
sent_at timestamp null
);
事务内:
@Transactional
public void create(Order order) {
orderRepository.save(order);
outboxRepository.save(OutboxEvent.from(order));
}
投递:
扫描未发送 outbox
-> 发送到 Kafka
-> 标记 sent
-> 失败递增 retry_count
相比事务内直接发消息,Outbox 保证业务成功与事件产生的原子性。CDC 方案可以进一步减少扫描压力。
13.8 Spring 调度
启用:
@EnableScheduling
@SpringBootApplication
public class Application {
}
任务:
@Component
public class SettlementJob {
@Scheduled(cron = "0 0 2 * * *", zone = "Asia/Shanghai")
public void settle() {
settlementService.settleYesterday();
}
}
固定延迟:
@Scheduled(fixedDelay = 5, timeUnit = TimeUnit.MINUTES)
public void reconcile() {
}
fixedRate 可能任务重叠,必须配置:
spring.task.scheduling.pool.size=4
13.9 分布式调度
多实例部署时,@Scheduled 默认每个实例都会执行。
方案:
| 方案 | 特点 |
|---|---|
| ShedLock | 数据库或 Redis 锁,简单 |
| XXL-Job / ElasticJob | 任务平台和分片 |
| Kubernetes CronJob | 平台调度 |
| Quartz 集群 | 传统数据库锁 |
ShedLock 示例:
@Scheduled(cron = "0 */5 * * * *")
@SchedulerLock(name = "settlement", lockAtMostFor = "10m", lockAtLeastFor = "1m")
public void settle() {
}
lockAtMostFor 必须大于最坏执行时间,否则锁提前释放导致重复执行。
13.10 任务可观测性
日志:
job=start jobId=...
job=end status=success cost=1200ms read=10000 written=9998
job=end status=failed retryable=true
指标:
job_start_total
job_success_total
job_failure_total
job_duration_seconds
job_last_success_timestamp
queue_lag
consumer_retry_total
dead_letter_total
告警:
- 任务超过窗口未成功;
- 消费 lag 持续增长;
- 死信增长;
- 消费失败率突增;
- 任务执行时间持续上升。
本章小结
消息系统应遵循至少一次投递、业务幂等、失败重试和死信治理。Spring 事件适合进程内解耦,跨服务事件使用 Kafka/RocketMQ 并结合 Outbox。调度任务在多实例环境必须加分布式锁或任务平台,并记录执行时长、结果和最后成功时间。
思考题
@EventListener和@TransactionalEventListener有什么区别?- 如何保证 Kafka 消费幂等?
- 为什么推荐 Outbox 模式?
fixedRate调度可能带来什么问题?- 多实例定时任务如何避免重复执行?