跳转至

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

1.1 事务消息与本地消息表

RocketMQ 事务消息两阶段提交

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

RocketMQ 事务消息是一种让“发送消息”和“本地数据库操作”成为一个原子整体的分布式事务方案。 它解决的是分布式系统中经典的“发消息与数据库操作不一致”的痛点。

举个场景:用户下单,订单服务需要同时做两件事——1)在自己的数据库中创建订单;2)发送一条“下单成功”的消息,让下游的库存系统去扣库存。

如果不用事务消息,最常见的两种错误是:

  • 先操作数据库,再发消息:数据库更新成功,但发消息时网络闪断,消息没发出去,下游永远不知道订单已创建,库存就不会扣。

  • 先发消息,再操作数据库:消息发出去了,库存已经开始扣减,但订单服务自己的数据库操作失败,最后订单没创建,库存却少了。

事务消息通过 “半消息” 这个巧妙的中间状态,让消息暂时不会被消费,直到本地事务尘埃落定,从而解决了这个原子性问题。


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

RocketMQ 事务消息采用 两阶段提交 + 事务回查 的机制。我把整个流程画成一张图,再逐步拆解。

image.png

步骤详解:

  1. 半消息发送:生产者发送一条“半消息”到 Broker。Broker 将消息持久化到磁盘(CommitLog),但将其标记为 “暂不可投递”,所以消费者完全看不到它。这是保证消息可靠的第一步,因为已经落盘了。

  2. 执行本地事务:半消息发送成功后,生产者执行本地数据库操作(例如插入订单)。

  3. 提交或回滚:根据本地事务的结果,生产者向 Broker 发送 CommitRollback 请求。

  4. 如果本地事务成功,Commit 后,Broker 将半消息标记为正常消息,下游消费者就可以消费了。
  5. 如果本地事务失败,Rollback 后,Broker 将删除这条半消息,下游永远感知不到。

  6. 事务回查:如果生产者因为网络超时、Full GC、或宕机,没有及时向 Broker 返回 Commit/Rollback,Broker 会定期(默认 6 秒起,指数退避)主动向生产者发起回查请求。 生产者在回查接口中,根据本地数据库的记录判断事务最终状态,然后返回 Commit 或 Rollback。

如何保证消息不丢失?

RocketMQ 通过多层防护确保消息从生产、存储到消费的可靠性:

  • 同步刷盘:Broker 将半消息同步写入磁盘(CommitLog)后才返回成功,避免内存丢失。

  • 主从复制:半消息会被同步到从节点(如果配置了 SYNC_MASTER),主宕机后从可接管。

  • 事务回查兜底:对于长时间未确认的半消息,Broker 会以指数退避策略反复回查,直到拿到明确结果。这就覆盖了“生产者发半消息后立刻宕机”这种极端情况。

  • 消费端确认:只有消费者主动返回消费成功,消息才会被标记为已消费,否则会重试。

示例代码:生产者发送事务消息(Spring Boot + RocketMQTemplate)

// 订单服务:发送事务消息
@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 事务消息的典型应用场景。我们要求:订单创建成功 → 库存一定被扣减;订单创建失败 → 库存一定不扣减。 绝不允许出现“订单有了,库存没扣”或“库存扣了,订单却失败”的数据不一致。

方案架构:

image.png

详细流程:

  1. 订单服务 发送一条半消息,内容包含订单 ID、商品 ID、数量等。

  2. 半消息发送成功后,订单服务在本地数据库事务中插入订单记录。

  3. 如果插入成功,向 Broker 提交(Commit)消息,Broker 将消息转为可投递,库存服务随后消费并扣减库存。

  4. 如果插入失败(例如商品超卖、数据库异常),向 Broker 回滚(Rollback)消息,Broker 删除半消息,库存服务永远收不到通知。

异常兜底——回查机制的作用:

假设订单服务在插入订单成功、但还没来得及向 Broker 发送 Commit 的瞬间,服务器宕机了。此时 Broker 发现这条半消息长时间未确认,会主动发起事务回查。订单服务重启后,回查接口通过查询数据库,发现订单已经存在,于是返回 Commit。消息最终还是会投递到库存服务,库存最终会被扣减。这就保证了最终一致性。

库存服务的幂等消费:

因为网络重试或回查可能导致消息重复投递,库存服务必须做好幂等——根据订单 ID 判断库存是否已经扣减过,如果已扣减则直接返回成功,避免重复扣库存。

