RocketMQ 系列文章(进阶篇第 3 篇):消息过滤与消息回溯实战指南
前言:从“泛消费”到“精准控”的进阶需求
在上一篇中,我们掌握了 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 测试验证(回溯功能)
-
先发送几条订单创建消息,确保物流服务已消费。
-
调用 /order/backtrack/time,传入 topic=order_filter_topic、consumerGroup=logistics-consumer-group、backtrackTime=消息发送时间,观察物流服务是否重新消费该时间点后的消息。
-
调用 /order/backtrack/offset,传入队列ID、目标偏移量(如1),观察物流服务是否从偏移量1开始重新消费。
-
调用 /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 生产环境落地能力。
更多推荐




所有评论(0)