人人都会AI编程

27.5 异步解耦与延迟消息业务场景

更新时间:2026-07-10

在微服务架构中,同步调用链路过长、核心流程与非关键流程耦合过紧,是导致系统响应慢、扩展难的主要原因之一。消息队列的引入,能够从架构层面实现 异步解耦,并借助 延迟消息 能力处理那些“稍后再执行”的业务场景。本节通过具体业务案例,展示如何利用消息队列将这两个能力落地。

27.5.1 异步解耦:把“非核心链路”剥离主流程

典型场景:用户注册送积分、发新人红包、写运营日志

在传统的同步实现中,用户注册接口往往会直接调用积分服务、优惠券服务、日志服务:

@Transactional
public void register(User user) {
    userDao.insert(user);           // 核心:保存用户
    pointService.addPoint(user);    // 非核心:送积分
    couponService.issueCoupon(user);// 非核心:发红包
    logService.record("REG",user);  // 非核心:运营日志
}

问题显而易见:

  • 响应时间变长:四个步骤串行执行,任何一个慢或被依赖服务抖动,都会拖死注册接口。
  • 可用性绑定:积分服务宕机?注册失败。优惠券服务超时?注册失败。核心流程的可用性被非核心功能绑架。
  • 紧耦合:每增加一个注册后处理(如发送欢迎短信),都必须修改注册代码,违反开闭原则。

引入消息队列后的异步解耦方案

只保留数据库写入这个核心动作,注册成功后发送一条 用户注册事件 到消息队列,其余业务由各自的消费者异步处理:

@Transactional
public void register(User user) {
    userDao.insert(user);
    // 发送注册事件,包含必要的用户信息
    RegisterEvent event = new RegisterEvent(user.getId(), user.getNickname());
    rabbitTemplate.convertAndSend("user.exchange", "user.register", event);
    // 主流程结束,立即返回成功
}

对应的消费者:

@RabbitListener(queues = "point.queue")
public void handleRegisterEventForPoint(RegisterEvent event) {
    pointService.addPoint(event.getUserId());
}

@RabbitListener(queues = "coupon.queue")
public void handleRegisterEventForCoupon(RegisterEvent event) {
    couponService.issueCoupon(event.getUserId());
}

改造后的效果:

  • 注册接口响应时间从几百毫秒下降到几十毫秒(仅剩数据库写入)。
  • 积分服务宕机不影响用户注册,消息会暂存队列,待恢复后继续消费。
  • 新增“欢迎短信”功能,只需新增一个消费者订阅同一事件,注册服务零改动。

实际落地要点

  • 消息可靠性保障:必须确保“用户插入成功”与“消息发送成功”的原子性。常用解决方案是 事务消息(如 RocketMQ) 或 发件箱模式(Outbox Pattern):先写入业务表,同时写入一张本地消息表,通过定时任务或 CDC 将消息投递到消息队列。
  • 消费者幂等:消息可能重复投递,积分发放等操作必须基于用户ID + 场景进行去重(如 Redis 记录已处理的 userId)。
  • 异步后的用户体验:注册后积分可能不会立即到账,产品层面需要告知用户“新人奖励将在5分钟内到账”。

27.5.2 延迟消息:处理定时执行类业务

典型场景:订单30分钟未支付自动取消、会议开始前15分钟提醒、延迟重试

这类场景的共同点是:某个操作需要在未来某个设定时刻触发,且触发时机精度要求不高(秒级或分钟级)。传统做法是启动定时任务轮询数据库:

@Scheduled(fixedDelay = 5000) // 每5秒扫一次
public void cancelUnpaidOrders() {
    List<Order> orders = orderDao.findUnpaidOlderThan(30, TimeUnit.MINUTES);
    orders.forEach(this::cancelOrder);
}

轮询的代价随着数据量增长而线性上升,大量无效扫描浪费资源。使用 延迟消息 可以将“定时检测”转变为“准时通知”:

方案一:RabbitMQ 死信队列实现延迟消息

RabbitMQ 自身不直接支持任意延迟,但可以利用消息的 TTL(生存时间)和死信队列(DLX)曲线救国:

  1. 创建一个带有 x-message-ttl 的队列(延迟队列),并指定死信交换机。
  2. 生产者将消息发送到该队列,不消费。
  3. 消息在队列中存活 TTL 时长后过期,自动转发到死信队列。
  4. 业务消费者监听死信队列,收到消息时执行取消订单等操作。

消息属性设置(发送时):