// 库存服务消费端
@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());
    }
}

通过 半消息 → 本地事务 → Commit/Rollback → 回查 这套组合拳,再加上消费端的幂等,我们就在“下单”和“扣库存”这两个分布式操作之间,实现了一种可落地的最终一致性。它不追求强同步的完美,但确保在任何网络抖动和宕机下,系统状态最终都会收敛到正确的平衡点。

本地消息表轮询重试

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

(本地消息表的定义、与事务消息的对比)

本地消息表是一种基于数据库的分布式事务解决方案,它将业务操作和消息发送都放在同一个本地事务中,通过定时任务轮扫描未发送的消息进行重试,最终保证消息的可靠性投递。

与 RocketMQ 事务消息的区别:

  • 实现复杂度:本地消息表更简单,不需要依赖 MQ 的特殊功能

  • 性能:本地消息表需要额外的数据库读写,性能略低

  • 适用场景:本地消息表适用于任何 MQ(Kafka、RabbitMQ 等),事务消息仅适用于 RocketMQ

  • 一致性保证:两者都能保证最终一致性,但事务消息的实时性更好

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

⭐⭐(本地消息表设计、轮询重试机制、消息可靠性)

1️⃣ Common Answer 本地消息表就是建一张表存消息,订单插入的时候同时插入消息记录,用同一个事务。然后定时任务扫描状态为待发送的消息,发送成功后更新状态。如果发送失败就重试,直到成功为止。

2️⃣ Impressive Answer 本地消息表的核心思想是将业务操作和消息发送绑定在同一个本地事务中,通过异步轮询重试保证消息的最终可靠投递。我给你详细设计一下:

数据库表设计

-- 订单表
CREATE TABLE orders (
    id BIGINT PRIMARY KEY,
    user_id BIGINT,
    total_amount DECIMAL,
    status VARCHAR(20),
    create_time DATETIME
);

-- 本地消息表
CREATE TABLE local_message (
    id BIGINT AUTO_INCREMENT PRIMARY KEY,
    biz_type VARCHAR(50) NOT NULL,  -- 业务类型:ORDER_CREATED
    biz_id BIGINT NOT NULL,          -- 业务ID:订单ID
    topic VARCHAR(100) NOT NULL,     -- MQ Topic
    tag VARCHAR(50),                 -- MQ Tag
    message_body TEXT NOT NULL,      -- 消息内容
    status TINYINT NOT NULL DEFAULT 0, -- 0:待发送 1:发送中 2:发送成功 3:发送失败
    send_times INT DEFAULT 0,        -- 发送次数
    next_send_time DATETIME,         -- 下次发送时间
    create_time DATETIME,
    update_time DATETIME,
    UNIQUE KEY uk_biz (biz_type, biz_id),  -- 防止重复
    INDEX idx_status_time (status, next_send_time) -- 轮询索引
);

核心业务流程

核心代码实现

订单服务(插入订单和消息):

@Service
public class OrderService {

    @Transactional
    public void createOrder(OrderDTO orderDTO) {
        // 1. 插入订单
        Order order = new Order();
        order.setUserId(orderDTO.getUserId());
        order.setTotalAmount(orderDTO.getTotalAmount());
        order.setStatus("SUCCESS");
        orderMapper.insert(order);

        // 2. 插入本地消息(同一事务)
        LocalMessage message = new LocalMessage();
        message.setBizType("ORDER_CREATED");
        message.setBizId(order.getId());
        message.setTopic("order_created");
        message.setTag("stock");
        message.setMessageBody(JSON.toJSONString(order));
        message.setStatus(0); // 待发送
        message.setSendTimes(0);
        message.setNextSendTime(new Date());
        localMessageMapper.insert(message);
    }
}

定时任务(轮询重试):

@Component
public class MessageSendScheduler {

