前言:从“泛消费”到“精准控”的进阶需求​

在上一篇中,我们掌握了 RocketMQ 事务消息的原理与落地实践,解决了分布式系统中“业务与消息”的原子性问题,确保了数据最终一致性。但在实际生产场景中,仅实现消息的可靠发送与消费还不够——随着业务规模扩大,Topic 中会积累大量不同类型、不同场景的消息,若消费者盲目消费所有消息,不仅会增加系统开销,还会导致业务逻辑冗余、处理效率低下。

例如:电商场景中,一个“order_topic”可能包含“订单创建”“订单支付”“订单取消”“订单退款”等多种消息,物流服务只需消费“订单创建”消息,支付服务只需消费“订单支付”“订单退款”消息,若不进行过滤,所有服务都会收到所有类型的消息,再手动筛选,会严重影响系统性能。​

同时,生产环境中难免出现异常:消费者消费消息时宕机、业务逻辑bug导致消息处理失败、误删除重要消息等,此时需要能够“回溯历史消息”,重新消费指定时间段或指定条件的消息,避免数据丢失或业务中断。​

为解决上述问题,RocketMQ 提供了消息过滤和消息回溯两大核心特性,前者实现“精准消费”,后者实现“历史消息找回”。本篇将从原理、特性、实操三个维度,深度解析这两个特性,结合实际业务场景,手把手教你落地精准消费与消息回溯方案,进一步提升 RocketMQ 在生产环境中的灵活性和可用性。

前置要求:​
已掌握 RocketMQ 基础 API 开发(生产者、消费者实现),了解 Topic、Tag、Group 核心概念​
已部署 RocketMQ 高可用集群(Dledger 或多主多从模式均可)​
已掌握 RocketMQ 事务消息的核心原理与实操,具备简单的 Java 开发能力

一、消息过滤:实现“精准消费”的核心原理

RocketMQ 的消息过滤,本质是“在消息发送时标记特征,在消息消费时根据特征筛选”,核心目标是让消费者只消费自己关心的消息,减少无效消费,提升系统吞吐量。RocketMQ 提供了两种核心过滤方式:Tag 过滤(最常用)和 SQL92 过滤(高级场景),同时支持自定义过滤(灵活扩展)。

1.1 核心概念:过滤的本质的是“特征匹配”
消息过滤的核心逻辑的是“生产者给消息打标签/加属性,消费者根据标签/属性设置筛选条件,Broker 或消费者根据条件筛选消息”。根据筛选发生的位置,可分为两种模式:

  • Broker 端过滤:消息发送到 Broker 后,Broker 先根据消费者的筛选条件,筛选出符合条件的消息,再推送给消费者(或供消费者拉取)。优点是减少无效消息的网络传输,降低消费者压力;缺点是增加 Broker 计算开销。

  • 消费者端过滤:Broker 将所有消息推送给消费者,消费者在本地根据筛选条件,过滤掉不需要的消息。优点是减轻 Broker 压力;缺点是无效消息会占用网络带宽,增加消费者处理开销。

RocketMQ 中,Tag 过滤支持 Broker 端过滤(默认),SQL92 过滤仅支持 Broker 端过滤,自定义过滤支持消费者端过滤,开发者可根据业务场景选择合适的方式。

1.2 两种核心过滤方式详解(Tag + SQL92)

1.2.1 Tag 过滤(最常用,推荐优先使用)

Tag(标签)是 RocketMQ 最基础、最常用的过滤方式,本质是给消息设置“分类标签”,消费者通过订阅指定 Tag,实现对消息的精准筛选。Tag 过滤的核心特点是“轻量、高效、易用”,适用于大多数简单过滤场景。

核心规则

  • 一个消息只能设置一个 Tag(Tag 是字符串类型,长度不超过 64 个字符,建议命名规范,如“order:create”“order:pay”)。

  • 消费者订阅时,可指定单个 Tag(如“order:create”)、多个 Tag(用“||”分隔,如“order:create||order:pay”),或订阅所有 Tag(用“*”表示)。

  • Tag 过滤发生在 Broker 端,Broker 会根据消费者订阅的 Tag,筛选出符合条件的消息,避免无效消息传输。

适用场景
消息分类明确、筛选条件简单的场景,如:订单 Topic 按“创建、支付、取消、退款”分类,物流服务订阅“order:create”,支付服务订阅“order:pay||order:refund”。

