资讯动态

RocketMQ消息过滤实战:Tag与SQL92过滤原理、代码与避坑指南

发布时间:2026/9/16 3:12:25 来源:尧图企业网站定制
1. 消息过滤到底在解决什么问题1.1 消费端为什么需要消息过滤做RocketMQ开发的朋友应该都有这种经历一个Topic里塞了多种业务事件比如订单创建、订单支付、订单超时关闭全都往一个Topic里发。消费端如果啥都接收就要写一堆if else去判断消息类型再决定走哪条处理逻辑。这样不是不行但代码会很臃肿而且网络带宽、存储、消费线程全部被没用的数据白白占用。RocketMQ原生API提供的消息过滤能力就是解决这个痛点的。它允许消费者在订阅时声明自己只关心哪些消息让Broker或消费端在投递前就把不匹配的消息剔除掉。这样一方面减少了无效网络传输另一方面让消费者代码更干净业务边界更清晰。我最早用RocketMQ时也天真地以为“过滤就是消费端收到消息后自己判断一下”。后来才发现原生API里已经帮你安排好了两套方案一套是按Tag过滤一套是按SQL92属性过滤。用好了它们系统资源省下来一大截线上问题也少很多。1.2 原生API中两种过滤方式的场景差异Tag过滤是RocketMQ最基础、也是大部分教程里最常见的一种消息过滤方式。每个消息可以打一个Tag比如OrderCreate、OrderPay。消费者订阅时指定Tag表达式Broker端在投递前直接比较Tag字符串能匹配才发给消费者。Tag过滤的特点是简单、性能高适合业务类型比较粗粒度分类的场景。比如一个订单Topic里区分创建、支付、关闭用Tag就能搞定。SQL92过滤则要更灵活一些它是基于消息属性来做条件匹配的。消息发送时可以附带一组自定义属性比如price 100、region CN、isVip true。消费者通过SQL92语法写过滤表达式比如TAGS IS NOT NULL AND price BETWEEN 100 AND 200。它解决的是更细粒度、依赖多个维度的过滤需求。SQL92过滤的代价是计算量更大而且对Broker配置和版本有要求。所以选择哪种过滤方式本质上是在灵活性和资源消耗之间做取舍。后面我会结合代码详细讲清楚两条路的差异。2. 原生API消息过滤核心代码样例2.1 生产者给消息打上过滤标记写过滤样例前我们先把生产者侧的准备工作做扎实。使用RocketMQ原生API时消息发送方要做的就是给消息设置Tag和属性这两个字段就是后续过滤的依据。先看一个最常规的生产者样例我习惯用4.9.x版本的rocketmq-clientimport org.apache.rocketmq.client.producer.DefaultMQProducer; import org.apache.rocketmq.common.message.Message; public class FilterProducerDemo { public static void main(String[] args) throws Exception { DefaultMQProducer producer new DefaultMQProducer(filter_producer_group); producer.setNamesrvAddr(127.0.0.1:9876); producer.start(); // 消息1普通订单创建事件TagOrderCreate Message msg1 new Message(OrderTopic, OrderCreate, order_001.getBytes()); msg1.putUserProperty(region, CN); msg1.putUserProperty(price, 199); producer.send(msg1); // 消息2VIP订单支付事件TagOrderPay增加isVip属性 Message msg2 new Message(OrderTopic, OrderPay, order_002.getBytes()); msg2.putUserProperty(region, US); msg2.putUserProperty(price, 899); msg2.putUserProperty(isVip, true); producer.send(msg2); // 消息3普通订单超时关闭事件TagOrderClose Message msg3 new Message(OrderTopic, OrderClose, order_003.getBytes()); msg3.putUserProperty(region, CN); msg3.putUserProperty(price, 59); producer.send(msg3); System.out.println(消息发送完成); producer.shutdown(); } }这段代码里new Message(topic, tag, body)的第二个参数就是Tag。putUserProperty则是设置自定义属性后续SQL92过滤会从这里取字段。这里有一个关键点Tag和属性都是字符串存储如果你要在SQL92里做数值比较属性值也得存成可被解析的格式。RocketMQ在解析时会尝试转换类型但如果你习惯存带单位或前导零的字符串过滤条件很容易失效。2.2 消费者按Tag订阅的核心写法接下来是消费者侧。使用原生API时最常用的是DefaultMQPushConsumer。它支持通过MessageSelector来声明订阅条件其中MessageSelector.byTag()就是Tag过滤入口。一个完整的Tag过滤消费者样例import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer; import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext; import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus; import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently; import org.apache.rocketmq.client.consumer.MessageSelector; import org.apache.rocketmq.common.message.MessageExt; import java.util.List; public class TagFilterConsumerDemo { public static void main(String[] args) throws Exception { DefaultMQPushConsumer consumer new DefaultMQPushConsumer(filter_consumer_group); consumer.setNamesrvAddr(127.0.0.1:9876); // 按Tag过滤只消费OrderCreate和OrderPay consumer.subscribe(OrderTopic, MessageSelector.byTag(OrderCreate || OrderPay)); consumer.registerMessageListener(new MessageListenerConcurrently() { Override public ConsumeConcurrentlyStatus consumeMessage(ListMessageExt msgs, ConsumeConcurrentlyContext context) { for (MessageExt msg : msgs) { System.out.printf(收到过滤后的消息: Tag%s, keys%s, body%s%n, msg.getTags(), msg.getKeys(), new String(msg.getBody())); } return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; } }); consumer.start(); System.out.println(Tag过滤消费者启动成功); } }注意MessageSelector.byTag(OrderCreate || OrderPay)这里用的是双竖线分隔多个Tag。很多新手会写成逗号或|单竖线结果消息一条都收不到。RocketMQ的Tag表达式语法比较严格||表示或的关系单个Tag不需要分隔符。Tag过滤是Broker端完成的也就是说消息在进入消费者之前已经被Broker筛掉了一部分。这样消费者拿到的消息Tag一定是你订阅表达式里的值。2.3 消费者按SQL92属性过滤的进阶写法如果Tag过滤满足不了你比如你要按价格区间、按区域代码、按多个业务属性组合过滤就该用MessageSelector.bySql()了。看下面这个SQL92过滤样例import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer; import org.apache.rocketmq.client.consumer.MessageSelector; import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext; import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus; import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently; import org.apache.rocketmq.common.message.MessageExt; import java.util.List; public class SqlFilterConsumerDemo { public static void main(String[] args) throws Exception { DefaultMQPushConsumer consumer new DefaultMQPushConsumer(sql_filter_consumer_group); consumer.setNamesrvAddr(127.0.0.1:9876); // 按SQL92过滤只消费CN地区、价格在100到500之间的非关闭消息 consumer.subscribe(OrderTopic, MessageSelector.bySql((TAGS IS NOT NULL AND TAGS NOT IN (OrderClose)) AND region CN AND price BETWEEN 100 AND 500)); consumer.registerMessageListener(new MessageListenerConcurrently() { Override public ConsumeConcurrentlyStatus consumeMessage(ListMessageExt msgs, ConsumeConcurrentlyContext context) { for (MessageExt msg : msgs) { System.out.printf(SQL过滤消息: Tag%s, region%s, price%s%n, msg.getTags(), msg.getProperty(region), msg.getProperty(price)); } return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; } }); consumer.start(); System.out.println(SQL92过滤消费者启动成功); } }这里有个细节TAGS是RocketMQ内置字段可以直接在SQL表达式里使用。也就是说你可以同时用Tag和自定义属性做组合过滤不是只能二选一。另外SQL92过滤对Broker版本有要求我记得4.4.0之后才逐渐完善。如果用的是老版本可能压根不支持bySql启动时会报错需要在服务端配置开启相关组件。这一点放到后面的坑位部分细说。3. 消息过滤背后的原理与性能考量3.1 Broker端过滤与客户端过滤的取舍很多人只知道“用Tag过滤更快、SQL过滤更慢”但不清楚为什么。要理解这个问题得先看看RocketMQ在投递消息时做了哪些事情。当消费者订阅了某个Topic时Broker会为它构建过滤条件。消息到达Broker后如果配置了过滤Broker会先检查这条消息的Tag或属性是否满足消费者的订阅表达式不满足就不投递给这个消费者。这就是Broker端过滤好处是无效消息不会占用Consumer与Broker之间的网络IO消费者拿到手的都是“精筛”过的数据。但它也有代价每个消息都要在Broker端执行一次匹配计算。Tag匹配本质上是字符串比较非常快而SQL92过滤要解析表达式、做类型转换、跑比较逻辑CPU消耗会明显上升。尤其在Topic消息量特别大、订阅关系特别多时Broker端过滤可能成为瓶颈。为了缓解Broker压力RocketMQ也设计了一种客户端过滤机制由Consumer拉取到本地后再筛。但原生API里并不推荐大家手动实现因为处理逻辑相当复杂。实际业务里如果过滤条件特别复杂、Broker端算不动通常的解法是拆分Topic把不同业务类型直接分到不同的Topic上而不是把所有消息塞到一个Topic里再猛做过滤。这实际上就是建议你能用Topic划分业务就用TopicTag是第二选择SQL92是兜底手段。3.2 订阅关系如何影响消费行为RocketMQ的消费者是主动拉取消息的但PushConsumer封装了长轮询看起来像服务端推送。过滤条件在订阅关系里被注册消费者发起拉取请求时Broker会拿着当前消费者的订阅关系去匹配。这里有一个很值得注意的点同一个消费组内的所有消费者订阅关系必须保持一致。比如你有两个实例订阅了相同的Topic但一个订阅OrderCreate || OrderPay另一个订阅OrderClose这样会导致消费行为异常甚至消息重复或丢失。我遇到过真实案例某团队上了多个消费者实例某个实例临时改了Tag表达式去调试结果整个消费组的行为变得混乱部分消息被错误分发给了一些本不该接到的实例。后来排查发现就是订阅关系不一致导致的。所以如果你在项目里动态修改订阅条件一定要确保整个消费组的策略同步不要只改一台机器。3.3 消息属性与过滤语法细节SQL92过滤虽然灵活但它的字段来源是消息属性。在生产者里通过putUserProperty设置的键值对会被保存为消息的properties。消费者可以使用普通属性来做过滤也可以使用几个内置属性如TAGS、KEYS、WAIT等。内置属性的具体用法RocketMQ官方文档给过一段参考TAGS IS NOT NULL AND region CN AND price 100还可以用IN、BETWEEN、LIKE、IS NULL、IS NOT NULL、逻辑运算符AND、OR、NOT。不过要注意SQL92过滤表达式的语法和能力范围是RocketMQ实现的一个子集不是完整的SQL。像REPLACE、字符串拼接这类函数就不要想了。属性类型转换也容易踩坑。price被存成字符串899你在SQL里写price BETWEEN 100 AND 500RocketMQ会尝试把属性值转成数值但如果你存了899元这种带中文的内容转换失败后的行为可能与你预期不同。最稳妥的做法是凡是参与数值过滤的属性在生产者侧确保存储的是纯数字字符串。4. 实战中的常见坑与调优建议4.1 Tag过滤的经典坑位Tag过滤看起来简单坑是真不少。我列几个自己踩过或被身边同事问过的典型问题。第一多个Tag的表达式必须是||连接不能写逗号不能写|单竖线。这个前面已经提过。另外Tag本身不能为空字符串如果你想过滤所有带Tag的消息不能写空串得写TAGS IS NOT NULL。第二Tag与Topic的关系要理顺。同一个消息只能有一个Tag但一个Tag可以代表一类业务事件。这种一对多的归类方式决定了Tag适合做粗粒度分类。如果你有“同时属于两个分类”的消息单一Tag做不到最好在属性里加一个category字段然后用SQL92过滤。第三订阅时如果传了*表示接收Topic下的所有消息。网上很多教程这么写但生产环境不建议滥用因为全量订阅会让一些原本不需要的脏数据也进入业务逻辑。真需要全量接收的场景不如直接换个Topic命名策略。第四Tag过滤是精确匹配还是模糊匹配我要强调它是精确匹配不能用通配符。之前有同事写byTag(Order*)想匹配Order开头的消息结果一条都收不到。RocketMQ的Tag表达式不支持*模糊匹配只有||这种多值匹配。4.2 SQL92过滤的兼容性坑位SQL92过滤最容易遇到的问题就是Broker配置不支持。很多生产环境用的是RocketMQ 4.x系列老版本默认没有开启SQL过滤属性。如果你启动消费者时报错第一种情况是message selector is not supported之类。解决办法通常需要修改Broker的配置文件加上enablePropertyFiltertrue enableCalcFilterBitMaptrue改完后要重启Broker。注意集群环境下所有Broker节点都要同步配置不然消费者可能在某台Broker上正常换一台就异常。另外SQL92过滤表达式的写法很严格。比如字符串必须用单引号括起来属性字段名区分大小写。regionCN和REGIONCN是不同的因为在SQL92语法里属性名被视为标识符RocketMQ实现时对大小写敏感。这一点相当坑因为生产者putUserProperty(region, CN)小写消费者表达式里写成REGION就会一直匹配不上。更隐蔽的一个坑是消息属性里如果存在多个同名属性怎么办RocketMQ的属性表是Map结构后写入的会覆盖前面的所以不用太担心重复键的问题。4.3 过滤与顺序消息、事务消息的组合在实践里过滤机制很少单打独斗它经常和顺序消息、事务消息一起用。比如订单状态流转你可能既要保证同一个订单的消息按顺序消费又希望只处理支付成功后的后续动作。这种情况下Tag或SQL过滤仍然有效但要注意顺序消息是基于队列维度的同一个订单的bin上的消息必须路由到同一个MessageQueue。过滤条件可能会影响消息的分配行为尤其是Tag不同但订单ID相同的情况如果生产者在发送时选择MessageQueue的规则不统一那么即使过滤正确顺序也可能乱。事务消息的消息在commit之后才会被消费者看到过滤机制对它同样生效。我常用的方式是在事务消息里也设置业务属性比如orderId、operator然后用SQL92过滤掉无效事件。不过要牢记事务消息的过滤条件是在消息可见性确认后才匹配的如果事务回滚了消费者根本不会看到那条消息过滤条件也无从谈起。4.4 过滤机制的性能调优建议在实际项目中如果Topic的消息量比较大过滤条件就不能只图写得开心还要考虑性能。我有几点经验可以参考优先用Tag过滤其次才是SQL92。能用两个Tag解决的问题不要写成SQL表达式。SQL92表达式尽量简单。避免使用大量NOT IN、OR组合因为每多一个判断条件Broker计算量都会上涨。如果确实需要复杂的SQL92过滤尽可能把“最容易过滤掉大多数消息”的条件放在最前面。虽然RocketMQ不一定按书写顺序短路执行但从编写习惯上保持这种倾向是好的。对消息量巨大的Topic考虑按业务拆分Topic而不是把过滤工作全压在Broker端。毕竟Topic拆分后消费组天然只收到对应业务的消息连过滤都不需要。消费者数量与队列数量的配比也要注意。过滤场景下如果Consumer实例过多而MessageQueue数量不足有些Consumer会空转显得消费“慢”或“不均衡”容易被误判成过滤条件有问题。5. 消息过滤问题排查实录5.1 常见问题速查表我把自己遇到的问题和从同事那边收集到的典型案例整理成了一张速查表遇到类似情况可以直接对照排查。现象可能原因解决办法消费者收不到任何消息订阅表达式写错如Tag用了逗号分隔检查subscribe条件和Tag表达式语法改为只收到部分Tag的消息多个消费者实例订阅关系不一致统一消费组内所有实例的订阅表达式SQL92过滤启动报错Broker未开启属性过滤配置在broker.conf中加入enablePropertyFiltertrue后重启SQL过滤结果与预期不符属性字段大小写不一致或属性值类型无法解析核对生产者putUserProperty字段名和表达式字段名大小写过滤后消息延迟明显表达式过于复杂Broker CPU被打满简化表达式或拆分Topic同一个消费者重复收到同一条消息消费失败重试机制触发业务代码抛异常排查消费逻辑异常检查消费返回状态Tag过滤后仍然收到不匹配的TAGConsumer端本地Filter模式与Broker端Filter不一致检查是否使用老版本客户端并参考官方升级建议这张表看着简单真正排查时往往要组合使用。我自己的习惯是优先看消费组订阅关系日志然后看Broker日志最后再怀疑业务代码。5.2 用mqadmin与日志定位过滤问题RocketMQ自带mqadmin命令在排查消息过滤问题时非常好用。比如你想确认某个Topic下到底有哪些消息、Tag是什么可以执行mqadmin queryMsgByKey -n 127.0.0.1:9876 -t OrderTopic -k order_001这会通过消息Key查询消息详情其中能看到消息的Tag和properties。如果消费者收不到消息但消息确实存在且Tag也对那问题很可能出在订阅关系上。想看某个消费组的订阅情况可以用mqadmin consumerStatus -n 127.0.0.1:9876 -g filter_consumer_group这条命令会显示消费组内各实例的客户端连接和订阅信息。如果发现订阅表达式和你预期的不一致基本就是消费组内有“异类”实例注册了不同订阅条件。再配合Broker日志在${storePathRootDir}/store或日志目录下搜索目标Topic的消费请求记录能更细粒度地看到过滤命中情况。日志里通常会出现类似consume message filter、tags等关键字耐心点总能找到线索。还有一个笨但有效的方法在消费者里临时把订阅改成MessageSelector.byTag(*)如果这样能收到消息而特定Tag条件收不到说明是过滤条件的问题如果全量订阅也收不到说明是消息本身没发到该Topic或消费组的问题。用二分法缩小范围比瞎猜高效得多。我在实际排查中遇到过最隐蔽的问题是消息Key为空。queryMsgByKey查不到结果一度以为是Broker丢消息了。后来发现是生产者发送时没有设置setKeys导致查询命令无法定位。所以再次提醒生产环境给消息加上业务唯一Key既能做链路追踪也能方便出问题时快速定位。

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价