    // 每5秒执行一次
    @Scheduled(fixedDelay = 5000)
    public void sendPendingMessages() {
        // 1. 查询待发送的消息(分页,避免一次查太多)
        int pageSize = 100;
        List<LocalMessage> messages = localMessageMapper.selectPendingMessages(pageSize);

        if (CollectionUtils.isEmpty(messages)) {
            return;
        }

        // 2. 批量发送消息
        for (LocalMessage message : messages) {
            try {
                // 2.1 更新状态为发送中(防止重复消费)
                localMessageMapper.updateStatusToSendIng(message.getId());

                // 2.2 发送消息到 MQ
                SendResult result = rocketMQTemplate.syncSend(
                    message.getTopic() + ":" + message.getTag(),
                    message.getMessageBody()
                );

                if (result.getSendStatus() == SendStatus.SEND_OK) {
                    // 2.3 发送成功,更新状态
                    localMessageMapper.updateStatusToSuccess(message.getId());
                    log.info("消息发送成功,id={}", message.getId());
                } else {
                    throw new RuntimeException("消息发送失败");
                }

            } catch (Exception e) {
                log.error("消息发送失败,id={}", message.getId(), e);

                // 2.4 发送失败,更新状态和下次发送时间(指数退避)
                int sendTimes = message.getSendTimes() + 1;
                long delaySeconds = (long) Math.pow(2, sendTimes); // 2s, 4s, 8s...
                Date nextSendTime = new Date(System.currentTimeMillis() + delaySeconds * 1000);

                localMessageMapper.updateStatusToPending(
                    message.getId(),
                    sendTimes,
                    nextSendTime
                );

                // 2.5 如果重试次数超过阈值,标记为失败,人工介入
                if (sendTimes >= 10) {
                    localMessageMapper.updateStatusToFailed(message.getId());
                    // 发送告警通知
                    alertService.sendAlert("消息发送失败超过阈值,id=" + message.getId());
                }
            }
        }
    }
}

Mapper SQL(关键查询):

-- 查询待发送的消息(使用索引,避免全表扫描)
SELECT * FROM local_message
WHERE status = 0
  AND next_send_time <= NOW()
ORDER BY create_time ASC
LIMIT #{pageSize};

-- 更新状态为发送中(乐观锁,防止并发重复发送)
UPDATE local_message
SET status = 1,
    update_time = NOW()
WHERE id = #{id} AND status = 0;

-- 更新状态为成功
UPDATE local_message
SET status = 2,
    update_time = NOW()
WHERE id = #{id};

-- 更新状态为待发送(重试)
UPDATE local_message
SET status = 0,
    send_times = #{sendTimes},
    next_send_time = #{nextSendTime},
    update_time = NOW()
WHERE id = #{id};

消息不丢失的保障

  1. 本地事务保证:订单插入和消息插入在同一个事务中,要么都成功,要么都失败

  2. 轮询重试机制:定时任务持续扫描待发送消息,直到成功

  3. 指数退避策略:避免短时间内频繁重试,减少系统压力

  4. 失败告警机制:重试次数超过阈值时人工介入

  5. 消息去重:消费者端需要实现幂等性(如用 biz_type + biz_id 做唯一索引)

与 RocketMQ 事务消息的对比

维度 本地消息表 RocketMQ 事务消息
实现复杂度 简单,不依赖 MQ 特殊功能 复杂,需要实现事务监听器
性能 需要额外的数据库读写,性能略低 性能更好,无需额外数据库操作
MQ 适配性 适用于所有 MQ(Kafka、RabbitMQ 等) 仅适用于 RocketMQ
实时性 依赖定时任务轮询,有延迟 实时性更好
一致性保证 最终一致性 最终一致性
适用场景 跨 MQ 场景、已有 MQ 不支持事务 RocketMQ 环境、对实时性要求高

3️⃣ Key Differences

维度 Common Answer Impressive Answer
技术深度 简单说轮询重试 详细设计表结构、定时任务、指数退避
实践经验 缺乏生产环境考虑 考虑了并发控制、性能优化、告警机制
思考维度 仅关注功能实现 关注可靠性、可扩展性、与事务消息对比
表达方式 口头描述 用时序图、代码示例、SQL、对比表格
面试官印象 了解基本概念 有完整的架构设计能力

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

⭐⭐⭐(本地消息表应用、多下游通知、失败处理)

1️⃣ Common Answer 用本地消息表,支付成功后插入多条消息记录,然后定时任务发送到不同的 MQ topic。如果某个下游失败了就重试,直到成功。

2️⃣ Impressive Answer 这是一个典型的一对多通知场景,我会用本地消息表 + 分发策略来保证所有下游都能收到通知。让我详细设计一下:

整体架构设计

数据库表设计

-- 支付表
CREATE TABLE payment (
    id BIGINT PRIMARY KEY,
    order_id BIGINT,
    user_id BIGINT,
    amount DECIMAL,
    status VARCHAR(20),
    create_time DATETIME
);

