消息可靠发送案例

wen java案例 1

本文目录导读:

消息可靠发送案例

  1. 核心痛点与解决思路
  2. 分层级的可靠发送
  3. 领域驱动设计:基于聚合根的事务性发件箱(重点案例)
  4. 代码示例(基于Java + Spring Boot + RabbitMQ)
  5. 高级方案:基于消息中间件本身的机制
  6. 总结与最佳实践

我将为你详细讲解消息可靠发送的完整案例,涵盖从理论基础到代码实现的各个方面。

消息可靠发送是分布式系统和微服务架构中的核心难题,这里的“可靠”通常指不丢失不重复(或至少能处理重复)且顺序正确


核心痛点与解决思路

在分布式系统中,消息发送面临三大挑战:

  1. 网络不可靠:发送方发出消息,可能因网络超时而不知道是否送达。
  2. 进程故障:发送方或接收方在发送过程中可能宕机。
  3. 消息重复:由于重试机制,接收方可能收到重复消息。

核心解决思路“三管齐下”

  • 确认机制:接收方收到消息后必须发送确认(ACK)。
  • 重试机制:发送方未收到ACK时,自动重发。
  • 幂等性设计:接收方必须有能力处理重复消息,保证最终结果一致。

分层级的可靠发送

根据业务需求,消息发送的可靠性可以分为三个层级:

最多一次 (At-most-once)

  • 场景:允许丢失少量数据(如日志、监控)。
  • 实现:发送后不等待确认,不重试。
  • 风险:消息可能丢失。

最少一次 (At-least-once) —— 最常用

  • 场景:核心业务数据(如订单、支付)。
  • 实现:发送后等待ACK,超时未收到则重试。
  • 优势:保证不丢。
  • 风险:可能产生重复消息,需要接收方做幂等处理。

恰好一次 (Exactly-once)

  • 场景:金融交易、高价值业务。
  • 实现:依赖底层系统(如Kafka的幂等生产者、RocketMQ的事务消息),或通过“最少一次”+“下游幂等”来逻辑上实现

领域驱动设计:基于聚合根的事务性发件箱(重点案例)

这是目前业界解决“本地事务与发送消息不一致”问题的最佳实践,也是实现可靠发送的典型案例。

问题背景: 用户下单场景,你需要在数据库中创建订单,并向MQ发送一条“订单创建”消息,如果先入库,再发消息,入库成功但发消息失败,就会导致订单数据不一致。

解决方案清单:

  1. 创建消息表:在业务数据库中建一张 outbox(发件箱)表。
  2. 同事务写库:在创建订单的同一个事务中,将订单数据和待发送的消息标记一起写入数据库。
  3. 独立定时任务扫描:一个后台定时任务(或CDC工具,如Debezium)扫描 outbox 表,将状态为 PENDING 的消息发送给MQ。
  4. 确认并标记:收到MQ确认后,更新 outbox 表的状态为 SENT

流程图如下:

sequenceDiagram
    participant C as 业务服务
    participant DB as 业务数据库
    participant S as 定时任务/CDC
    participant MQ as 消息队列
    C->>DB: 1.开启事务
    C->>DB: 2.插入订单表
    C->>DB: 3.插入Outbox表(状态:待发送)
    C->>DB: 4.提交事务
    S->>DB: 5.查询待发送消息
    alt 消息存在
        S->>MQ: 6.发送消息
        MQ-->>S: 7.返回ACK
        S->>DB: 8.更新Outbox状态为已发送
    else 发送失败
        S->>S: 等待下次轮询重试
        Note over S: 由于消息在DB中,不会丢失
    end

代码示例(基于Java + Spring Boot + RabbitMQ)

这是一个简化的案例,演示了发件箱与本地事务的结合。

第一步:依赖配置

// pom.xml 关键依赖 (略)
// application.yml
server:
  port: 8081
spring:
  rabbitmq:
    host: localhost
    port: 5672
    username: guest
    password: guest
  datasource:
    url: jdbc:mysql://localhost:3306/test?useSSL=false
    username: root
    password: root

第二步:实现发件箱机制

实体类与Mapper(MyBatis-Plus)

import com.baomidou.mybatisplus.annotation.*;
import java.time.LocalDateTime;
// 数据库表对应实体
@Data
@TableName("outbox")
public class OutboxMessage {
    @TableId(value = "id", type = IdType.ASSIGN_ID)
    private Long id;
    private String aggregateType; // 聚合类型,"Order"
    private String aggregateId;   // 聚合ID,例如订单号
    private String payload;       // 要发送的消息体(JSON格式)
    @TableField(fill = FieldFill.INSERT)
    private LocalDateTime createTime;
    @TableField(fill = FieldFill.INSERT)
    private String status; // PENDING -> SENDING -> SENT
}

核心服务:订单创建(关键点)