1.2.2 SQL92 过滤(高级场景,灵活筛选)

Tag 过滤仅能根据单个标签筛选,无法满足复杂筛选需求(如根据消息属性范围、多条件组合筛选)。此时可使用 SQL92 过滤,通过 SQL 语句的条件表达式,筛选出符合条件的消息,支持多条件组合、范围筛选等。

核心规则

  • SQL92 过滤基于消息的用户属性(UserProperty),生产者发送消息时,给消息设置自定义属性(如“price”“userId”“createTime”),消费者通过 SQL 语句筛选这些属性。

  • 支持的 SQL 语法:=、!=、>、<、>=、<=、IN、NOT IN、LIKE、AND、OR 等(不支持复杂 SQL,如 JOIN、GROUP BY)。

  • SQL92 过滤仅支持 Broker 端过滤,且需要在 Broker 端开启 SQL 过滤功能(默认关闭)。

适用场景
筛选条件复杂的场景,如:电商优惠券消息,筛选“面值>10元”且“有效期>7天”的消息;订单消息,筛选“用户ID为1001、1002”且“订单金额>1000元”的消息。

1.3 两种过滤方式对比

过滤方式 核心特点 适用场景 筛选位置
Tag 过滤 轻量、高效、易用,一个消息一个 Tag 消息分类明确、筛选条件简单 Broker 端(默认)
SQL92 过滤 灵活,支持多条件、范围筛选,基于消息属性 筛选条件复杂,需多属性组合 Broker 端(需开启)

核心建议
优先使用 Tag 过滤,满足大多数业务场景,兼顾性能和易用性;若 Tag 过滤无法满足需求(如多条件筛选),再使用 SQL92 过滤,注意开启 Broker 端配置并控制筛选复杂度,避免影响 Broker 性能。

二、消息回溯:找回历史消息的核心原理与场景

消息回溯,指消费者能够重新消费“已经消费过”或“未消费但错过”的历史消息,核心解决生产环境中“消息处理失败、误删除、系统宕机”等异常场景下的数据恢复问题。RocketMQ 基于消息的存储机制(CommitLog + ConsumeQueue),支持多种回溯方式,满足不同场景的需求。

2.1 消息回溯的核心前提

RocketMQ 的消息回溯依赖于消息的持久化存储——Broker 会将消息持久化到 CommitLog 中,即使消息被消费者消费,只要未超过存储时间(默认 72 小时,可配置),就可以通过回溯机制重新消费。

关键注意点:

  • 消息回溯的前提是消息仍在 Broker 的存储时间范围内(超过存储时间,消息会被自动清理,无法回溯)。

  • 消息回溯会导致消费者重复消费消息,因此必须保证消费者的幂等性(与事务消息的幂等性要求一致)。

2.2 三种核心消息回溯方式
2.2.1 按时间回溯(最常用)

按时间回溯是最常用的回溯方式,消费者可以指定一个具体的时间点(如“2026-03-28 10:00:00”),重新消费该时间点之后的所有消息(或该时间点之前的历史消息)。

核心原理:RocketMQ 的 ConsumeQueue 中记录了消息的存储时间戳,Broker 可根据时间戳快速定位到对应的消息,消费者从该位置开始重新消费。

适用场景:消费者宕机、业务逻辑 bug 导致某一时间段的消息处理失败,需要重新消费该时间段的所有消息(如:上午 10 点到 11 点的订单消息未处理成功,回溯到 10 点重新消费)。

2.2.2 按偏移量回溯(精准回溯)
每个消费者组在消费消息时,会记录自己的消费偏移量(Offset)——即当前消费到的消息在 ConsumeQueue 中的位置。按偏移量回溯,就是让消费者从指定的偏移量位置开始重新消费。

核心原理:偏移量是消息在 ConsumeQueue 中的唯一标识,每个消息对应一个唯一的偏移量,消费者可通过指定偏移量,精准定位到某一条消息,从该消息开始重新消费。

适用场景:精准定位某一条或某一段消息的消费异常,如:已知某条消息(偏移量为 1000)处理失败,可直接回溯到偏移量 1000,重新消费该消息及之后的消息。

2.2.3 按消息 ID 回溯(单个消息回溯)
按消息 ID 回溯,指通过消息的唯一 ID(Message ID),精准找到该消息,重新消费该消息。这种方式适用于单个消息处理失败的场景,无需回溯所有消息,效率更高。