-- 本地消息表(支持多下游)
CREATE TABLE local_message (
    id BIGINT AUTO_INCREMENT PRIMARY KEY,
    biz_type VARCHAR(50) NOT NULL,      -- 业务类型:PAYMENT_SUCCESS
    biz_id BIGINT NOT NULL,             -- 业务ID:支付ID
    target_system VARCHAR(50) NOT NULL, -- 目标系统:POINT, COUPON, RISK, ORDER
    topic VARCHAR(100) NOT NULL,        -- MQ Topic
    tag VARCHAR(50),
    message_body TEXT NOT NULL,
    status TINYINT NOT NULL DEFAULT 0,  -- 0:待发送 1:发送中 2:发送成功 3:发送失败
    send_times INT DEFAULT 0,
    next_send_time DATETIME,
    create_time DATETIME,
    update_time DATETIME,
    UNIQUE KEY uk_biz_target (biz_type, biz_id, target_system),
    INDEX idx_status_time (status, next_send_time)
);

核心代码实现

支付服务(插入支付和多条消息):

@Service
public class PaymentService {

    @Transactional
    public void processPaymentSuccess(PaymentDTO paymentDTO) {
        // 1. 插入支付记录
        Payment payment = new Payment();
        payment.setOrderId(paymentDTO.getOrderId());
        payment.setUserId(paymentDTO.getUserId());
        payment.setAmount(paymentDTO.getAmount());
        payment.setStatus("SUCCESS");
        paymentMapper.insert(payment);

        Long paymentId = payment.getId();

        // 2. 插入多条本地消息(同一事务)
        List<LocalMessage> messages = buildMessages(paymentId, paymentDTO);
        for (LocalMessage message : messages) {
            localMessageMapper.insert(message);
        }
    }

    private List<LocalMessage> buildMessages(Long paymentId, PaymentDTO paymentDTO) {
        List<LocalMessage> messages = new ArrayList<>();

        // 积分消息
        messages.add(createMessage(paymentId, "POINT", "point_topic", "add",
            buildPointMessage(paymentDTO)));

        // 优惠券消息
        messages.add(createMessage(paymentId, "COUPON", "coupon_topic", "use",
            buildCouponMessage(paymentDTO)));

        // 风控消息
        messages.add(createMessage(paymentId, "RISK", "risk_topic", "check",
            buildRiskMessage(paymentDTO)));

        // 订单消息
        messages.add(createMessage(paymentId, "ORDER", "order_topic", "update",
            buildOrderMessage(paymentDTO)));

        return messages;
    }

    private LocalMessage createMessage(Long bizId, String targetSystem,
                                      String topic, String tag, String body) {
        LocalMessage message = new LocalMessage();
        message.setBizType("PAYMENT_SUCCESS");
        message.setBizId(bizId);
        message.setTargetSystem(targetSystem);
        message.setTopic(topic);
        message.setTag(tag);
        message.setMessageBody(body);
        message.setStatus(0);
        message.setSendTimes(0);
        message.setNextSendTime(new Date());
        return message;
    }
}

定时任务(按目标系统分组批量发送):

@Component
public class MessageSendScheduler {

    @Scheduled(fixedDelay = 5000)
    public void sendPendingMessages() {
        // 1. 按目标系统分组查询,避免跨系统影响
        List<String> targetSystems = Arrays.asList("POINT", "COUPON", "RISK", "ORDER");

        for (String targetSystem : targetSystems) {
            sendMessagesByTargetSystem(targetSystem);
        }
    }

