SpringNotes

第 13 章:消息与调度

zjc 于 2026-01-13 发布

这是《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);
    }
}

顺序通常依赖:

  1. 相同 key 写同一分区;
  2. 生产端幂等和 max in flight 配置;
  3. 不无限重试乱序;
  4. 消费端按分区顺序处理。

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");
}

注意:

  1. 业务写和幂等记录要在一个本地事务;
  2. 幂等表需要清理策略;
  3. 并发同一 key 可能需要唯一约束兜底;
  4. 重试不会破坏结果;
  5. 消费位移提交要晚于业务事务。

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

告警:

  1. 任务超过窗口未成功;
  2. 消费 lag 持续增长;
  3. 死信增长;
  4. 消费失败率突增;
  5. 任务执行时间持续上升。

本章小结

消息系统应遵循至少一次投递、业务幂等、失败重试和死信治理。Spring 事件适合进程内解耦,跨服务事件使用 Kafka/RocketMQ 并结合 Outbox。调度任务在多实例环境必须加分布式锁或任务平台,并记录执行时长、结果和最后成功时间。

思考题

  1. @EventListener@TransactionalEventListener 有什么区别?
  2. 如何保证 Kafka 消费幂等?
  3. 为什么推荐 Outbox 模式?
  4. fixedRate 调度可能带来什么问题?
  5. 多实例定时任务如何避免重复执行?