核心原理:RocketMQ 的 CommitLog 中记录了消息 ID 与消息存储位置的映射关系,可通过消息 ID 快速查询到消息的存储位置,进而重新消费该消息。

适用场景:单个消息处理失败(如:某一条订单消息因参数异常处理失败),只需重新消费该条消息,无需影响其他消息的消费。

2.3 消息回溯的注意事项

  • 消息回溯会触发重复消费,必须保证消费者的幂等性(如:基于消息 ID 或业务 ID 做幂等校验),避免重复处理业务导致异常。

  • 按时间回溯时,时间点的精度为毫秒级,需注意时间格式的正确性(如“yyyy-MM-dd HH:mm:ss”)。

  • 偏移量是消费者组级别的,不同消费者组的偏移量相互独立,回溯某一个消费者组的偏移量,不会影响其他消费者组。

  • 消息回溯仅能回溯“未被清理”的消息,若消息已超过 Broker 的存储时间(默认 72 小时),无法回溯,需提前配置合适的存储时间(生产环境可根据业务需求调整为 7 天或更久)。

三、实操:消息过滤与消息回溯代码实现(Java 版)

本次实操基于 RocketMQ 4.8.0 版本,结合 SpringBoot 整合方式,延续上一篇的电商场景,实现“订单消息过滤”和“订单消息回溯”,分为 4 部分:Tag 过滤实操、SQL92 过滤实操、三种消息回溯实操、幂等性保障。

环境准备:沿用上一篇的 SpringBoot 项目、RocketMQ 集群、数据库(无需额外新增依赖和配置,仅需开启 SQL92 过滤功能)。

3.1 环境准备:开启 SQL92 过滤功能(Broker 端)
SQL92 过滤默认关闭,需修改 Broker 配置文件(broker.conf),添加如下配置,重启 Broker 生效:

3.2 消息过滤实操(Tag + SQL92)

3.2.1 Tag 过滤实操(生产者+消费者)

步骤 1:生产者发送消息(设置 Tag)

生产者发送订单消息时,给不同类型的订单消息设置不同的 Tag(如“order:create”“order:pay”“order:cancel”):

import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import org.springframework.messaging.Message;
import org.springframework.messaging.support.MessageBuilder;

import java.util.UUID;

@Component
public class OrderFilterProducer {

    @Autowired
    private RocketMQTemplate rocketMQTemplate;

    /**
     * 发送订单消息(带 Tag)
     * @param orderId 订单ID
     * @param tag 消息 Tag(如 order:create、order:pay)
     * @param orderAmount 订单金额(用于后续 SQL 过滤)
     */
    public void sendOrderMessageWithTag(String orderId, String tag, Double orderAmount) {
        // 构建消息体(订单ID)
        Message<String> message = MessageBuilder
                .withPayload(orderId)
                // 设置消息用户属性(用于 SQL 过滤)
                .setHeader("orderAmount", orderAmount)
                .setHeader("userId", "user001") // 模拟用户ID
                .build();

        // 发送消息:Topic:Tag 格式,指定 Tag
        rocketMQTemplate.syncSend("order_filter_topic:" + tag, message);
        System.out.println("发送订单消息成功,订单ID:" + orderId + ",Tag:" + tag);
    }
}

步骤 2:消费者订阅指定 Tag(精准过滤)

物流服务仅订阅“order:create” Tag 的消息,支付服务仅订阅“order:pay”和“order:cancel” Tag 的消息:

// 物流服务消费者(仅消费 order:create 消息)
@Component
@RocketMQMessageListener(
        topic = "order_filter_topic",
        selectorExpression = "order:create", // 订阅单个 Tag
        consumerGroup = "logistics-consumer-group"
)
public class LogisticsConsumer implements RocketMQListener<String> {

    @Override
    public void onMessage(String orderId) {
        System.out.println("物流服务收到订单创建消息,订单ID:" + orderId + ",开始安排物流");
        // 执行物流相关业务(省略)
    }
}

// 支付服务消费者(消费 order:pay 和 order:cancel 消息)
@Component
@RocketMQMessageListener(
        topic = "order_filter_topic",
        selectorExpression = "order:pay||order:cancel", // 订阅多个 Tag,用 || 分隔
        consumerGroup = "payment-consumer-group"
)
public class PaymentConsumer implements RocketMQListener<String> {