    private void sendMessagesByTargetSystem(String targetSystem) {
        int pageSize = 50;
        List<LocalMessage> messages = localMessageMapper.selectPendingMessages(
            targetSystem, pageSize);

        if (CollectionUtils.isEmpty(messages)) {
            return;
        }

        for (LocalMessage message : messages) {
            try {
                // 更新状态为发送中
                localMessageMapper.updateStatusToSendIng(message.getId());

                // 发送消息
                SendResult result = rocketMQTemplate.syncSend(
                    message.getTopic() + ":" + message.getTag(),
                    message.getMessageBody()
                );

                if (result.getSendStatus() == SendStatus.SEND_OK) {
                    localMessageMapper.updateStatusToSuccess(message.getId());
                    log.info("消息发送成功,system={}, id={}", targetSystem, message.getId());
                } else {
                    throw new RuntimeException("消息发送失败");
                }

            } catch (Exception e) {
                log.error("消息发送失败,system={}, id={}", targetSystem, message.getId(), e);

                // 失败重试
                int sendTimes = message.getSendTimes() + 1;
                long delaySeconds = calculateDelaySeconds(targetSystem, sendTimes);
                Date nextSendTime = new Date(System.currentTimeMillis() + delaySeconds * 1000);

                localMessageMapper.updateStatusToPending(
                    message.getId(), sendTimes, nextSendTime);

                if (sendTimes >= getMaxRetryTimes(targetSystem)) {
                    localMessageMapper.updateStatusToFailed(message.getId());
                    alertService.sendAlert(String.format(
                        "消息发送失败超过阈值,system=%s, id=%d",
                        targetSystem, message.getId()));
                }
            }
        }
    }

    // 不同系统配置不同的重试策略
    private long calculateDelaySeconds(String targetSystem, int sendTimes) {
        switch (targetSystem) {
            case "POINT":
                return (long) Math.pow(2, sendTimes); // 2s, 4s, 8s...
            case "COUPON":
                return sendTimes * 5; // 5s, 10s, 15s...
            case "RISK":
                return 1; // 风控需要实时,快速重试
            case "ORDER":
                return (long) Math.pow(2, sendTimes);
            default:
                return (long) Math.pow(2, sendTimes);
        }
    }

    private int getMaxRetryTimes(String targetSystem) {
        switch (targetSystem) {
            case "RISK":
                return 20; // 风控允许更多重试
            default:
                return 10;
        }
    }
}

异常场景处理

异常场景 处理方案
某个下游 MQ 不可用 其他下游正常发送,失败的按系统独立重试
某个下游消费失败 该下游消息持续重试,不影响其他下游
所有下游都失败 按系统独立重试,互不影响
某个下游重试次数超限 标记为失败,发送告警,人工介入
消息重复投递 下游服务实现幂等性(用 biz_type + biz_id + target_system 做唯一索引)

监控告警指标

  1. 待发送消息堆积量:按目标系统分组监控

  2. 发送失败率:按目标系统统计

  3. 平均发送耗时:监控定时任务执行时间

  4. 死信消息数:重试超限的消息数量

降级策略 如果本地消息表出现严重堆积(如超过 10 万条),可以:

  1. 增加定时任务线程数:并行处理不同系统的消息

  2. 增加批量大小:从 50 条增加到 200 条

  3. 暂停非核心系统:暂停优惠券、风控等非核心系统的消息发送

  4. 人工干预:对堆积消息进行批量处理

3️⃣ Key Differences

维度 Common Answer Impressive Answer
技术深度 简单说插入多条消息 详细设计多下游分发、按系统分组、差异化重试策略
实践经验 缺乏异常场景考虑 考虑了系统隔离、监控告警、降级策略
思考维度 仅关注功能实现 关注高可用、可扩展性、生产运维
表达方式 口头描述 用架构图、代码示例、异常处理表格
面试官印象 了解基本概念 有完整的系统设计能力

消息幂等性

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

幂等性 原本是数学概念:f(f(x)) = f(x)。在消息队列中,它意味着同一条消息被消费一次或多次,最终的业务结果完全相同。换句话说,即使消息被重复投递,系统也不会产生副作用,不会重复扣款、重复发积分。

为什么需要?

因为消息队列本身只能保证“至少一次投递”,无法保证“恰好一次”。在以下场景中,消息必然会出现重复:

场景一:生产者重复发送
  网络超时 → 生产者重试 → 同一条订单消息发送了两次

场景二:消费者重复消费
  消费者处理完消息,在提交 Offset 前宕机 → 重启后从上次 Offset 重新消费,导致消息重复处理

场景三:Rebalance 重平衡
  Kafka 消费者组重平衡 → 分区分配给新消费者,新消费者可能从旧 Offset 开始消费,重复消息

如果你不做幂等处理,积分系统可能会给用户加两次积分,库存系统可能会重复扣库存。幂等性就是让这些重复的消息“失效”,让系统在不可靠的网络和重试机制下,依然保持数据正确。


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

实现幂等性的核心思路是:记录每条消息的“已处理”状态,在处理前检查。根据记录介质和去重方式的不同,常见方案如下:

image.png

方案一:数据库唯一索引(最可靠)

