跳转至

一、消息队列与数据库协同

1.1 事务消息与本地消息表

RocketMQ 事务消息两阶段提交

📌 1. 什么是 RocketMQ 事务消息?它解决了什么问题?

事务消息是 RocketMQ 提供的一种分布式事务解决方案。它让“发送消息”这个动作与本地数据库操作在逻辑上绑定为一个整体——要么消息发成功且数据库操作提交,要么数据库回滚、消息也不可见。

传统消息队列困境:
  1. 先操作数据库,再发送消息
      → 数据库提交后,发送消息失败 → 数据不一致
  2. 先发送消息,再操作数据库
      → 消息发出去了,但数据库操作失败 → 下游消费了不该存在的数据

事务消息:
  半消息先发送 → 执行本地事务 → 根据本地事务结果决定消息是提交还是回滚

它解决的核心问题是分布式系统下,本地事务和消息发送的原子性。比如下单后需要通知库存系统扣减库存,必须保证“订单创建成功”与“库存扣减消息被可靠投递”这两个动作要么都成功,要么都失败,而不会出现一个成功另一个失败的不一致状态。


⚙️ 2. 详细说明 RocketMQ 事务消息的两阶段提交流程,以及如何保证消息不丢失?

RocketMQ 事务消息采用 两阶段提交 + 事务回查 的机制。整个流程包含:发送半消息 → 执行本地事务 → 提交/回滚消息,以及 Broker 的 回查机制。

流程图示:

 生产者 (Producer)                   Broker (MQ服务器)              消费者 (Consumer)
     │                                     │                              │
     │ ① 发送半消息 (Half Message)          │                              │
     ├────────────────────────────────────→│                              │
     │     (此时消息对消费者不可见)          │ 保存半消息,标记为“暂不可投递”  │
     │                                     │                              │
     │ ② 半消息发送成功,返回半消息 offset    │                              │
     │←────────────────────────────────────┤                              │
     │                                     │                              │
     │ ③ 执行本地事务 (例如:创建订单)        │                              │
     │    ┌─ 本地事务成功 ──────────────┐    │                              │
     │    │ 提交消息 (Commit)            │    │                              │
     │    ├──────────────────────────→  │ 标记半消息为可投递 → 消费者可见   │
     │    └───────────────────────────┘    │                              │
     │    ┌─ 本地事务失败 ──────────────┐    │                              │
     │    │ 回滚消息 (Rollback)          │    │                              │
     │    ├──────────────────────────→  │ 删除半消息                     │
     │    └───────────────────────────┘    │                              │
     │                                     │                              │
     │ ④ 如果 Producer 未返回 Commit/Rollback (断网/宕机)                   │
     │    Broker 会定时回查 (Checkback)      │                              │
     │←────────────────────────────────────┤                              │
     │ ⑤ Producer 收到回查,检查本地事务状态   │                              │
     │    → 提交/回滚消息                    │                              │
     └─────────────────────────────────────┘                              │

步骤详解:

  1. 半消息发送:生产者向 Broker 发送一条“半消息”。Broker 持久化该消息,但标记为 暂不可投递,消费者无法拉取。此步骤保证了消息不会因为网络闪断而丢失,因为 Broker 已经落盘。

  2. 执行本地事务:生产者收到半消息的成功响应后,执行本地数据库操作(如创建订单)。本地事务的结果会决定消息的最终命运。

  3. 如果本地事务成功,生产者向 Broker 发送 Commit 请求,Broker 将半消息标记为可投递。
  4. 如果本地事务失败,生产者发送 Rollback 请求,Broker 删除该半消息。

  5. Broker 回查:如果生产者因为网络超时、宕机等原因没有及时返回 CommitRollback,Broker 会定期扫描长时间未确认的半消息,主动发起事务回查。生产者需要在回查接口中根据本地事务的最终状态,回复 CommitRollback。这个机制是保证“不丢失”的最终防线。

如何保证消息不丢失?

RocketMQ 通过层层把关来保障消息的可靠性:

  • 同步刷盘:Broker 将半消息写入磁盘后再返回成功,保证消息在 Broker 侧不丢失。

  • 主从复制:半消息会被同步到从节点,主节点宕机后从节点可接管,消息依然存在。

  • 事务回查:对于未确认的半消息,Broker 会以指数退避的策略反复回查,直到获得明确的提交或回滚结果。回查次数和间隔均可配置,保证即使在极端宕机场景下,消息状态最终也能确定。

  • 消费端 ACK:消费者只有成功处理消息后才会返回确认,否则 RocketMQ 会重试投递,确保消息被成功消费。

示例代码:实现一个带事务消息的生产者

// 下单服务:事务消息发送
@Service
public class OrderService {

    @Autowired
    private RocketMQTemplate rocketMQTemplate;

    @Transactional
    public void createOrder(Order order) {
        // 构造消息
        String msg = JSON.toJSONString(order);
        Message<String> message = MessageBuilder.withPayload(msg).build();

        // 发送事务消息
        TransactionSendResult result = rocketMQTemplate.sendMessageInTransaction(
            "order_topic:inventory", message, order);

        if (result.getLocalTransactionState() == LocalTransactionState.COMMIT_MESSAGE) {
            log.info("订单创建成功,事务消息已提交");
        } else {
            throw new RuntimeException("订单创建失败,事务消息已回滚");
        }
    }
}

// 事务监听器:执行本地事务 + 回查
@RocketMQTransactionListener
public class OrderTransactionListener implements RocketMQLocalTransactionListener {

    @Autowired
    private OrderDao orderDao;

    @Override
    public RocketMQLocalTransactionState executeLocalTransaction(Message msg, Object arg) {
        try {
            Order order = (Order) arg;
            orderDao.insert(order);          // 本地事务:插入订单
            return RocketMQLocalTransactionState.COMMIT;
        } catch (Exception e) {
            return RocketMQLocalTransactionState.ROLLBACK;
        }
    }

    @Override
    public RocketMQLocalTransactionState checkLocalTransaction(Message msg) {
        // Broker 回查时,根据订单是否存在来判断事务是否成功
        Order order = JSON.parseObject(new String((byte[]) msg.getPayload()), Order.class);
        if (orderDao.findById(order.getId()) != null) {
            return RocketMQLocalTransactionState.COMMIT;
        } else {
            return RocketMQLocalTransactionState.ROLLBACK;
        }
    }
}

🛒 3. 在电商订单系统中,用户下单成功后需要异步通知库存系统扣减库存,如何保证下单和通知库存的一致性?

这是一个典型的分布式最终一致性问题。我们需要保证:订单创建成功 → 库存一定被扣减;订单创建失败 → 库存一定不扣减。

解决方案:使用 RocketMQ 事务消息 + 库存系统幂等消费

用户下单
┌─────────────────┐
│  订单服务        │  ① 发送半消息 (OrderCreated)  →  Broker
│  执行本地事务     │  ② 插入订单数据 (本地数据库)
│  提交事务消息     │  ③ 本地事务成功 → 提交半消息
└─────────────────┘
                        ▼ ④ 消息对消费者可见
                ┌──────────────┐
                │  库存服务     │
                │  消费消息     │
                │  扣减库存     │  (幂等处理:根据 orderId 去重)
                └──────────────┘

流程详解:

  1. 订单服务发送一条半消息,内容为“订单已创建,请扣减库存”,携带订单 ID。

  2. 半消息发送成功后,订单服务执行本地事务——向订单表插入一条记录。此时订单状态可以是“待支付”或“已创建”。

  3. 本地事务成功 → 订单服务向 Broker 提交事务消息,Broker 标记消息为可投递。若本地事务失败(如数据库主键冲突),订单服务回滚消息,Broker 删除该半消息。

  4. 库存服务订阅该消息。当收到消息后,执行库存扣减操作。为了应对重复消费(网络重试或回查导致),库存服务必须基于 orderId 做幂等处理——比如在库存表中记录一条扣减流水,以订单 ID 为唯一索引,重复消费时忽略。

为什么必须用事务消息?

如果订单服务和库存服务之间只用普通消息,会出现:

  • 先插订单后发消息:如果发消息失败,订单创建了但库存没扣。

  • 先发消息后插订单:如果插订单失败,消息已发出,库存白白扣减。

事务消息让“插订单”与“发消息”在同一个本地事务边界内,通过半消息的回查机制,保证了这两个操作的原子性。即便订单服务在发送半消息后宕机,Broker 也会通过回查来最终确认订单是否真实存在,从而决定是否投递扣库存的消息。

库存服务的幂等消费代码示例:

@Service
@RocketMQMessageListener(topic = "order_topic", consumerGroup = "inventory_group")
public class InventoryConsumer implements RocketMQListener<String> {

    @Autowired
    private InventoryService inventoryService;

    @Override
    public void onMessage(String message) {
        OrderEvent event = JSON.parseObject(message, OrderEvent.class);
        String orderId = event.getOrderId();

        // 幂等处理:根据订单ID查询是否已经扣减过库存
        if (inventoryService.isAlreadyDeducted(orderId)) {
            log.info("库存已扣减,忽略重复消息");
            return;
        }

        // 执行扣减库存,并记录扣减流水(在同一本地事务中)
        inventoryService.deductAndRecord(orderId, event.getSkuItems());
    }
}

异常场景覆盖:

  • 订单插入成功,提交消息失败:Broker 回查发现订单已存在 → 提交消息,库存最终扣减。

  • 订单插入失败,但半消息已发:生产者回滚消息,Broker 删除;即使生产者宕机,回查发现订单不存在 → 回滚消息,库存不会扣减。

  • 库存服务消费失败:RocketMQ 会重试消费,直到成功(保证幂等)。

  • 网络超时导致半消息未确认:Broker 按固定间隔(如 6 秒开始,指数退避到 60 秒)持续回查,直到获得确认。这个机制保证了临时故障不会导致消息丢失。