    @Override
    public void onMessage(String orderId) {
        System.out.println("支付服务收到消息,订单ID:" + orderId + ",处理支付/取消逻辑");
        // 执行支付相关业务(省略)
    }
}

3.2.2 SQL92 过滤实操(生产者+消费者)

步骤 1:生产者发送消息(设置用户属性)

沿用上面的生产者代码,发送消息时已设置“orderAmount”(订单金额)和“userId”(用户ID)两个用户属性,无需额外修改。

步骤 2:消费者使用 SQL 语句筛选消息
假设财务服务需要消费“订单金额>1000元”且“用户ID为user001”的订单消息,使用 SQL92 过滤实现:

// 财务服务消费者(SQL92 过滤)
@Component
@RocketMQMessageListener(
        topic = "order_filter_topic",
        selectorType = SelectorType.SQL92, // 指定过滤方式为 SQL92
        selectorExpression = "orderAmount > 1000 AND userId = 'user001'", // SQL 筛选条件
        consumerGroup = "finance-consumer-group"
)
public class FinanceConsumer implements RocketMQListener<String> {

    @Override
    public void onMessage(String orderId) {
        System.out.println("财务服务收到高价值订单消息,订单ID:" + orderId + ",开始对账");
        // 执行财务对账相关业务(省略)
    }
}

3.2.3 测试验证(过滤功能)
编写测试接口,发送不同 Tag、不同属性的消息,验证过滤功能:

@RestController
@RequestMapping("/order/filter")
public class OrderFilterController {

    @Autowired
    private OrderFilterProducer filterProducer;

    // 发送订单创建消息(Tag:order:create,金额 800 元)
    @PostMapping("/create")
    public String sendCreateOrder(@RequestParam String orderId) {
        filterProducer.sendOrderMessageWithTag(orderId, "order:create", 800.0);
        return "订单创建消息发送成功";
    }

    // 发送订单支付消息(Tag:order:pay,金额 1200 元)
    @PostMapping("/pay")
    public String sendPayOrder(@RequestParam String orderId) {
        filterProducer.sendOrderMessageWithTag(orderId, "order:pay", 1200.0);
        return "订单支付消息发送成功";
    }

    // 发送订单取消消息(Tag:order:cancel,金额 500 元)
    @PostMapping("/cancel")
    public String sendCancelOrder(@RequestParam String orderId) {
        filterProducer.sendOrderMessageWithTag(orderId, "order:cancel", 500.0);
        return "订单取消消息发送成功";
    }
}

预期结果:

  • 调用 /order/filter/create:仅物流服务收到消息。

  • 调用 /order/filter/pay:支付服务收到消息,财务服务(金额>1000)也收到消息。

  • 调用 /order/filter/cancel:仅支付服务收到消息。

3.3 消息回溯实操(三种方式)
本次实操以“物流服务消费者(logistics-consumer-group)”为例,实现三种消息回溯方式,核心是通过 RocketMQTemplate 或 RocketMQAdmin 操作消费偏移量。

3.3.1 按时间回溯(最常用)
实现逻辑:通过 RocketMQAdmin 工具,将消费者组的消费偏移量重置到指定时间点,消费者从该时间点开始重新消费。

import org.apache.rocketmq.client.admin.RocketMQAdmin;
import org.apache.rocketmq.client.admin.ResetOffsetRequest;
import org.apache.rocketmq.client.exception.MQClientException;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;

import java.util.Date;
import java.util.Set;

@RestController
@RequestMapping("/order/backtrack")
public class OrderBacktrackController {

    @Autowired
    private RocketMQTemplate rocketMQTemplate;