原理:在数据库中建一张“消息去重表”,将消息 ID(或业务唯一键)作为唯一索引。处理消息时,先向这张表插入记录,如果插入成功说明是首次处理,执行业务逻辑;如果插入失败(唯一键冲突),说明消息已处理过,直接忽略。

示例代码:

-- 消息去重表
CREATE TABLE msg_dedup (
    msg_id VARCHAR(64) PRIMARY KEY,
    create_time DATETIME DEFAULT CURRENT_TIMESTAMP
) ENGINE=InnoDB;
@Transactional
public void processMessage(String msgId, OrderEvent event) {
    try {
        // 1. 尝试插入去重记录,利用主键冲突
        msgDedupMapper.insert(msgId);
    } catch (DuplicateKeyException e) {
        // 已处理过,直接返回
        return;
    }
    // 2. 执行业务逻辑:增加积分
    pointsService.addPoints(event.getUserId(), event.getAmount());
}

优点:天然可靠,依赖数据库的 ACID 保证,无需额外组件。

缺点:增加数据库写压力;如果业务操作和去重记录不在同一个数据库,需要跨库事务。


方案二:Redis + 业务 Token(高性能)

原理:每条消息生成一个唯一的 Token(通常是 msgId业务Id+版本号),写入 Redis,并设置过期时间(如 24 小时)。处理前用 SETNX 命令尝试写入,成功则处理,失败则跳过。

示例代码:

public void processMessage(String msgId, OrderEvent event) {
    String key = "msg:dedup:" + msgId;
    // SETNX: 如果 key 不存在则设置并返回 true,否则 false
    Boolean isFirst = redisTemplate.opsForValue().setIfAbsent(key, "1", Duration.ofHours(24));
    if (Boolean.FALSE.equals(isFirst)) {
        return; // 已处理
    }
    pointsService.addPoints(event.getUserId(), event.getAmount());
}

优点:高性能,Redis 读写极快,适合高并发;自动过期,节省空间。

缺点:依赖 Redis 的高可用;Redis 数据丢失可能导致重复消费(需要评估风险)。


方案三:业务状态机 + 乐观锁(最轻量)

原理:利用业务本身的状态字段(如订单状态从“已支付”变成“已发放积分”),通过更新状态的条件(乐观锁)来保证只执行一次。

示例代码:

// 通过更新积分发放状态来防重,该操作本身也是业务的一部分
int rows = db.execute(
    "UPDATE orders SET points_status='GRANTED' WHERE order_id=? AND points_status='PENDING'",
    event.getOrderId()
);
if (rows > 0) {
    pointsService.addPoints(event.getUserId(), event.getAmount());
}

优点:零额外存储,性能最好,复用业务状态。

缺点:要求业务本身有明确的状态流转;如果业务操作本身不具备状态字段(如“发送通知”),则不适用。


方案对比总结:

方案 可靠性 性能 额外存储 适用场景
数据库唯一索引 极高 需要去重表 对一致性要求极高的核心业务
Redis + Token 极高 需要 Redis 高并发、允许少量重复风险的场景
业务状态机 极高 极高 业务有明确状态机、字段能承载去重的场景

推荐组合:核心资金类用“数据库唯一索引 + 业务状态机”双保险;非核心通知类用 Redis Token;同时所有消费者都应实现幂等,形成纵深防御。


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

这是一个典型的“支付成功 → 发放积分”的幂等场景。我们采用 数据库唯一索引 + 业务状态机 的双重保证,确保积分只增加一次。

整体流程:

image.png

步骤详解:

  1. 流水检查:每次收到支付成功消息,先根据订单 ID 查询积分流水表。如果发现已有发放记录,直接返回成功。这层是轻量级的查重。

  2. 唯一索引拦截:插入积分流水记录时,以 order_id 作为唯一索引。如果并发或重试导致同一条订单消息被两个线程同时处理,其中一个线程的 INSERT 会因为主键冲突而失败,从而阻止重复处理。

  3. 余额更新:在插入流水成功的同一事务中,更新用户的积分余额(SET points = points + amount),并更新订单表上的积分发放状态为“已发放”。所有操作在一个数据库事务中,保证原子性。

代码示例(Spring Boot + MyBatis):

@Service
public class PointsGrantService {

    @Autowired
    private PointsFlowMapper flowMapper;
    @Autowired
    private UserPointsMapper userPointsMapper;
    @Autowired
    private OrderMapper orderMapper;