import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import cn.hutool.json.JSONUtil;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import java.util.List;
import java.util.stream.Collectors;
@Service
public class OrderService {
    private final OrderMapper orderMapper;
    private final OutboxMapper outboxMapper;
    private final RabbitTemplate rabbitTemplate;
    // 构造函数注入...(省略)
    /**
     * 创建订单,并在同一个事务中写入发件箱
     * 最终一致性:保证订单数据与消息绝对一致
     */
    @Transactional(rollbackFor = Exception.class)
    public void createOrder(Order order) {
        // 1. 核心业务逻辑:保存订单
        orderMapper.insert(order);
        // 2. 专门发送消息:构建消息体
        OrderCreatedEvent event = new OrderCreatedEvent(
                order.getId(),
                order.getGoodsId(),
                order.getPrice()
        );
        // 3. 写入发件箱表(PENDING状态)
        OutboxMessage message = new OutboxMessage();
        message.setAggregateType("Order");
        message.setAggregateId(order.getId().toString());
        message.setPayload(JSONUtil.toJsonStr(event));
        message.setStatus("PENDING");
        // 注意:这两步在同一个数据库事务中
        // 如果插入Order失败,Outbox也不会写入
        outboxMapper.insert(message);
        // 事务提交后,任务将由定时任务执行发送
    }
}

发件箱轮询发送器(定时任务)

import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import java.util.List;
@Component
public class OutboxPoller {
    private final OutboxMapper outboxMapper;
    private final RabbitTemplate rabbitTemplate;
    // 注入...
    @Scheduled(fixedDelay = 1000) // 每秒扫描一次
    public void processMessages() {
        // 1. 查询待发送数据(为了防止多实例并发,通常加 limit 和 乐观锁/悲观锁)
        List<OutboxMessage> pendingMessages = outboxMapper.findByStatus("PENDING");
        for (OutboxMessage message : pendingMessages) {
            try {
                // 2. 发送前更新状态(防止重复消费)
                message.setStatus("SENDING");
                outboxMapper.updateById(message);
                // 3. 真正发送消息
                rabbitTemplate.convertAndSend(
                    "order.exchange", 
                    "order.created", 
                    message.getPayload()
                );
                // 4. 发送成功,更新状态
                message.setStatus("SENT");
                outboxMapper.updateById(message);
            } catch (Exception e) {
                // 记录日志并跳过(等待下次轮询重试)
                log.error("消息发送失败: {}", message.getId(), e);
                // 保留状态为SENDING,为了防止无限重试,可设定重试次数上限
            }
        }
    }
}

第三步:消息接收方(保证幂等性)

在消费者端,为了防止消息重复,必须在业务逻辑中进行幂等校验。

import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;
@Component
public class OrderConsumer {
    @Autowired
    private OrderStatusMapper orderStatusMapper;
    @Autowired
    private MessageLogMapper messageLogMapper;
    @RabbitListener(queues = "order.created.queue")
    public void handleOrderCreated(OrderCreatedEvent event) {
        Long eventId = event.getId();
        // 幂等性检查:确保当前消息没有被处理过
        // 方法1:使用Redis
        // if (stringRedisTemplate.hasKey("message:processed:" + eventId)) { return; }
        // 方法2:使用数据库消息表
        if (messageLogMapper.findByMsgId(eventId) != null) {
            log.info("消息已处理,跳过: {}", eventId);
            return;
        }
        // 核心业务逻辑处理
        // 这里执行更新订单状态等操作
        // 记录已处理的消息ID
        messageLogMapper.insert(new MessageLog(eventId));
    }
}

高级方案:基于消息中间件本身的机制

如果不想自己写发件箱逻辑,可以直接利用成熟中间件的高级特性:

RocketMQ 事务消息

  • 发送半消息。
  • 执行本地事务确认。
  • 回查与最终发送。
  • 实现:与发件箱方案类似,但已由中间件封装,不需要自己写轮询。

Kafka 幂等生产者

  • 设置 enable.idempotence=true,Kafka会在内部处理重试和去重,保证分区内不重复且顺序不丢失。

总结与最佳实践

场景 推荐方案
小规模单体应用 发件箱(Outbox)定时任务方案
中大型微服务 事务性发件箱 + CDC(如Debezium)
高并发海量数据 Kafka 事务 + 幂等生产者
金融分布式事务 阿里的 Seata(AT模式) 或 TCC 模式

关于可靠发送,有几个核心结论:

  1. 绝对不丢是核心目标,通过 数据落盘 + 重试 保证。
  2. 绝对不重是伪命题,网络重传无法避免,必须通过业务层幂等解决
  3. 推荐使用发件箱模式将数据一致性牢牢掌握在自己手中,这也是当前后端架构师岗位的核心考点。

希望这个案例能帮你全面理解消息可靠发送机制,如果有具体的中间件(如Kafka或RabbitMQ)细节需要深入,可以继续提问。

抱歉,评论功能暂时关闭!