    /**
     * 按时间回溯消息
     * @param topic 主题
     * @param consumerGroup 消费者组
     * @param backtrackTime 回溯时间(格式:yyyy-MM-dd HH:mm:ss)
     */
    @PostMapping("/time")
    public String backtrackByTime(
            @RequestParam String topic,
            @RequestParam String consumerGroup,
            @RequestParam String backtrackTime
    ) throws Exception {
        // 1. 获取 RocketMQAdmin 实例
        RocketMQAdmin rocketMQAdmin = rocketMQTemplate.getRocketMQAdmin();

        // 2. 解析回溯时间(转换为时间戳,毫秒级)
        SimpleDateFormat sdf = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss");
        Date date = sdf.parse(backtrackTime);
        long timestamp = date.getTime();

        // 3. 构建重置偏移量请求
        ResetOffsetRequest request = new ResetOffsetRequest();
        request.setTopic(topic);
        request.setConsumerGroup(consumerGroup);
        request.setResetTimestamp(timestamp); // 回溯到指定时间点
        request.setForce(true); // 强制重置(即使消费者正在消费)

        // 4. 执行偏移量重置(回溯)
        rocketMQAdmin.resetOffset(request);

        return "按时间回溯成功,回溯时间:" + backtrackTime;
    }
}

3.3.2 按偏移量回溯(精准回溯)
实现逻辑:先查询消费者组当前的消费偏移量,再将偏移量重置到指定值,消费者从该偏移量开始重新消费。

/**
 * 按偏移量回溯消息
 * @param topic 主题
 * @param consumerGroup 消费者组
 * @param queueId 队列ID(RocketMQ 每个 Topic 默认4个队列)
 * @param targetOffset 目标偏移量
 */
@PostMapping("/offset")
public String backtrackByOffset(
        @RequestParam String topic,
        @RequestParam String consumerGroup,
        @RequestParam int queueId,
        @RequestParam long targetOffset
) throws Exception {
    RocketMQAdmin rocketMQAdmin = rocketMQTemplate.getRocketMQAdmin();

    // 构建重置偏移量请求
    ResetOffsetRequest request = new ResetOffsetRequest();
    request.setTopic(topic);
    request.setConsumerGroup(consumerGroup);
    request.setQueueId(queueId); // 指定队列(每个队列的偏移量独立)
    request.setTargetOffset(targetOffset); // 目标偏移量
    request.setForce(true);

    // 执行偏移量重置
    rocketMQAdmin.resetOffset(request);

    return "按偏移量回溯成功,队列ID:" + queueId + ",目标偏移量:" + targetOffset;
}

3.3.3 按消息 ID 回溯(单个消息)
实现逻辑:通过消息 ID 查询消息的存储位置(队列ID、偏移量),再按偏移量回溯到该消息,重新消费。

在这里插入代码片/**
 * 按消息ID回溯消息(单个消息)
 * @param msgId 消息ID
 * @param consumerGroup 消费者组
 */
@PostMapping("/msgId")
public String backtrackByMsgId(
        @RequestParam String msgId,
        @RequestParam String consumerGroup
) throws Exception {
    RocketMQAdmin rocketMQAdmin = rocketMQTemplate.getRocketMQAdmin();

    // 1. 根据消息ID查询消息详情(获取队列ID和偏移量)
    MessageExt messageExt = rocketMQAdmin.viewMessage(msgId);
    String topic = messageExt.getTopic();
    int queueId = messageExt.getQueueId();
    long offset = messageExt.getQueueOffset();

    // 2. 按偏移量回溯到该消息
    ResetOffsetRequest request = new ResetOffsetRequest();
    request.setTopic(topic);
    request.setConsumerGroup(consumerGroup);
    request.setQueueId(queueId);
    request.setTargetOffset(offset); // 消息对应的偏移量
    request.setForce(true);

    rocketMQAdmin.resetOffset(request);

    return "按消息ID回溯成功,消息ID:" + msgId + ",队列ID:" + queueId + ",偏移量:" + offset;
}

3.3.4 测试验证(回溯功能)

  1. 先发送几条订单创建消息,确保物流服务已消费。

  2. 调用 /order/backtrack/time,传入 topic=order_filter_topic、consumerGroup=logistics-consumer-group、backtrackTime=消息发送时间,观察物流服务是否重新消费该时间点后的消息。

  3. 调用 /order/backtrack/offset,传入队列ID、目标偏移量(如1),观察物流服务是否从偏移量1开始重新消费。

  4. 调用 /order/backtrack/msgId,传入某条消息的ID,观察物流服务是否重新消费该条消息。

3.4 幂等性保障(回溯必做)
消息回溯会导致重复消费,需给消费者添加幂等性校验,沿用之前的 Redis 幂等方案(实际开发中需引入 Redis 依赖):