    @Transactional(rollbackFor = Exception.class)
    public void grantPoints(String orderId, Long userId, int points) {
        // 1. 快速检查:查询流水表
        PointsFlow existing = flowMapper.selectByOrderId(orderId);
        if (existing != null) {
            return; // 已发放
        }

        // 2. 插入流水,order_id 为唯一键,如果重复则抛出 DuplicateKeyException
        PointsFlow flow = PointsFlow.builder()
                .orderId(orderId)
                .userId(userId)
                .points(points)
                .build();
        try {
            flowMapper.insert(flow);
        } catch (DuplicateKeyException e) {
            return; // 并发写入,重复处理
        }

        // 3. 增加用户积分(原子更新)
        userPointsMapper.addPoints(userId, points);

        // 4. 更新订单表积分发放状态,同时做状态机防重
        orderMapper.updatePointsStatus(orderId, "GRANTED");
    }
}
-- 积分流水表 DDL
CREATE TABLE points_flow (
    id BIGINT AUTO_INCREMENT PRIMARY KEY,
    order_id VARCHAR(64) NOT NULL,
    user_id BIGINT NOT NULL,
    points INT NOT NULL,
    create_time DATETIME DEFAULT CURRENT_TIMESTAMP,
    UNIQUE KEY uk_order_id (order_id)
);

关键设计点:

  • 唯一索引保护:即使没有第一步的查询,第二步的 INSERT 也会因为唯一键冲突而阻止重复。查询是一种优化,避免无效的插入异常。

  • 事务边界:流水插入和余额更新在同一个本地事务中,要么都成功,要么都回滚。如果余额更新失败(如用户不存在),流水也会回滚,不会产生脏数据。

  • 外部兜底:如果业务允许极其罕见的 Redis/DB 故障,可以添加一个定时对账任务,每天对比支付单和积分流水,发现差异自动补发或告警。

收束:保证积分不重复增加的秘密,不在于某一种技术,而在于让重复消息在数据库的强约束面前变成“无操作”。唯一索引是守门员,流水表是检查站,而事务则确保每一步都走得稳当。这种防御思路,同样适用于库存、优惠券、余额等任何需要精确计数的场景。


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 格式的场景


三种格式对比

维度 STATEMENT ROW MIXED
数据一致性 ⭐⭐ ⭐⭐⭐⭐⭐ ⭐⭐⭐⭐
空间占用 ⭐⭐⭐⭐⭐ ⭐⭐ ⭐⭐⭐⭐
网络传输 ⭐⭐⭐⭐⭐ ⭐⭐ ⭐⭐⭐⭐
可读性 ⭐⭐⭐⭐⭐ ⭐⭐ ⭐⭐⭐⭐
恢复精确度 ⭐⭐ ⭐⭐⭐⭐⭐ ⭐⭐⭐⭐
配置复杂度 ⭐⭐⭐⭐⭐ ⭐⭐⭐⭐⭐ ⭐⭐⭐⭐⭐

生产环境配置建议

-- 查看当前 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

维度 Common Answer Impressive Answer
技术深度 简单说三种格式的区别 详细对比优缺点,给出配置建议
实践经验 缺乏场景选择经验 能根据不同场景选择合适的格式
思考维度 仅关注概念 关注数据一致性、性能、空间占用的权衡
表达方式 口头描述 用代码示例、对比表格、配置命令
面试官印象 了解基本概念 有生产环境配置经验

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

异常场景处理

异常场景 处理方案
ES 同步失败 消息重试,直到成功
ES 不可用 消息堆积在 MQ,ES 恢复后自动消费
Binlog 丢失 Canal 记录消费位点,可以从上次位点继续消费
数据不一致 定时任务全量比对 DB 和 ES 数据,修复不一致数据

监控告警指标

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

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

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

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

3️⃣ Key Differences

维度 Common Answer Impressive Answer
技术深度 简单说用 ROW 格式 详细分析 STATEMENT 和 MIXED 的问题,说明 ROW 的优势
实践经验 缺乏生产环境考虑 考虑了配置、异常处理、监控告警
思考维度 仅关注格式选择 关注数据一致性、生产实践、异常处理
表达方式 口头描述 用代码示例、配置命令、异常处理表格
面试官印象 了解基本概念 有完整的 CDC 实战经验

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 整体架构

核心组件详解

  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 消费失败 不发送 ACK,下次继续消费