通过这种设计,订单服务和库存服务在分布式环境下实现了最终一致性——要么订单与库存扣减同时成功,要么二者都不发生,系统绝不会出现“订单创建了但库存没扣”或“库存扣了但订单没创建”的不一致状态。

本地消息表轮询重试


📌 1. 什么是本地消息表?它和 RocketMQ 事务消息有什么区别?

本地消息表 是一种分布式事务的最终一致性方案。它的核心思想很简单:把要发送的消息先持久化到业务数据库的一张“消息表”里,和业务数据在同一个本地事务中落库。然后再用一个后台任务轮询这张表,把未发送的消息投递到消息队列。 这样,消息就“沾”上了业务数据的 ACID 特性,不会凭空消失。

image.png

RocketMQ 事务消息 则是 RocketMQ 自带的分布式事务能力,它通过“半消息 + 本地事务 + 回查”的机制,把本地事务和消息发送绑定在一起,无需开发者自己维护消息表。

两者的区别一目了然:

查看内嵌表格

一句话总结: 本地消息表是自食其力的“手工作坊”,RocketMQ 事务消息是现代化的“自动化流水线”。如果你已经用了 RocketMQ,首选事务消息;如果 MQ 不支持事务消息,或者你想完全掌控流程、不想依赖特定中间件,本地消息表依然是一个可靠的后备方案。


🛠️ 2. 请设计一个基于本地消息表的订单通知方案,并说明如何保证消息不丢失?

下面以“订单创建后通知库存系统”为例,设计一个基于本地消息表的可靠通知方案。

架构流程:

image.png

如何保证消息不丢失?

消息不丢失需要贯穿“生产 → 传输 → 消费”全链路。

① 生产端:消息与业务数据同库、同事务

-- 伪代码:订单创建时的本地事务
BEGIN TRANSACTION
    INSERT INTO orders (id, user_id, amount, status) VALUES (...);
    INSERT INTO outbox (id, order_id, topic, message_body, status, create_time)
         VALUES (uuid, order_id, 'order_created', '{"orderId":123}', 'PENDING', now());
COMMIT

因为两行 INSERT 在同一个数据库事务里,要么都成功,要么都失败。订单创建成功,则消息必然写入磁盘;订单创建失败,消息也不会残留。这是防止消息丢失的第一道防线。

② 发送端:可靠的轮询投递 定时任务每隔几秒扫描 status = 'PENDING' 的记录。为防止并发和超时,需要在消息表中加锁或使用乐观锁:

// 使用乐观锁防止重复发送
int updated = db.update(
    "UPDATE outbox SET status='SENDING', version=version+1 WHERE id=? AND status='PENDING' AND version=?",
    recordId, oldVersion
);
if (updated > 0) {
    // 成功抢占到该消息的发送权
    sendToMQ(record);
    db.update("UPDATE outbox SET status='SENT' WHERE id=?", recordId);
}

如果发送成功,标记 SENT。如果发送失败(如网络超时),可以标记 PENDING 等待下次重试,或者设置重试次数上限,超过后标记 FAILED 并告警。

③ 传输端:MQ 自身的可靠性

消息队列(如 RocketMQ)自带同步刷盘、主从复制、消息持久化等机制。只要消息成功投递到 MQ,传输过程中的可靠性由 MQ 保证。

④ 消费端:幂等消费 + 手动 ACK + 死信队列

  • 幂等:消费者根据 order_id 判断库存是否已扣减,避免重复消费。

  • 手动 ACK:消费成功后才确认消息,如果业务处理失败,不确认,MQ 会重试。

  • 死信队列:重试多次仍失败的消息转入死信队列,由人工或监控系统处理,确保消息不丢失。

⑤ 兜底:定时任务对账

除了实时通知,还可以在每天低峰期进行一次全量对账。例如对比订单表状态与库存扣减流水,发现不一致的订单自动补发通知或触发人工修复。

示例代码:轮询发送任务 (Spring Scheduler)

@Scheduled(fixedDelay = 5000)
public void sendPendingMessages() {
    List<OutboxRecord> records = outboxMapper.selectPending(100); // 每次取100条
    for (OutboxRecord record : records) {
        // 乐观锁更新状态
        boolean locked = outboxMapper.tryLock(record.getId(), record.getVersion());
        if (!locked) continue;

        try {
            Message<String> msg = MessageBuilder.withPayload(record.getMessageBody()).build();
            rocketMQTemplate.send(record.getTopic(), msg);
            outboxMapper.updateStatus(record.getId(), "SENT");
        } catch (Exception e) {
            log.error("发送消息失败: {}", record.getId(), e);
            outboxMapper.updateStatus(record.getId(), "PENDING"); // 回退状态,等待重试
        }
    }
}

这个设计的好处是:无论何种异常(数据库宕机、MQ 宕机、网络分区),最终都会通过轮询和重试把消息送到 MQ,不会丢失一条已创建订单的通知。


🔗 3. 在支付系统中,支付成功后需要通知多个下游系统(如积分、优惠券、风控),如何保证所有通知都成功?

这是一个典型的“一对多可靠通知”问题。核心挑战:一个支付成功事件,必须可靠地送达积分、优惠券、风控三个系统,且不能丢失、不能重复、不能部分成功。

推荐方案:本地消息表 + 订阅者模式 + 幂等消费

支付服务 (支付成功后)
    └─ 本地事务: 更新支付单 + 插入一条 outbox 记录
                  ┌───────────────┐
                  │  本地消息表    │
                  │  topic: payment_paid
                  │  body: {"orderId":123,"amount":100}
                  └───────┬───────┘
                          │ 定时任务投递
                  ┌───────────────┐
                  │   消息队列     │
                  └───────┬───────┘
           ┌──────────────┼──────────────┐
           ▼              ▼              ▼
      积分服务        优惠券服务       风控服务
   (消费 + 幂等)    (消费 + 幂等)    (消费 + 幂等)

关键设计点:

① 主服务只发一条消息,多系统订阅

支付服务不关心下游有几个系统,它只负责发出一个“支付成功”的领域事件。积分、优惠券、风控各自订阅同一个 Topic,各自消费。这样支付服务完全解耦,未来新增下游(如消息推送)只需要增加一个新的消费者,不影响支付逻辑。

② 每个消费者独立 ACK 和重试

消息队列支持多个消费者组。每个下游系统属于不同的消费者组,各自独立地消费同一条消息。积分服务消费失败不会影响优惠券服务的消费。消息队列会为每个消费者组分别重试失败的消息。

③ 每个消费者必须实现幂等

以积分服务为例:

@RocketMQMessageListener(topic = "payment_paid", consumerGroup = "points_group")
public class PointsConsumer implements RocketMQListener<String> {
    @Autowired
    private PointsService pointsService;

    @Override
    public void onMessage(String message) {
        PaymentEvent event = JSON.parseObject(message, PaymentEvent.class);
        String orderId = event.getOrderId();
        // 幂等判断
        if (pointsService.isAlreadyAwarded(orderId)) {
            log.info("积分已发放,忽略重复消息");
            return;
        }
        // 发放积分 + 记录发放流水 (同一本地事务)
        pointsService.awardAndRecord(orderId, event.getAmount());
    }
}

每个下游系统都基于 orderId(或支付流水号)做幂等校验,保证多次重试不会重复发积分/扣券。

④ 本地消息表的轮询任务保证“至少一次投递”

支付服务通过本地消息表保证消息一定到达 MQ。MQ 通过消费者组重试保证消息一定被每个下游系统消费成功。这样即使积分服务挂了半小时,重启后仍然能继续消费积压的消息。

⑤ 死信队列与人工兜底

对于多次重试(如 16 次)仍然消费失败的消息,RocketMQ 会将其转入死信队列。运维人员可以监控死信队列,分析失败原因后,修正数据并重新投递,或手动补偿。

⑥ 最终一致性核对

可以建立一个每日对账任务,分别查询支付单状态、积分流水、优惠券使用记录、风控记录,对比是否存在“支付成功但某系统未处理”的差异,生成日报并告警。

总结:

通过 本地消息表保证事件不丢 + MQ 多消费者组独立投递 + 各系统幂等消费 + 死信队列兜底,我们就能在一个不太可靠的分布式环境里,实现“支付成功事件最终通知到所有下游”的目标。这套架构不追求强一致性,但追求最终一致、绝对可靠、异常可追溯,是电商支付场景的经典实践。

消息幂等性

1、什么是消息幂等性?为什么需要消息幂等性?

(幂等性定义、产生重复消息的原因)

消息幂等性是指:无论消息被消费多少次,产生的结果都是一样的。即多次消费和一次消费的效果相同。

需要消息幂等性的原因:

  1. 生产者重复发送:网络抖动导致生产者没有收到 Broker 的 ACK,重试发送

  2. 消费者重复消费:消费者消费成功后 ACK 丢失,Broker 重新投递

  3. 事务消息回查:RocketMQ 事务消息的回查机制可能导致消息重复投递

  4. 本地消息表重试:定时任务轮询重试可能导致消息重复发送

2、请设计几种消息幂等性方案,并对比它们的优缺点?

⭐⭐(幂等性方案设计、方案对比)

1️⃣ Common Answer 幂等性就是用数据库唯一索引,或者用 Redis 存一个 key 表示已经处理过。消费的时候先查一下,如果已经处理过就跳过。

2️⃣ Impressive Answer 消息幂等性是分布式系统中的核心问题,我会从业务幂等性技术幂等性两个维度来设计。给你详细讲一下:

方案一:数据库唯一索引(推荐)

这是最可靠、最简单的方案,适用于所有需要持久化的业务场景。