MessagePostProcessor postProcessor = message -> {
    message.getMessageProperties().setExpiration(String.valueOf(30 * 60 * 1000)); // 30分钟
    return message;
};
rabbitTemplate.convertAndSend("order.delay.exchange", "order.delay", orderId, postProcessor);

方案二:RocketMQ 延迟消息

RocketMQ 天然支持18个级别的延迟消息(从1s到2h),生产者只需指定延迟级别:

Message message = new Message("ORDER_TOPIC", orderId.getBytes());
// 设置延迟级别:3 对应 10s,4 对应 30s,5 对应 1m,以此类推
message.setDelayTimeLevel(4); 
rocketMQTemplate.syncSend("ORDER_TOPIC", message);

消费者对接收到的时间消息进行处理,无需额外队列配置。

方案三:定时任务 + 消息队列(折中)

如果不想引入复杂的延迟消息机制,也可以使用定时任务 + 消息队列实现离散化扫描:定时任务只负责扫出即将到期的订单,然后每条订单发送一条消息到队列,由消费者并行处理。这种方式将“集中处理”转为“分散消费”,提高处理效率和容错性。

27.5.3 延迟消息实战:订单30分钟未支付自动取消

以 RabbitMQ 死信队列为例,给出关键配置和代码。

1. 声明交换机和队列

@Configuration
public class DelayQueueConfig {
    @Bean
    public DirectExchange orderDelayExchange() {
        return new DirectExchange("order.delay.exchange");
    }

    @Bean
    public Queue orderDelayQueue() {
        return QueueBuilder.durable("order.delay.queue")
                .withArgument("x-dead-letter-exchange", "order.dlx.exchange")
                .withArgument("x-dead-letter-routing-key", "order.dlx")
                .withArgument("x-message-ttl", 30 * 60 * 1000) // 30分钟
                .build();
    }

    @Bean
    public Binding delayBinding() {
        return BindingBuilder.bind(orderDelayQueue()).to(orderDelayExchange()).with("order.delay");
    }

    // 死信交换机及队列
    @Bean
    public DirectExchange orderDlxExchange() {
        return new DirectExchange("order.dlx.exchange");
    }

    @Bean
    public Queue orderDlxQueue() {
        return QueueBuilder.durable("order.dlx.queue").build();
    }

    @Bean
    public Binding dlxBinding() {
        return BindingBuilder.bind(orderDlxQueue()).to(orderDlxExchange()).with("order.dlx");
    }
}

2. 订单创建后发送延迟消息

@Service
public class OrderService {
    @Autowired
    private RabbitTemplate rabbitTemplate;

    @Transactional
    public void createOrder(Order order) {
        orderDao.insert(order);
        // 发送延迟消息,30分钟后触发
        rabbitTemplate.convertAndSend("order.delay.exchange", "order.delay", order.getId(), msg -> {
            msg.getMessageProperties().setExpiration(String.valueOf(30 * 60 * 1000));
            return msg;
        });
    }
}

3. 消费者处理过期订单

@Component
public class OrderCancelConsumer {
    @Autowired
    private OrderDao orderDao;

    @RabbitListener(queues = "order.dlx.queue")
    public void cancelOrder(Long orderId) {
        Order order = orderDao.findById(orderId);
        if (order != null && order.getStatus() == OrderStatus.UNPAID) {
            order.setStatus(OrderStatus.CANCELED);
            orderDao.update(order);
            log.info("订单 {} 超时未支付,已自动取消", orderId);
        }
    }
}

4. 注意事项

  • 幂等检查:消息可能重复,必须检查订单状态,避免重复取消或关闭已支付订单。
  • 延迟时间限制:RabbitMQ 死信方案中,消息一旦入队,TTL 就无法修改。若订单支付后需要取消延迟消息,可以记录消息ID,消费者判断状态后直接忽略。
  • 生产级替代:若延迟场景复杂(动态修改到期时间、大量延迟消息),建议使用专门任务调度系统(如 SchedulerX、XXL-JOB 的延迟任务),或直接采用 RocketMQ 的延迟消息。

27.5.4 综合场景:下单未支付倒计时通知

实际业务中往往不是单一操作,而是 阶梯式触发:下单后15分钟未支付发短信提醒,30分钟未支付自动取消。可以通过发送两条不同延迟时间的消息实现:

  • 订单创建时发送两条消息:一条 TTL=15分钟,路由到“提醒”消费者;另一条 TTL=30分钟,路由到“取消”消费者。
  • 消费者同样需要检查订单当前状态,以决定是否继续执行。

这样就将一个原本需要轮询+定时任务才能实现的复杂流程,完全转化为基于消息的异步驱动链路,极大降低了定时调度的压力,提升了系统响应能力和可扩展性。