Canal 重启 从位点文件恢复消费位点
MySQL Master 切换 Canal 自动重新连接新的 Master
消息队列不可用 Canal 暂停消费,MQ 恢复后继续
Binlog 丢失 Canal 记录位点,可以从上次位点继续

生产环境配置建议

# 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

维度 Common Answer Impressive Answer
技术深度 简单说模拟 Slave 详细说明 EventParser、EventSink、EventStore、EventAck
实践经验 缺乏生产环境考虑 考虑了位点持久化、HA 高可用、异常处理
思考维度 仅关注基本原理 关注数据可靠性、性能优化、监控告警
表达方式 口头描述 用架构图、代码示例、配置文件、异常处理表格
面试官印象 了解基本概念 有完整的 Canal 实战经验

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

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

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

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

整体架构设计

核心实现步骤

步骤 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...
}

异常场景处理

异常场景 处理方案
ES 不可用 消息堆积在 MQ,ES 恢复后自动消费
ES 同步失败 消息重试(默认 16 次),重试失败进入死信队列
Canal 消费延迟 监控告警, Canal 自动从上次位点继续消费
数据不一致 定时任务全量比对 DB 和 ES 数据,修复不一致数据
消息重复消费 ES 使用 id 作为文档 ID,天然支持幂等性

监控告警指标

  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

维度 Common Answer Impressive Answer
技术深度 简单说用 Canal + MQ 详细设计 MySQL 配置、Canal 配置、消费者代码
实践经验 缺乏生产环境考虑 考虑了异常处理、监控告警、性能优化
思维维度 仅关注功能实现 关注数据一致性、可靠性、可维护性
表达方式 口头描述 用架构图、代码示例、配置文件、异常处理表格
面试官印象 了解基本概念 有完整的 CDC 实战经验

4、容易一起考的题

关联题 和本题的关系
MySQL Binlog 三种格式 Canal 基于 Binlog 实现,格式选择影响 Canal 功能
如何保证数据一致性 Canal 通过位点记录保证数据一致性
Elasticsearch 的数据同步方案 Canal + MQ 是常用的 ES 同步方案
分布式系统的数据同步 Canal 是 CDC 技术的典型实现

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 和缓存。让我详细讲一下:

整体架构设计

核心实现步骤

步骤 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);
        }
    }
}

异常场景处理

异常场景 处理方案
ES 不可用 消息堆积在 MQ,ES 恢复后自动消费
缓存不可用 消息堆积在 MQ,缓存恢复后自动消费
ES 同步失败 消息重试(默认 16 次),重试失败进入死信队列
缓存同步失败 消息重试(默认 16 次),重试失败进入死信队列
消息重复消费 ES 和缓存使用唯一 ID,天然支持幂等性
数据不一致 定时任务校验并修复不一致数据

性能优化建议

  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

维度 Common Answer Impressive Answer
技术深度 简单说用 Canal + MQ 详细设计多消费者架构、一致性保证机制
实践经验 缺乏异常场景考虑 考虑了死信队列、数据一致性校验、性能优化
思维维度 仅关注功能实现 关注数据一致性、可靠性、性能优化
表达方式 口头描述 用架构图、代码示例、异常处理表格、监控指标
面试官印象 了解基本概念 有完整的 CDC 实战经验

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

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

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

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

场景需求分析

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

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

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

整体架构设计

核心实现步骤

步骤 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;
    }
}

为什么需要延迟双删?

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

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

  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 并回写了缓存,延迟双删会再次删除缓存,确保下次查询时获取最新数据。

异常场景处理

异常场景 处理方案
Canal 消费延迟 延迟双删保证缓存最终一致性
缓存删除失败 消息重试,直到成功
消息重复消费 缓存删除是幂等操作,重复删除无影响
缓存和 DB 不一致 定时任务校验并修复不一致数据
并发读写 延迟双删 + Cache Aside 模式保证一致性

监控告警指标

  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

维度 Common Answer Impressive Answer
技术深度 简单说用 Canal 更新缓存 详细设计延迟双删方案,说明为什么需要延迟双删
实践经验 缺乏并发场景考虑 考虑了并发读写、消息延迟、异常处理
思维维度 仅关注功能实现 关注数据一致性、并发安全、性能优化
表达方式 口头描述 用时序图、代码示例、异常处理表格
面试官印象 了解基本概念 有完整的缓存一致性实战经验