-- 业务表 + 唯一索引
CREATE TABLE stock_deduct (
    id BIGINT AUTO_INCREMENT PRIMARY KEY,
    order_id BIGINT NOT NULL,
    sku_id BIGINT NOT NULL,
    quantity INT NOT NULL,
    create_time DATETIME,
    UNIQUE KEY uk_order_sku (order_id, sku_id) -- 幂等性关键
);
@RocketMQMessageListener(topic = "order_stock", consumerGroup = "stock_group")
public class StockConsumer implements RocketMQListener<OrderMessage> {

    @Override
    public void onMessage(OrderMessage message) {
        try {
            // 直接插入,如果已存在会抛出 DuplicateKeyException
            StockDeduct record = new StockDeduct();
            record.setOrderId(message.getOrderId());
            record.setSkuId(message.getSkuId());
            record.setQuantity(message.getQuantity());
            stockMapper.insert(record);

            // 执行库存扣减
            stockService.deductStock(message.getSkuId(), message.getQuantity());

        } catch (DuplicateKeyException e) {
            log.warn("消息已处理过,跳过重复消费,orderId={}", message.getOrderId());
        }
    }
}

优点

  • 可靠性高,不依赖外部系统

  • 实现简单,无需额外代码

  • 天然支持事务一致性

缺点

  • 需要修改业务表结构

  • 不适用于非持久化场景


方案二:Redis 去重表

适用于高并发、对性能要求高的场景。

@RocketMQMessageListener(topic = "order_stock", consumerGroup = "stock_group")
public class StockConsumer implements RocketMQListener<OrderMessage> {

    @Autowired
    private RedisTemplate<String, String> redisTemplate;

    @Override
    public void onMessage(OrderMessage message) {
        String key = "msg:dedup:" + message.getOrderId() + ":" + message.getSkuId();

        // SETNX:如果 key 不存在则设置成功,返回 true;否则返回 false
        Boolean success = redisTemplate.opsForValue().setIfAbsent(key, "1", 24, TimeUnit.HOURS);

        if (Boolean.TRUE.equals(success)) {
            // 首次消费,执行业务逻辑
            stockService.deductStock(message.getSkuId(), message.getQuantity());
            log.info("消息处理成功,orderId={}", message.getOrderId());
        } else {
            // 重复消费,直接跳过
            log.warn("消息已处理过,跳过重复消费,orderId={}", message.getOrderId());
        }
    }
}

优点

  • 性能高,Redis 操作快

  • 不需要修改业务表结构

缺点

  • 依赖 Redis 的可用性

  • 需要设置合理的过期时间

  • Redis 宕机可能导致重复消费


方案三:状态机 + 数据库乐观锁

适用于业务状态流转的场景。

-- 订单表(带版本号)
CREATE TABLE orders (
    id BIGINT PRIMARY KEY,
    status VARCHAR(20) NOT NULL,
    version INT NOT NULL DEFAULT 0, -- 乐观锁版本号
    update_time DATETIME
);
@RocketMQMessageListener(topic = "order_update", consumerGroup = "order_group")
public class OrderConsumer implements RocketMQListener<OrderMessage> {

    @Override
    public void onMessage(OrderMessage message) {
        // 使用乐观锁更新状态
        int updated = orderMapper.updateStatusWithVersion(
            message.getOrderId(),
            "PAID",
            "UNPAID",  // 只有从 UNPAID 状态才能更新为 PAID
            message.getVersion()
        );

        if (updated > 0) {
            // 更新成功,执行后续逻辑
            log.info("订单状态更新成功,orderId={}", message.getOrderId());
        } else {
            // 更新失败,说明状态已变更或重复消费
            log.warn("订单状态更新失败或重复消费,orderId={}", message.getOrderId());
        }
    }
}
-- Mapper SQL
UPDATE orders
SET status = #{newStatus},
    version = version + 1,
    update_time = NOW()
WHERE id = #{orderId}
  AND status = #{oldStatus}
  AND version = #{version};

优点

  • 天然支持业务状态流转

  • 避免并发问题

缺点

  • 需要业务表有版本号字段

  • 实现相对复杂


方案四:分布式锁

适用于需要保证同一时间只有一个消费者处理消息的场景。

@RocketMQMessageListener(topic = "order_stock", consumerGroup = "stock_group")
public class StockConsumer implements RocketMQListener<OrderMessage> {

    @Autowired
    private RedissonClient redissonClient;

    @Override
    public void onMessage(OrderMessage message) {
        String lockKey = "lock:order:" + message.getOrderId();
        RLock lock = redissonClient.getLock(lockKey);

        try {
            // 尝试获取锁,最多等待 0 秒,锁 10 秒后自动释放
            boolean acquired = lock.tryLock(0, 10, TimeUnit.SECONDS);

            if (acquired) {
                try {
                    // 获取锁成功,执行业务逻辑
                    // 先查一下是否已经处理过
                    boolean processed = checkIfProcessed(message.getOrderId());
                    if (!processed) {
                        stockService.deductStock(message.getSkuId(), message.getQuantity());
                        markAsProcessed(message.getOrderId());
                    }
                } finally {
                    lock.unlock();
                }
            } else {
                // 获取锁失败,说明其他消费者正在处理
                log.warn("获取锁失败,跳过重复消费,orderId={}", message.getOrderId());
            }
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            log.error("获取锁被中断", e);
        }
    }
}

优点

  • 避免并发重复消费

  • 可以控制并发度

缺点

  • 实现复杂,依赖 Redis

  • 性能相对较低


方案对比

查看内嵌表格

生产环境推荐组合

  1. 核心业务:数据库唯一索引(最可靠)

  2. 高并发场景:Redis 去重表 + 数据库唯一索引(双重保障)

  3. 状态流转:状态机 + 乐观锁

  4. 特殊场景:分布式锁 + 数据库唯一索引

3️⃣ Key Differences

查看内嵌表格

3、在积分系统中,用户支付成功后需要增加积分,如何保证积分不会重复增加?

⭐⭐⭐(幂等性实战、业务幂等性设计)

1️⃣ Common Answer 用 Redis 存一个 key,消费的时候先查一下,如果已经处理过就跳过。或者用数据库唯一索引,插入一条积分记录,如果重复就报错。

2️⃣ Impressive Answer 积分系统的幂等性设计需要考虑业务幂等性技术幂等性两个层面。我会用数据库唯一索引 + 业务状态校验的双重保障方案。让我详细设计一下:

业务场景分析

  • 用户支付成功后,需要给用户增加积分

  • 积分增加后,用户可以用积分兑换商品

  • 如果积分重复增加,会导致用户积分异常,造成资损

数据库表设计

-- 用户积分表
CREATE TABLE user_point (
    id BIGINT AUTO_INCREMENT PRIMARY KEY,
    user_id BIGINT NOT NULL,
    total_point INT NOT NULL DEFAULT 0,
    update_time DATETIME,
    UNIQUE KEY uk_user_id (user_id)
);

-- 积分明细表(幂等性表)
CREATE TABLE point_detail (
    id BIGINT AUTO_INCREMENT PRIMARY KEY,
    user_id BIGINT NOT NULL,
    order_id BIGINT NOT NULL,          -- 关联订单
    point INT NOT NULL,                -- 积分变化量
    type VARCHAR(20) NOT NULL,         -- 类型:PAYMENT_ADD, EXCHANGE_USE
    status TINYINT NOT NULL DEFAULT 0, -- 0:待生效 1:已生效
    create_time DATETIME,
    UNIQUE KEY uk_order (order_id),    -- 幂等性关键:同一订单只能有一条记录
    INDEX idx_user_time (user_id, create_time)
);

核心代码实现

积分消费者(双重幂等性保障):

@RocketMQMessageListener(topic = "payment_success", consumerGroup = "point_group")
public class PointConsumer implements RocketMQListener<PaymentMessage> {

    @Autowired
    private UserPointMapper userPointMapper;

    @Autowired
    private PointDetailMapper pointDetailMapper;

    @Override
    public void onMessage(PaymentMessage message) {
        Long userId = message.getUserId();
        Long orderId = message.getOrderId();
        Integer point = calculatePoint(message.getAmount()); // 根据金额计算积分

        try {
            // 第一重幂等性:插入积分明细(数据库唯一索引)
            PointDetail detail = new PointDetail();
            detail.setUserId(userId);
            detail.setOrderId(orderId);
            detail.setPoint(point);
            detail.setType("PAYMENT_ADD");
            detail.setStatus(0);
            pointDetailMapper.insert(detail);

            // 第二重幂等性:校验订单是否已处理(防止并发重复消费)
            PointDetail existingDetail = pointDetailMapper.selectByOrderId(orderId);
            if (existingDetail != null && existingDetail.getStatus() == 1) {
                log.warn("订单已处理过,跳过重复消费,orderId={}", orderId);
                return;
            }

            // 增加用户积分
            int updated = userPointMapper.addPoint(userId, point);
            if (updated > 0) {
                // 更新积分明细状态为已生效
                pointDetailMapper.updateStatus(detail.getId(), 1);
                log.info("积分增加成功,userId={}, orderId={}, point={}", userId, orderId, point);
            } else {
                throw new RuntimeException("增加积分失败");
            }

        } catch (DuplicateKeyException e) {
            // 唯一索引冲突,说明订单已处理过
            log.warn("订单已处理过(唯一索引冲突),跳过重复消费,orderId={}", orderId);

            // 二次确认:查询订单状态,确保积分已生效
            PointDetail detail = pointDetailMapper.selectByOrderId(orderId);
            if (detail != null && detail.getStatus() == 1) {
                log.info("二次确认:订单积分已生效,orderId={}", orderId);
            } else {
                // 异常情况:订单记录存在但状态未生效,需要人工介入
                log.error("异常:订单记录存在但状态未生效,orderId={}", orderId);
                alertService.sendAlert("积分状态异常,orderId=" + orderId);
            }

        } catch (Exception e) {
            log.error("积分处理失败,orderId={}", orderId, e);
            // 抛出异常,触发 RocketMQ 重试
            throw e;
        }
    }