// 物流服务消费者(添加幂等性校验)
@Component
@RocketMQMessageListener(
        topic = "order_filter_topic",
        selectorExpression = "order:create",
        consumerGroup = "logistics-consumer-group"
)
public class LogisticsConsumer implements RocketMQListener<String> {

    @Autowired
    private StringRedisTemplate redisTemplate;

    // 幂等性 key 前缀
    private static final String IDEMPOTENT_KEY_PREFIX = "order:logistics:processed:";

    @Override
    public void onMessage(String orderId) {
        // 1. 幂等性校验:查询该订单是否已处理
        String key = IDEMPOTENT_KEY_PREFIX + orderId;
        Boolean isProcessed = redisTemplate.hasKey(key);
        if (Boolean.TRUE.equals(isProcessed)) {
            System.out.println("订单" + orderId + "已处理,跳过重复消费");
            return;
        }

        // 2. 执行物流相关业务
        System.out.println("物流服务收到订单创建消息,订单ID:" + orderId + ",开始安排物流");

        // 3. 标记订单已处理(设置过期时间,避免 Redis 内存溢出)
        redisTemplate.opsForValue().set(key, "1", 24, TimeUnit.HOURS);
    }
}

四、生产环境落地注意事项(避坑指南)

4.1 消息过滤注意事项

  • Tag 命名规范:Tag 建议采用“业务模块:消息类型”的格式(如“order:create”“user:register”),避免杂乱无章,便于运维和排查。

  • SQL92 过滤性能:SQL 筛选条件不宜过于复杂(如避免使用多个 OR、LIKE 模糊查询),否则会增加 Broker 计算开销,影响集群性能;同时避免频繁使用 SQL92 过滤,优先使用 Tag 过滤。

  • 过滤方式选择:简单筛选用 Tag,复杂筛选用 SQL92,自定义过滤(如基于消息内容过滤)仅在特殊场景使用,且需在消费者端实现,避免影响 Broker 性能。

4.2 消息回溯注意事项

  • 幂等性是前提:无论哪种回溯方式,都必须保证消费者的幂等性,否则会导致重复处理业务(如重复发货、重复对账)。

  • 控制回溯范围:避免大范围回溯(如回溯几天前的所有消息),会导致消费者压力骤增,影响业务正常运行;建议精准回溯(如按时间范围、按偏移量范围)。

  • 存储时间配置:根据业务需求调整 Broker 消息存储时间(默认 72 小时),核心业务可配置为 7 天或更久,确保有足够的时间回溯消息;同时注意磁盘空间,避免消息过多导致磁盘溢出。

  • 回溯权限控制:消息回溯属于敏感操作(可能影响业务),生产环境中需添加权限控制,仅允许管理员操作,同时做好操作日志记录,便于追溯。

4.3 性能优化建议

  • Tag 过滤优化:同一个 Topic 下的 Tag 数量不宜过多(建议不超过 100 个),否则会增加 Broker 筛选压力;若 Tag 数量过多,建议拆分 Topic。

  • 回溯性能优化:回溯时尽量避开业务高峰期,选择低峰期操作;同时控制回溯的消息数量,避免一次性回溯大量消息导致消费者线程阻塞。

  • 监控过滤与回溯:通过 RocketMQ 控制台,监控消息过滤的命中率、回溯操作的频率和范围,及时发现异常(如过滤命中率过低、频繁大范围回溯)。

五、本篇核心总结及下一篇预告

RocketMQ 消息过滤的核心是“特征匹配”,提供 Tag 过滤(轻量高效)和 SQL92 过滤(灵活复杂)两种核心方式,分别适用于不同场景,优先使用 Tag 过滤。

  • 消息回溯的核心是“重置消费偏移量”,支持按时间、按偏移量、按消息 ID 三种方式,解决历史消息找回问题,前提是消息未被清理且消费者保证幂等性。

  • 实操关键:Tag 过滤通过“生产者设 Tag、消费者订阅 Tag”实现,SQL92 过滤需开启 Broker 配置并设置消息属性,消息回溯通过 RocketMQAdmin 重置偏移量实现。

  • 生产落地:重点关注过滤性能、回溯范围、幂等性保障,做好权限控制和监控,避免影响系统稳定性和业务正常运行。

下一篇,我们将讲解 RocketMQ 进阶特性——死信队列与延迟消息,解决“消息消费失败处理”和“定时任务”问题,进一步完善 RocketMQ 生产环境落地能力。

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