    private Integer calculatePoint(BigDecimal amount) {
        // 积分规则:每消费 1 元获得 1 积分
        return amount.intValue();
    }
}

Mapper SQL:

-- 插入积分明细(唯一索引 uk_order)
INSERT INTO point_detail (user_id, order_id, point, type, status, create_time)
VALUES (#{userId}, #{orderId}, #{point}, #{type}, 0, NOW());

-- 查询订单明细
SELECT * FROM point_detail WHERE order_id = #{orderId};

-- 增加用户积分
UPDATE user_point
SET total_point = total_point + #{point},
    update_time = NOW()
WHERE user_id = #{userId};

-- 更新积分明细状态
UPDATE point_detail
SET status = 1
WHERE id = #{id};

异常场景处理

查看内嵌表格

监控告警指标

  1. 重复消费次数:监控 DuplicateKeyException 异常次数

  2. 积分明细状态异常:监控 status=0 超过 5 分钟的记录

  3. 用户积分异常:监控用户积分突增或突减

  4. 消费失败率:监控 RocketMQ 消费失败率

测试用例

@SpringBootTest
public class PointConsumerTest {

    @Autowired
    private PointConsumer pointConsumer;

    @Test
    public void testIdempotency() {
        PaymentMessage message = new PaymentMessage();
        message.setUserId(1001L);
        message.setOrderId(2001L);
        message.setAmount(new BigDecimal("100"));

        // 第一次消费,应该成功
        pointConsumer.onMessage(message);

        UserPoint userPoint = userPointMapper.selectByUserId(1001L);
        assertEquals(100, userPoint.getTotalPoint());

        // 第二次消费(重复),应该跳过
        pointConsumer.onMessage(message);

        userPoint = userPointMapper.selectByUserId(1001L);
        assertEquals(100, userPoint.getTotalPoint()); // 积分不应该增加
    }
}

生产环境优化建议

  1. 积分明细表分区:按用户 ID 分区,提升查询性能

  2. 积分明细表归档:将 3 个月前的明细数据归档到历史表

  3. 积分增加批量处理:如果消息量大,可以批量处理多个订单

  4. 积分变更异步通知:积分变更后异步通知用户(如推送、短信)

3️⃣ Key Differences

查看内嵌表格


1.2 CDC 变更数据捕获

MySQL Binlog 三种格式

1、什么是 MySQL Binlog?它有哪几种格式?

(Binlog 定义、三种格式)

MySQL Binlog(Binary Log)是 MySQL 的二进制日志,记录了所有对数据库数据进行修改的操作(INSERT、UPDATE、DELETE),主要用于数据恢复、主从复制、数据同步等场景。

Binlog 有三种格式:

  1. STATEMENT:基于 SQL 语句的复制,记录执行的 SQL 语句

  2. ROW:基于行的复制,记录每一行数据的变化

  3. MIXED:混合模式,默认使用 STATEMENT,特殊情况下自动切换到 ROW

2、请详细说明 MySQL Binlog 三种格式的区别、优缺点及适用场景?

⭐⭐(三种格式的对比、优缺点分析)

1️⃣ Common Answer STATEMENT 记录 SQL 语句,ROW 记录每一行的变化,MIXED 是混合模式。STATEMENT 节省空间但可能有问题,ROW 准确但空间大,MIXED 是折中方案。

2️⃣ Impressive Answer MySQL Binlog 的三种格式各有优劣,我会从数据一致性、性能、空间占用等维度详细对比。给你详细讲一下:

格式一:STATEMENT(基于 SQL 语句)

记录执行的 SQL 语句本身。

-- 示例:Binlog 记录的内容
UPDATE user SET point = point + 100 WHERE id = 1;

优点

  • 空间占用小:只记录 SQL 语句,不记录具体数据变化

  • 网络传输快:传输的日志量小

  • 可读性强:可以直接查看执行的 SQL

缺点

  • 数据一致性风险:某些函数可能导致主从数据不一致 ```sql-- 问题示例:NOW() 函数在主从执行时间不同UPDATE user SET last_login = NOW() WHERE id = 1;

-- 问题示例:UUID() 函数每次生成不同的值 INSERT INTO user (id, name, token) VALUES (1, 'Alice', UUID());

```

  • 不确定性操作:使用 LIMIT、随机函数等可能导致不一致 ```sql-- 问题示例:LIMIT 在不同数据量下结果不同DELETE FROM log WHERE create_time < '2024-01-01' LIMIT 1000;

```

适用场景

  • 对数据一致性要求不高的场景

  • 读写分离的主从复制

  • 简单的 CRUD 操作


格式二:ROW(基于行)

记录每一行数据的变化。

-- 示例:Binlog 记录的内容
### UPDATE test.user
### WHERE
###   @1=1                    -- id
###   @2='Alice'              -- name
###   @3=100                  -- point (修改前)
### SET
###   @1=1
###   @2='Alice'
###   @3=200                  -- point (修改后)

优点

  • 数据一致性高:记录每一行的具体变化,主从数据完全一致

  • 支持所有操作:不受函数、随机值等影响

  • 精确恢复:可以精确恢复到某一行数据的状态

缺点

  • 空间占用大:记录所有行的变化,日志量巨大

  • 网络传输慢:传输的日志量大

  • 可读性差:无法直接查看执行的 SQL

适用场景

  • 对数据一致性要求高的场景(如金融、电商)

  • 数据同步到 ES、缓存等

  • 数据恢复、审计


格式三:MIXED(混合模式)

默认使用 STATEMENT,在以下情况自动切换到 ROW:

  1. 使用了不确定函数(NOW()、UUID()、RAND() 等)

  2. 使用了 LIMIT 且没有 ORDER BY

  3. 使用了 LOAD DATA INFILE

  4. 使用了 INSERT ... SELECT

  5. 使用了用户自定义函数

-- 示例:MIXED 模式下的行为
-- 正常情况:使用 STATEMENT
UPDATE user SET point = point + 100 WHERE id = 1;

-- 自动切换到 ROW:使用 NOW() 函数
UPDATE user SET last_login = NOW() WHERE id = 1;

优点

  • 兼顾性能和一致性:大部分情况用 STATEMENT 节省空间,特殊情况用 ROW 保证一致性

  • 自动优化:无需手动配置,MySQL 自动判断

缺点

  • 不可控:切换逻辑由 MySQL 控制,无法精确控制

  • 不确定性:某些情况下可能无法准确预测使用哪种格式

适用场景

  • 对性能和一致性都有一定要求的场景

  • 不想手动管理 Binlog 格式的场景


三种格式对比

查看内嵌表格

生产环境配置建议

-- 查看当前 Binlog 格式
SHOW VARIABLES LIKE 'binlog_format';

-- 设置 Binlog 格式(需要重启 MySQL)
SET GLOBAL binlog_format = 'ROW';  -- 推荐:CDC 场景
SET GLOBAL binlog_format = 'STATEMENT';  -- 主从复制场景
SET GLOBAL binlog_format = 'MIXED';  -- 折中方案

-- 推荐配置:CDC 场景
SET GLOBAL binlog_format = 'ROW';
SET GLOBAL binlog_row_image = 'FULL';  -- 记录修改前后的完整数据

CDC 场景推荐使用 ROW 格式的原因

  1. 数据一致性:保证同步到 ES、缓存的数据准确

  2. 精确变更:可以获取修改前后的完整数据

  3. 支持所有操作:不受函数、随机值等影响

  4. Canal、Debezium 等工具:都推荐使用 ROW 格式

3️⃣ Key Differences

查看内嵌表格

3、在电商系统中,需要将订单数据实时同步到 Elasticsearch,应该选择哪种 Binlog 格式?为什么?

⭐⭐⭐(Binlog 格式选择、CDC 实践)

1️⃣ Common Answer 用 ROW 格式,因为 ROW 记录每一行的变化,比较准确。STATEMENT 可能有问题,MIXED 不太好控制。

2️⃣ Impressive Answer 在电商订单同步到 ES 的场景中,必须使用 ROW 格式。让我详细分析一下原因:

场景需求分析

  • 订单数据需要实时同步到 ES,用于商品搜索、订单查询

  • 订单状态变更(如待支付 → 已支付 → 已发货)需要实时反映到 ES

  • 订单金额、用户信息等变更需要精确同步

  • 数据一致性要求高,不能出现订单状态不一致的情况

为什么不能用 STATEMENT 格式?

STATEMENT 格式记录的是 SQL 语句,存在以下问题:

-- 问题 1:使用 NOW() 函数,主从执行时间不同
UPDATE orders
SET status = 'PAID',
    pay_time = NOW()  -- 主从执行时间不同,导致 pay_time 不一致
WHERE id = 1001;

-- 问题 2:使用 UUID() 函数,每次生成不同的值
UPDATE orders
SET order_no = CONCAT(order_no, '_', UUID())  -- 主从生成不同的 UUID
WHERE id = 1001;

-- 问题 3:使用 LIMIT,可能影响不同数量的行
UPDATE orders
SET status = 'CANCELLED'
WHERE create_time < '2024-01-01'
LIMIT 1000;  -- 主从数据量不同,可能影响不同数量的行

这些问题会导致:

  • ES 数据不准确:订单状态、支付时间等数据不一致

  • 搜索结果错误:用户搜索订单时可能查不到或查到错误数据

  • 业务异常:订单状态判断错误,导致业务流程异常

为什么不能用 MIXED 格式?

MIXED 格式在大多数情况下使用 STATEMENT,只在特殊情况下切换到 ROW。但问题是:

  • 不可控:MySQL 自动判断切换逻辑,无法精确控制

  • 不确定性:某些情况下可能无法准确预测使用哪种格式

  • 风险高:如果 MySQL 判断失误,会导致数据不一致

为什么必须使用 ROW 格式?

ROW 格式记录每一行的具体变化,可以获取修改前后的完整数据:

-- 示例:订单状态变更
### UPDATE test.orders
### WHERE
###   @1=1001                -- id
###   @2='UNPAID'            -- status (修改前)
###   @3=NULL                -- pay_time (修改前)
### SET
###   @1=1001
###   @2='PAID'              -- status (修改后)
###   @3='2024-03-29 15:00:00'  -- pay_time (修改后)

ROW 格式的优势:

  1. 数据一致性高:主从数据完全一致,ES 数据准确

  2. 精确变更:可以获取修改前后的完整数据,便于 ES 更新

  3. 支持所有操作:不受函数、随机值等影响

  4. 精确恢复:如果 ES 同步失败,可以精确恢复到某个时间点

生产环境配置

-- 设置 Binlog 格式为 ROW
SET GLOBAL binlog_format = 'ROW';

-- 设置行镜像为 FULL(记录修改前后的完整数据)
SET GLOBAL binlog_row_image = 'FULL';

-- 验证配置
SHOW VARIABLES LIKE 'binlog_format';
SHOW VARIABLES LIKE 'binlog_row_image';

Canal 配置示例

# canal.properties
canal.serverMode = rocketMQ
canal.mq.servers = 127.0.0.1:9876
canal.mq.producerGroup = canal_producer

# instance.properties
canal.instance.master.address=127.0.0.1:3306
canal.instance.master.journal.name=
canal.instance.master.position=
canal.instance.master.timestamp=
canal.instance.dbUsername=canal
canal.instance.dbPassword=canal
canal.instance.connectionCharset=UTF-8
canal.instance.filter.regex=.*\\..*  # 监听所有表
canal.instance.filter.black.regex=

消费者代码示例(同步订单到 ES)

@RocketMQMessageListener(topic = "canal_order", consumerGroup = "es_sync_group")
public class OrderToEsConsumer implements RocketMQListener<CanalMessage> {

    @Autowired
    private ElasticsearchRestTemplate esTemplate;

    @Override
    public void onMessage(CanalMessage message) {
        List<CanalEntry.Entry> entries = message.getEntries();

        for (CanalEntry.Entry entry : entries) {
            if (entry.getEntryType() == CanalEntry.EntryType.ROWDATA) {
                CanalEntry.RowChange rowChange = CanalEntry.RowChange.parseFrom(entry.getStoreValue());

                for (CanalEntry.RowData rowData : rowChange.getRowDatasList()) {
                    String tableName = entry.getHeader().getTableName();

                    if ("orders".equals(tableName)) {
                        if (rowChange.getEventType() == CanalEntry.EventType.INSERT) {
                            // 新增订单,插入到 ES
                            Order order = parseOrder(rowData.getAfterColumnsList());
                            esTemplate.save(order);

                        } else if (rowChange.getEventType() == CanalEntry.EventType.UPDATE) {
                            // 更新订单,更新到 ES
                            Order order = parseOrder(rowData.getAfterColumnsList());
                            esTemplate.save(order);

                        } else if (rowChange.getEventType() == CanalEntry.EventType.DELETE) {
                            // 删除订单,从 ES 删除
                            Long orderId = parseOrderId(rowData.getBeforeColumnsList());
                            esTemplate.delete(orderId, Order.class);
                        }
                    }
                }
            }
        }
    }
}

异常场景处理

查看内嵌表格

监控告警指标

  1. Binlog 延迟:监控 Canal 消费 Binlog 的延迟时间

  2. ES 同步延迟:监控消息从 MQ 到 ES 的延迟时间

  3. 同步失败率:监控 ES 同步失败的次数

  4. 数据一致性:定时比对 DB 和 ES 的数据一致性

3️⃣ Key Differences

查看内嵌表格

4、容易一起考的题

查看内嵌表格


Canal 架构原理

1、基础题:什么是 Canal?它解决了什么问题?

(Canal 定义、应用场景)

Canal 是阿里巴巴开源的 MySQL Binlog 增量订阅&消费组件,通过模拟 MySQL Slave 的交互协议,伪装成 MySQL Slave,向 MySQL Master 发送 dump 协议,MySQL Master 接收到 dump 请求后推送 Binlog 给 Canal,Canal 解析 Binlog 后发送到消息队列或存储到其他系统。

主要应用场景:

  • 数据库镜像:数据库实时备份

  • 数据同步:将 MySQL 数据同步到 ES、Redis、MongoDB 等

  • 缓存更新:数据库变更后自动更新缓存

  • 数据解耦:将数据库变更事件化,解耦业务系统

2、进阶题:请详细说明 Canal 的架构原理,以及它如何保证数据不丢失?

⭐⭐(Canal 架构、数据可靠性保证)

1️⃣ Common Answer Canal 就是模拟 MySQL Slave,向 Master 发送 dump 请求,Master 推送 Binlog 给 Canal,Canal 解析后发送到 MQ。为了保证数据不丢失,Canal 会记录消费位点,如果消费失败可以重试。

2️⃣ Impressive Answer Canal 的核心原理是模拟 MySQL Slave 的交互协议,通过EventParser 解析 BinlogEventSink 过滤和路由EventStore 存储事件EventAck 确认机制来保证数据可靠性。让我详细讲一下:

Canal 整体架构

graph TD
    A[MySQL Master] -->|Binlog| B[Canal Server]
    B --> C[EventParser]
    C --> D[EventSink]
    D --> E[EventStore]
    E --> F[Canal Client]
    E --> G[RocketMQ/Kafka]
    F --> H[业务系统]
    G --> H

    subgraph Canal Server
    B
    C
    D
    E
    end

    subgraph Canal Client
    F
    end

核心组件详解

  1. EventParser(事件解析器)

负责模拟 MySQL Slave,向 MySQL Master 发送 dump 协议,接收并解析 Binlog。

// EventParser 核心逻辑
public class EventParser {

    private MySQLConnection mysqlConnection;

    public void start() {
        // 1. 连接 MySQL Master
        mysqlConnection.connect();

        // 2. 发送 dump 协议,请求 Binlog
        BinlogDumpCommand dumpCommand = new BinlogDumpCommand();
        dumpCommand.setBinlogFileName("mysql-bin.000001");
        dumpCommand.setBinlogPosition(1000);
        mysqlConnection.sendCommand(dumpCommand);

        // 3. 接收 Binlog Event
        while (true) {
            BinlogEvent event = mysqlConnection.receiveEvent();

            // 4. 解析 Binlog Event
            if (event instanceof QueryEvent) {
                // QUERY 事件(DDL)
                parseQueryEvent((QueryEvent) event);
            } else if (event instanceof TableMapEvent) {
                // TABLE_MAP 事件(表结构映射)
                parseTableMapEvent((TableMapEvent) event);
            } else if (event instanceof WriteRowsEvent) {
                // WRITE_ROWS 事件(INSERT)
                parseWriteRowsEvent((WriteRowsEvent) event);
            } else if (event instanceof UpdateRowsEvent) {
                // UPDATE_ROWS 事件(UPDATE)
                parseUpdateRowsEvent((UpdateRowsEvent) event);
            } else if (event instanceof DeleteRowsEvent) {
                // DELETE_ROWS 事件(DELETE)
                parseDeleteRowsEvent((DeleteRowsEvent) event);
            }
        }
    }
}
  1. EventSink(事件过滤器)

负责过滤和路由 Binlog 事件,支持表名过滤、字段过滤等。

// EventSink 核心逻辑
public class EventSink {

    private List<CanalEventFilter> filters;
    private EventStore eventStore;

    public boolean sink(CanalEntry.Entry entry) {
        // 1. 表名过滤
        String tableName = entry.getHeader().getTableName();
        if (!isTableAllowed(tableName)) {
            return false;
        }

        // 2. 字段过滤
        if (entry.getEntryType() == CanalEntry.EntryType.ROWDATA) {
            CanalEntry.RowChange rowChange = CanalEntry.RowChange.parseFrom(entry.getStoreValue());
            filterColumns(rowChange);
        }

        // 3. 路由到不同的 EventStore
        routeToEventStore(entry);

        return true;
    }
}
  1. EventStore(事件存储器)

负责存储 Binlog 事件,支持内存存储、文件存储、混合存储。

// EventStore 核心逻辑
public class MemoryEventStore implements EventStore {

    private RingBuffer<CanalEntry.Entry> ringBuffer;
    private PositionManager positionManager;

    public void put(CanalEntry.Entry entry) {
        // 1. 存储事件到 RingBuffer
        ringBuffer.put(entry);

        // 2. 记录位点
        LogPosition position = buildPosition(entry);
        positionManager.persist(position);
    }

    public CanalEntry.Entry get(Position position) {
        // 1. 从 RingBuffer 获取事件
        CanalEntry.Entry entry = ringBuffer.get(position);

        // 2. 更新消费位点
        positionManager.update(position);

        return entry;
    }
}
  1. EventAck(确认机制)

负责确认消费位点,保证数据不丢失。

// EventAck 核心逻辑
public class EventAck {

    private PositionManager positionManager;

    public void ack(LogPosition position) {
        // 1. 更新消费位点
        positionManager.update(position);

        // 2. 持久化位点(防止重启丢失)
        positionManager.persist(position);

        // 3. 定时清理已确认的事件
        cleanConfirmedEvents(position);
    }
}

数据可靠性保证机制

  1. 位点记录(Position)

Canal 记录了三个关键位点:

  • Binlog 文件名mysql-bin.000001

  • Binlog 位置1000

  • 时间戳1648560000000

// 位点记录
public class LogPosition {
    private String journalName;      // Binlog 文件名
    private Long position;           // Binlog 位置
    private Long timestamp;          // 时间戳
    private String gtid;             // GTID(MySQL 5.6+)
}
  1. 位点持久化

位点存储在文件中,防止 Canal 重启后丢失:

# canal.properties
canal.instance.master.journal.name=mysql-bin.000001
canal.instance.master.position=1000
canal.instance.master.timestamp=1648560000000
  1. 消费确认机制

Canal Client 消费成功后发送 ACK 确认:

// Canal Client 消费代码
CanalConnector connector = CanalConnectors.newSingleConnector(
    new InetSocketAddress("127.0.0.1", 11111),
    "example",
    "",
    ""
);

connector.connect();
connector.subscribe(".*\\..*");

while (true) {
    Message message = connector.getWithoutAck(100);  // 获取 100 条消息
    long batchId = message.getId();

    try {
        List<CanalEntry.Entry> entries = message.getEntries();

        // 处理消息
        for (CanalEntry.Entry entry : entries) {
            // 业务处理
        }

        // 处理成功,ACK 确认
        connector.ack(batchId);

    } catch (Exception e) {
        // 处理失败,不 ACK,下次继续消费
        connector.rollback(batchId);
    }
}
  1. HA 高可用

Canal 支持 HA 高可用,主节点宕机后自动切换:

# canal.properties
canal.register.ip = 127.0.0.1
canal.port = 11111
canal.zkServers = 127.0.0.1:2181  # ZooKeeper 地址
canal.cluster.mode = true  # 集群模式

异常场景处理

查看内嵌表格

生产环境配置建议

# canal.properties
canal.serverMode = rocketMQ  # 使用 RocketMQ
canal.mq.servers = 127.0.0.1:9876
canal.mq.producerGroup = canal_producer
canal.mq.topic = canal_order  # Topic 名称
canal.mq.partition = 0  # 分区数

# instance.properties
canal.instance.master.address=127.0.0.1:3306
canal.instance.dbUsername=canal
canal.instance.dbPassword=canal
canal.instance.connectionCharset=UTF-8
canal.instance.filter.regex=.*\\..*  # 监听所有表
canal.instance.filter.black.regex=mysql\\.slave_.*  # 黑名单

# 性能优化
canal.instance.network.receiveBufferSize = 16384
canal.instance.network.sendBufferSize = 16384
canal.instance.memory.batch.mode = MEMSIZE  # 内存批处理模式
canal.instance.memory.buffer.size = 32768  # 缓冲区大小

监控告警指标

  1. Binlog 延迟:监控 Canal 消费 Binlog 的延迟时间

  2. 消费堆积:监控 EventStore 的堆积量

  3. 消费失败率:监控消费失败的次数

  4. HA 切换次数:监控主节点切换次数

3️⃣ Key Differences

查看内嵌表格

3、场景题:在电商系统中,需要将订单数据实时同步到 Elasticsearch,如何用 Canal 实现?

⭐⭐⭐(Canal 实战、数据同步方案)

1️⃣ Common Answer 用 Canal 监听 MySQL 的 Binlog,解析订单表的变更,然后发送到 RocketMQ,消费者消费消息后更新到 ES。如果消费失败就重试。

2️⃣ Impressive Answer 这是一个典型的 CDC 场景,我会用 Canal + RocketMQ + Elasticsearch 的架构来实现实时同步。让我详细设计一下:

整体架构设计

graph TD
    A[MySQL Master] -->|Binlog| B[Canal Server]
    B -->|解析 Binlog| C[RocketMQ]
    C -->|消息| D[ES Sync Consumer]
    D -->|更新| E[Elasticsearch]
    F[订单服务] -->|写入订单| A

    subgraph Canal Server
    B
    end

    subgraph 消息队列
    C
    end

    subgraph 消费者
    D
    end

    subgraph 搜索引擎
    E
    end

核心实现步骤

步骤 1:配置 MySQL

-- 确保 MySQL 开启 Binlog
SHOW VARIABLES LIKE 'log_bin';

-- 设置 Binlog 格式为 ROW
SET GLOBAL binlog_format = 'ROW';

-- 设置行镜像为 FULL
SET GLOBAL binlog_row_image = 'FULL';

-- 创建 Canal 用户
CREATE USER 'canal'@'%' IDENTIFIED BY 'canal';
GRANT SELECT, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO 'canal'@'%';
FLUSH PRIVILEGES;

-- 验证权限
SHOW GRANTS FOR 'canal'@'%';

步骤 2:配置 Canal Server

# canal.properties
canal.serverMode = rocketMQ
canal.mq.servers = 127.0.0.1:9876
canal.mq.producerGroup = canal_producer
canal.mq.topic = canal_order
canal.mq.partition = 0
canal.mq.partitionHash = .*\\..*:id  # 按 id 分区

# instance.properties
canal.instance.master.address=127.0.0.1:3306
canal.instance.master.journal.name=
canal.instance.master.position=
canal.instance.master.timestamp=
canal.instance.dbUsername=canal
canal.instance.dbPassword=canal
canal.instance.connectionCharset=UTF-8
canal.instance.filter.regex=shop\\.orders  # 只监听 shop.orders 表
canal.instance.filter.black.regex=
canal.instance.network.receiveBufferSize = 16384
canal.instance.network.sendBufferSize = 16384
canal.instance.memory.batch.mode = MEMSIZE
canal.instance.memory.buffer.size = 32768

步骤 3:消费者代码(同步订单到 ES)

@RocketMQMessageListener(topic = "canal_order", consumerGroup = "es_sync_group")
public class OrderToEsConsumer implements RocketMQListener<MessageExt> {

    @Autowired
    private ElasticsearchRestTemplate esTemplate;

    @Autowired
    private OrderMapper orderMapper;

    @Override
    public void onMessage(MessageExt message) {
        try {
            // 1. 解析 Canal 消息
            CanalMessage canalMessage = JSON.parseObject(new String(message.getBody()), CanalMessage.class);

            for (CanalEntry.Entry entry : canalMessage.getEntries()) {
                if (entry.getEntryType() == CanalEntry.EntryType.ROWDATA) {
                    CanalEntry.RowChange rowChange = CanalEntry.RowChange.parseFrom(entry.getStoreValue());
                    String tableName = entry.getHeader().getTableName();

                    if ("orders".equals(tableName)) {
                        handleOrderChange(rowChange);
                    }
                }
            }

        } catch (Exception e) {
            log.error("同步订单到 ES 失败", e);
            // 抛出异常,触发 RocketMQ 重试
            throw e;
        }
    }

    private void handleOrderChange(CanalEntry.RowChange rowChange) {
        CanalEntry.EventType eventType = rowChange.getEventType();

        if (eventType == CanalEntry.EventType.INSERT) {
            // 新增订单
            for (CanalEntry.RowData rowData : rowChange.getRowDatasList()) {
                Order order = parseOrder(rowData.getAfterColumnsList());
                esTemplate.save(order);
                log.info("新增订单同步到 ES 成功,orderId={}", order.getId());
            }

        } else if (eventType == CanalEntry.EventType.UPDATE) {
            // 更新订单
            for (CanalEntry.RowData rowData : rowChange.getRowDatasList()) {
                Order order = parseOrder(rowData.getAfterColumnsList());
                esTemplate.save(order);
                log.info("更新订单同步到 ES 成功,orderId={}", order.getId());
            }

        } else if (eventType == CanalEntry.EventType.DELETE) {
            // 删除订单
            for (CanalEntry.RowData rowData : rowChange.getRowDatasList()) {
                Long orderId = parseOrderId(rowData.getBeforeColumnsList());
                esTemplate.delete(orderId, Order.class);
                log.info("删除订单同步到 ES 成功,orderId={}", orderId);
            }
        }
    }

    private Order parseOrder(List<CanalEntry.Column> columns) {
        Order order = new Order();
        for (CanalEntry.Column column : columns) {
            String name = column.getName();
            String value = column.getValue();

            switch (name) {
                case "id":
                    order.setId(Long.parseLong(value));
                    break;
                case "order_no":
                    order.setOrderNo(value);
                    break;
                case "user_id":
                    order.setUserId(Long.parseLong(value));
                    break;
                case "status":
                    order.setStatus(value);
                    break;
                case "total_amount":
                    order.setTotalAmount(new BigDecimal(value));
                    break;
                case "create_time":
                    order.setCreateTime(parseDateTime(value));
                    break;
                case "update_time":
                    order.setUpdateTime(parseDateTime(value));
                    break;
            }
        }
        return order;
    }

    private Long parseOrderId(List<CanalEntry.Column> columns) {
        for (CanalEntry.Column column : columns) {
            if ("id".equals(column.getName())) {
                return Long.parseLong(column.getValue());
            }
        }
        return null;
    }
}

步骤 4:Elasticsearch 文档设计

@Document(indexName = "orders")
public class Order {

    @Id
    private Long id;

    @Field(type = FieldType.Keyword)
    private String orderNo;

    @Field(type = FieldType.Long)
    private Long userId;

    @Field(type = FieldType.Keyword)
    private String status;

    @Field(type = FieldType.Double)
    private BigDecimal totalAmount;

    @Field(type = FieldType.Date, format = DateFormat.date_hour_minute_second)
    private Date createTime;

    @Field(type = FieldType.Date, format = DateFormat.date_hour_minute_second)
    private Date updateTime;

    // getter/setter...
}

异常场景处理

查看内嵌表格

监控告警指标

  1. Canal 延迟:监控 Canal 消费 Binlog 的延迟时间

  2. ES 同步延迟:监控消息从 MQ 到 ES 的延迟时间

  3. 同步失败率:监控 ES 同步失败的次数

  4. 数据一致性:定时比对 DB 和 ES 的数据一致性

性能优化建议

  1. 批量消费:设置 consumeMessageBatchMaxSize=10,批量消费消息

  2. 异步消费:使用 @RocketMQMessageListener(consumeMode = ConsumeMode.CONCURRENTLY) 并发消费

  3. ES 批量写入:使用 esTemplate.saveAll() 批量写入 ES

  4. 索引分区:按时间分区索引,如 orders_202403orders_202404

测试用例

@SpringBootTest
public class OrderToEsConsumerTest {

    @Autowired
    private OrderMapper orderMapper;

    @Autowired
    private ElasticsearchRestTemplate esTemplate;

    @Test
    public void testSyncOrderToEs() {
        // 1. 插入订单到 MySQL
        Order order = new Order();
        order.setOrderNo("ORD20240329001");
        order.setUserId(1001L);
        order.setStatus("UNPAID");
        order.setTotalAmount(new BigDecimal("100.00"));
        order.setCreateTime(new Date());
        order.setUpdateTime(new Date());
        orderMapper.insert(order);

        // 2. 等待 Canal 同步
        Thread.sleep(5000);

        // 3. 查询 ES
        Order esOrder = esTemplate.get(order.getId(), Order.class);

        // 4. 验证数据
        assertNotNull(esOrder);
        assertEquals(order.getOrderNo(), esOrder.getOrderNo());
        assertEquals(order.getUserId(), esOrder.getUserId());
        assertEquals(order.getStatus(), esOrder.getStatus());
    }
}

3️⃣ Key Differences

查看内嵌表格

4、容易一起考的题

查看内嵌表格


DB 变更同步 ES/缓存

1、基础题:为什么要将数据库变更同步到 Elasticsearch 和缓存?

(同步原因、应用场景)

将数据库变更同步到 Elasticsearch 和缓存的主要原因:

  1. 性能优化
  2. Elasticsearch 支持全文搜索、复杂查询,查询性能远高于 MySQL
  3. 缓存(Redis)读取速度极快,可以减轻数据库压力

  4. 读写分离

  5. MySQL 用于写操作(OLTP)
  6. Elasticsearch 用于读操作(OLAP)
  7. 缓存用于热点数据读取

  8. 业务解耦

  9. 数据库变更后自动同步到 ES 和缓存,业务代码无需关心同步逻辑
  10. 避免业务代码中耦合同步逻辑

  11. 数据一致性

  12. 通过 CDC(Canal)实时同步,保证数据一致性
  13. 避免定时同步带来的数据延迟

2、进阶题:请设计一个数据库变更同步到 Elasticsearch 和缓存的方案,并说明如何保证数据一致性?

⭐⭐(同步方案设计、数据一致性保证)

1️⃣ Common Answer 用 Canal 监听 MySQL 的 Binlog,解析变更后发送到 RocketMQ,然后有两个消费者,一个更新 ES,一个更新缓存。如果消费失败就重试,保证最终一致性。

2️⃣ Impressive Answer 我会设计一个 Canal + RocketMQ + 多消费者 的架构,实现数据库变更实时同步到 ES 和缓存。让我详细讲一下:

整体架构设计

graph TD
    A[MySQL Master] -->|Binlog| B[Canal Server]
    B -->|解析 Binlog| C[RocketMQ Topic: canal_order]
    C -->|消息| D[ES Sync Consumer]
    C -->|消息| E[Cache Sync Consumer]
    D -->|更新| F[Elasticsearch]
    E -->|更新| G[Redis Cache]
    H[订单服务] -->|写入订单| A

    subgraph Canal Server
    B
    end

    subgraph 消息队列
    C
    end

    subgraph 消费者
    D
    E
    end

    subgraph 存储系统
    F
    G
    end

核心实现步骤

步骤 1:配置 Canal

# canal.properties
canal.serverMode = rocketMQ
canal.mq.servers = 127.0.0.1:9876
canal.mq.producerGroup = canal_producer
canal.mq.topic = canal_order
canal.mq.partition = 0
canal.mq.partitionHash = .*\\..*:id

# instance.properties
canal.instance.master.address=127.0.0.1:3306
canal.instance.dbUsername=canal
canal.instance.dbPassword=canal
canal.instance.connectionCharset=UTF-8
canal.instance.filter.regex=shop\\.orders

步骤 2:ES 同步消费者

@RocketMQMessageListener(topic = "canal_order", consumerGroup = "es_sync_group")
public class OrderToEsConsumer implements RocketMQListener<MessageExt> {

    @Autowired
    private ElasticsearchRestTemplate esTemplate;

    @Override
    public void onMessage(MessageExt message) {
        try {
            CanalMessage canalMessage = JSON.parseObject(new String(message.getBody()), CanalMessage.class);

            for (CanalEntry.Entry entry : canalMessage.getEntries()) {
                if (entry.getEntryType() == CanalEntry.EntryType.ROWDATA) {
                    CanalEntry.RowChange rowChange = CanalEntry.RowChange.parseFrom(entry.getStoreValue());
                    String tableName = entry.getHeader().getTableName();

                    if ("orders".equals(tableName)) {
                        handleOrderChange(rowChange);
                    }
                }
            }

        } catch (Exception e) {
            log.error("同步订单到 ES 失败", e);
            throw e;
        }
    }

    private void handleOrderChange(CanalEntry.RowChange rowChange) {
        CanalEntry.EventType eventType = rowChange.getEventType();

        for (CanalEntry.RowData rowData : rowChange.getRowDatasList()) {
            if (eventType == CanalEntry.EventType.INSERT || eventType == CanalEntry.EventType.UPDATE) {
                // 新增或更新订单
                Order order = parseOrder(rowData.getAfterColumnsList());
                esTemplate.save(order);

            } else if (eventType == CanalEntry.EventType.DELETE) {
                // 删除订单
                Long orderId = parseOrderId(rowData.getBeforeColumnsList());
                esTemplate.delete(orderId, Order.class);
            }
        }
    }
}

步骤 3:缓存同步消费者

@RocketMQMessageListener(topic = "canal_order", consumerGroup = "cache_sync_group")
public class OrderToCacheConsumer implements RocketMQListener<MessageExt> {

    @Autowired
    private RedisTemplate<String, Object> redisTemplate;

    @Override
    public void onMessage(MessageExt message) {
        try {
            CanalMessage canalMessage = JSON.parseObject(new String(message.getBody()), CanalMessage.class);

            for (CanalEntry.Entry entry : canalMessage.getEntries()) {
                if (entry.getEntryType() == CanalEntry.EntryType.ROWDATA) {
                    CanalEntry.RowChange rowChange = CanalEntry.RowChange.parseFrom(entry.getStoreValue());
                    String tableName = entry.getHeader().getTableName();

                    if ("orders".equals(tableName)) {
                        handleOrderChange(rowChange);
                    }
                }
            }

        } catch (Exception e) {
            log.error("同步订单到缓存失败", e);
            throw e;
        }
    }

    private void handleOrderChange(CanalEntry.RowChange rowChange) {
        CanalEntry.EventType eventType = rowChange.getEventType();

        for (CanalEntry.RowData rowData : rowChange.getRowDatasList()) {
            String orderKey = "order:" + parseOrderId(rowData.getAfterColumnsList());

            if (eventType == CanalEntry.EventType.INSERT || eventType == CanalEntry.EventType.UPDATE) {
                // 新增或更新订单
                Order order = parseOrder(rowData.getAfterColumnsList());
                redisTemplate.opsForValue().set(orderKey, order, 1, TimeUnit.HOURS);

            } else if (eventType == CanalEntry.EventType.DELETE) {
                // 删除订单
                redisTemplate.delete(orderKey);
            }
        }
    }
}

数据一致性保证机制

  1. 消息幂等性

ES 和缓存都使用唯一 ID 作为主键,天然支持幂等性:

// ES 使用文档 ID
@Document(indexName = "orders")
public class Order {
    @Id
    private Long id;  // 唯一 ID
}

// 缓存使用唯一 Key
String orderKey = "order:" + orderId;
  1. 消息重试机制

RocketMQ 自动重试,直到成功或达到最大重试次数:

// 消费失败时抛出异常,触发重试
catch (Exception e) {
    log.error("同步失败", e);
    throw e;  // 触发 RocketMQ 重试
}
  1. 死信队列处理

重试失败的消息进入死信队列,人工介入处理:

// 监听死信队列
@RocketMQMessageListener(topic = "%DLQ%canal_order", consumerGroup = "dlq_handler_group")
public class DlqHandler implements RocketMQListener<MessageExt> {

    @Override
    public void onMessage(MessageExt message) {
        // 人工介入处理死信消息
        log.error("死信消息,需要人工处理,message={}", message);
        alertService.sendAlert("死信消息:" + message.getMsgId());
    }
}
  1. 数据一致性校验

定时任务校验 DB、ES、缓存的数据一致性:

@Component
public class DataConsistencyCheck {

    @Scheduled(cron = "0 0 2 * * ?")  // 每天凌晨 2 点执行
    public void checkConsistency() {
        // 1. 查询 DB 中的订单
        List<Order> dbOrders = orderMapper.selectAll();

        // 2. 查询 ES 中的订单
        List<Order> esOrders = esTemplate.searchForList(Query.findAll(), Order.class);

        // 3. 查询缓存中的订单
        List<Order> cacheOrders = new ArrayList<>();
        for (Order order : dbOrders) {
            String orderKey = "order:" + order.getId();
            Order cacheOrder = (Order) redisTemplate.opsForValue().get(orderKey);
            if (cacheOrder != null) {
                cacheOrders.add(cacheOrder);
            }
        }

        // 4. 比对数据一致性
        List<Order> inconsistentOrders = findInconsistentOrders(dbOrders, esOrders, cacheOrders);

        // 5. 修复不一致数据
        for (Order order : inconsistentOrders) {
            fixInconsistentOrder(order);
        }
    }
}

异常场景处理

查看内嵌表格

性能优化建议

  1. 批量消费:设置 consumeMessageBatchMaxSize=10,批量消费消息

  2. 异步消费:使用 @RocketMQMessageListener(consumeMode = ConsumeMode.CONCURRENTLY) 并发消费

  3. ES 批量写入:使用 esTemplate.saveAll() 批量写入 ES

  4. 缓存批量更新:使用 redisTemplate.opsForValue().multiSet() 批量更新缓存

  5. 消息分区:按订单 ID 分区,保证同一订单的消息顺序性

// 批量消费配置
@RocketMQMessageListener(
    topic = "canal_order",
    consumerGroup = "es_sync_group",
    consumeMode = ConsumeMode.CONCURRENTLY,
    consumeThreadMax = 10  // 最大消费线程数
)

监控告警指标

  1. Canal 延迟:监控 Canal 消费 Binlog 的延迟时间

  2. ES 同步延迟:监控消息从 MQ 到 ES 的延迟时间

  3. 缓存同步延迟:监控消息从 MQ 到缓存的延迟时间

  4. 同步失败率:监控 ES 和缓存同步失败的次数

  5. 数据一致性:定时比对 DB、ES、缓存的数据一致性

3️⃣ Key Differences

查看内嵌表格

3、场景题:在电商系统中,订单状态变更后需要实时更新缓存,如何保证缓存和数据库的一致性?

⭐⭐⭐(缓存一致性、实战场景)

1️⃣ Common Answer 用 Canal 监听订单表的变更,状态变更后发送消息到 MQ,消费者更新缓存。如果缓存更新失败就重试,保证最终一致性。

2️⃣ Impressive Answer 订单状态变更的缓存一致性是一个经典问题,我会用 Canal + Redis + 延迟双删 的方案来保证强一致性。让我详细设计一下:

场景需求分析

  • 订单状态变更(如待支付 → 已支付 → 已发货)需要实时更新缓存

  • 用户查询订单时优先从缓存读取,缓存不存在时从 DB 读取并回写缓存

  • 需要保证缓存和数据库的强一致性,避免用户看到过期数据

整体架构设计

sequenceDiagram
    participant User as 用户
    participant Order as 订单服务
    participant DB as MySQL
    participant Canal as Canal Server
    participant MQ as RocketMQ
    participant Consumer as 缓存消费者
    participant Cache as Redis Cache

    User->>Order: 1. 更新订单状态
    Order->>DB: 2. UPDATE orders SET status='PAID' WHERE id=1001
    DB-->>Order: 3. 更新成功
    DB->>Canal: 4. Binlog 事件
    Canal->>MQ: 5. 发送消息
    MQ->>Consumer: 6. 消费消息
    Consumer->>Cache: 7. 延迟双删<br/>先删除缓存
    Consumer->>Cache: 8. 延迟 500ms 后<br/>再次删除缓存
    Cache-->>Consumer: 9. 删除成功
    Order-->>User: 10. 返回更新成功

    Note over User,Cache: 用户下次查询时<br/>从 DB 读取最新数据<br/>并回写缓存

核心实现步骤

步骤 1:配置 Canal 监听订单表

# instance.properties
canal.instance.master.address=127.0.0.1:3306
canal.instance.dbUsername=canal
canal.instance.dbPassword=canal
canal.instance.connectionCharset=UTF-8
canal.instance.filter.regex=shop\\.orders  # 只监听订单表

步骤 2:缓存更新消费者(延迟双删)

@RocketMQMessageListener(topic = "canal_order", consumerGroup = "cache_sync_group")
public class OrderCacheConsumer implements RocketMQListener<MessageExt> {

    @Autowired
    private RedisTemplate<String, Object> redisTemplate;

    @Autowired
    private OrderMapper orderMapper;

    @Override
    public void onMessage(MessageExt message) {
        try {
            CanalMessage canalMessage = JSON.parseObject(new String(message.getBody()), CanalMessage.class);

            for (CanalEntry.Entry entry : canalMessage.getEntries()) {
                if (entry.getEntryType() == CanalEntry.EntryType.ROWDATA) {
                    CanalEntry.RowChange rowChange = CanalEntry.RowChange.parseFrom(entry.getStoreValue());
                    String tableName = entry.getHeader().getTableName();

                    if ("orders".equals(tableName)) {
                        handleOrderChange(rowChange);
                    }
                }
            }

        } catch (Exception e) {
            log.error("更新缓存失败", e);
            throw e;
        }
    }

    private void handleOrderChange(CanalEntry.RowChange rowChange) {
        CanalEntry.EventType eventType = rowChange.getEventType();

        for (CanalEntry.RowData rowData : rowChange.getRowDatasList()) {
            Long orderId = parseOrderId(rowData.getAfterColumnsList());
            String orderKey = "order:" + orderId;

            if (eventType == CanalEntry.EventType.UPDATE) {
                // 更新订单,执行延迟双删
                delayedDoubleDelete(orderId, orderKey);

            } else if (eventType == CanalEntry.EventType.DELETE) {
                // 删除订单,直接删除缓存
                redisTemplate.delete(orderKey);
            }
        }
    }

    private void delayedDoubleDelete(Long orderId, String orderKey) {
        // 第一次删除缓存
        redisTemplate.delete(orderKey);
        log.info("第一次删除缓存成功,orderId={}", orderId);

        // 延迟 500ms 后第二次删除缓存
        ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor();
        scheduler.schedule(() -> {
            try {
                // 第二次删除缓存
                redisTemplate.delete(orderKey);
                log.info("第二次删除缓存成功,orderId={}", orderId);
            } catch (Exception e) {
                log.error("第二次删除缓存失败,orderId={}", orderId, e);
            }
        }, 500, TimeUnit.MILLISECONDS);
    }
}

步骤 3:订单查询接口(Cache Aside 模式)

@Service
public class OrderService {

    @Autowired
    private OrderMapper orderMapper;

    @Autowired
    private RedisTemplate<String, Object> redisTemplate;

    public Order getOrderById(Long orderId) {
        String orderKey = "order:" + orderId;

        // 1. 先查缓存
        Order order = (Order) redisTemplate.opsForValue().get(orderKey);

        if (order != null) {
            log.info("缓存命中,orderId={}", orderId);
            return order;
        }

        // 2. 缓存不存在,查数据库
        log.info("缓存未命中,查询数据库,orderId={}", orderId);
        order = orderMapper.selectById(orderId);

        if (order != null) {
            // 3. 回写缓存
            redisTemplate.opsForValue().set(orderKey, order, 1, TimeUnit.HOURS);
            log.info("回写缓存成功,orderId={}", orderId);
        }

        return order;
    }
}

为什么需要延迟双删?

延迟双删是为了解决并发读写导致的数据不一致问题:

sequenceDiagram
    participant Thread1 as 线程1(写)
    participant DB as 数据库
    participant Cache as 缓存
    participant Thread2 as 线程2(读)

    Thread1->>DB: 1. 更新 DB(status=PAID)
    Thread1->>Cache: 2. 删除缓存
    Thread2->>Cache: 3. 查询缓存(未命中)
    Thread2->>DB: 4. 查询 DB(status=PAID)
    DB-->>Thread2: 5. 返回 PAID
    Thread2->>Cache: 6. 回写缓存(status=PAID)
    Note over Thread1,Thread2: 此时缓存和 DB 一致

    Thread1->>Cache: 7. 延迟双删<br/>第二次删除缓存
    Note over Thread1,Thread2: 此时缓存被删除<br/>下次查询时会重新回写

如果没有延迟双删,可能出现以下问题:

  1. 线程 1 更新 DB(status=PAID)

  2. 线程 1 删除缓存

  3. 线程 2 查询缓存(未命中)

  4. 线程 2 查询 DB(status=PAID)

  5. 线程 2 回写缓存(status=PAID)

  6. 此时缓存和 DB 一致

但如果线程 1 的 Canal 消息消费延迟,可能出现:

  1. 线程 1 更新 DB(status=PAID)

  2. 线程 1 删除缓存

  3. 线程 2 查询缓存(未命中)

  4. 线程 2 查询 DB(status=PAID)

  5. 线程 2 回写缓存(status=PAID)

  6. 线程 1 的 Canal 消息消费成功,删除缓存

  7. 此时缓存被删除,下次查询时会重新回写

延迟双删可以保证在 Canal 消息消费期间,如果其他线程查询了 DB 并回写了缓存,延迟双删会再次删除缓存,确保下次查询时获取最新数据。

异常场景处理

查看内嵌表格

监控告警指标

  1. Canal 延迟:监控 Canal 消费 Binlog 的延迟时间

  2. 缓存更新延迟:监控消息从 MQ 到缓存的延迟时间

  3. 缓存命中率:监控缓存的命中率

  4. 数据一致性:定时比对 DB 和缓存的数据一致性

性能优化建议

  1. 批量删除缓存:使用 redisTemplate.delete(keys) 批量删除缓存

  2. 异步删除:使用线程池异步执行延迟双删

  3. 缓存预热:系统启动时预热热点订单数据

  4. 缓存分区:按订单 ID 分区缓存,提升查询性能

// 批量删除缓存
public void batchDeleteCache(List<Long> orderIds) {
    List<String> keys = orderIds.stream()
        .map(orderId -> "order:" + orderId)
        .collect(Collectors.toList());
    redisTemplate.delete(keys);
}

3️⃣ Key Differences

查看内嵌表格

4、容易一起考的题

查看内嵌表格