RocketMQ 消息过滤:Tag 过滤、SQL92 表达式与过滤原理
一句话结论(30s)
RocketMQ 消息过滤本质是三级漏斗——Tag 过滤、SQL92 表达式、消费端二次过滤——因为在线业务里一个 Topic 承载多种业务、消费者只关心其中一种,Broker 端过滤能省掉不相关消息的传输和消费成本。关键设计是 Tag 的 hashCode 存在 ConsumeQueue 的 20 字节条目里,Broker 用 O(1) 整数比较完成匹配且不读 CommitLog,这是它比「全量投递给客户端过滤」更高效的根本。核心权衡:SQL92 灵活但必须先读 CommitLog 解析属性(有磁盘 IO 与 CPU 开销),所以能用 Tag 就不用 SQL92,消费端二次过滤是防哈希冲突的最后防线。
核心原理(2min)
主流程:生产者给消息打 Tag(或设置用户属性),消费者订阅 Tag 表达式或 SQL92 表达式,Broker 端完成过滤后只推送匹配的消息。关键机制有三层:Tag 过滤在 ConsumeQueue 层用 8 字节 Tag Hash 做整数比较,O(1) 且不读 CommitLog、不解析消息体;SQL92 过滤要先读 CommitLog 取 properties 再执行表达式(支持比较/逻辑/IN/LIKE 前缀,不支持函数和 JOIN,需 enablePropertyFilter=true 开启);消费端二次过滤兜底 Tag 哈希冲突和不同版本 SQL 引擎的偏差。注意同一 ConsumerGroup 内订阅关系必须完全一致,否则 Rebalance 时会被过滤掉本属于对方的队列导致消息丢失。
底层深入(5-10min)
为什么 Kafka 不过滤,而 RocketMQ 要过滤?
Kafka 的设计原则是”Broker 只管存,消费者自己决定要什么”——所有消息全部投递给消费者,消费者自己过滤。这在海量日志场景下合理——反正每条日志都要处理,不需要过滤。
但在线业务场景不同。一个 Topic 可能承载多种业务类型,消费者只关心其中一种。如果不做服务端过滤,不相关的消息也要传输和处理,浪费带宽和 CPU。
RocketMQ 的过滤在 Broker 端完成,只推送消费者真正关心的消息。
🤔 思考穿插:为什么 Broker 端过滤能省成本,Kafka 却坚持全量投递?—— 因为 Kafka 面向日志聚合,反正每条消息都要被处理,过滤没有收益;而在线业务一个 Topic 承载多种业务,消费者只关心一种,不相关消息若全量传输就是纯浪费带宽和 CPU。所以「要不要过滤」取决于「消息是否都需要被处理」。
三层过滤体系
Layer 1: Tag 过滤 → ConsumeQueue 的 Tags Hash 匹配(最快,O(1) 内存比较)
Layer 2: SQL92 表达式 → 解析消息属性,执行 SQL 逻辑(灵活,但有性能开销)
Layer 3: 消费端二次过滤 → 消费者拉取后再次检查(兜底,防止哈希冲突和 Broker 过滤漏洞)
一、Tag 过滤:最轻量、最常用
使用方式
// 生产者:给消息打 Tag
Message msg = new Message("OrderTopic", "PAID", body); // Tag = "PAID"
producer.send(msg);
// 消费者:订阅时指定 Tag
consumer.subscribe("OrderTopic", "PAID || SHIPPED || DELIVERED");
// ↑ 只接收这三种 Tag 的消息
过滤原理
Tag 过滤的核心机制在 ConsumeQueue 中:
还记得 ConsumeQueue 的 20 字节条目吗?
┌──────────────────────┬──────────────┬──────────────────────┐
│ CommitLog Offset │ Size │ Tag Hash Code │
│ (8 bytes) │ (4 bytes) │ (8 bytes) │
└──────────────────────┴──────────────┴──────────────────────┘
↑
生产者发送时,RocketMQ 计算 Tag 的
hashCode 并写入这个字段
消费者拉取消息时,Broker 的处理流程:
1. 从 ConsumeQueue 中读取 20 字节条目
2. 提取 Tags Hash Code(8 字节)
3. 与消费者订阅的 Tag 表达式的 Hash Code 集合匹配
4. 匹配成功 → 去 CommitLog 读取完整消息 → 返回消费者
5. 匹配失败 → 跳过该条目,不读 CommitLog(节省 IO)
Tag 表达式
// 支持的操作符
consumer.subscribe("Topic", "TagA || TagB"); // 或
consumer.subscribe("Topic", "*"); // 所有 Tag(通配符)
consumer.subscribe("Topic", "TagA || *"); // 等价于 *
// 不支持 TagA && TagB —— 一条消息只能有一个 Tag
性能特点
- 极快:Hash Code 是整数比较,一个 CPU 周期完成。
- 不需要解析消息体:只需要 ConsumeQueue 条目中的 8 字节。
- 不读 CommitLog:过滤掉的消息完全不产生磁盘 IO。
局限性
- 每个消息只能有一个 Tag:不能做复杂的多维度过滤。
- 哈希冲突:不同 Tag 的 hashCode 相同 → 消费者收到不该收到的消息 → 消费端需要二次过滤。
🤔 思考穿插:为什么 Tag 用 hashCode 过滤就快,却要消费端二次过滤兜底?—— 因为 hashCode 是 8 字节整数,整数比较是 O(1)、不读磁盘;但不同 Tag 可能哈希冲突导致误投,所以消费者要用完整 Tag 字符串再校验一次。所以「快」和「准」之间存在冲突,兜底才能两全。
- 无法基于消息内容过滤:Tag 只是标签,不是消息体中的字段。
二、SQL92 表达式过滤:灵活但有代价
使用方式
当 Tag 不够用时,SQL92 表达式提供了更灵活的过滤。
// 生产者:设置消息属性
Message msg = new Message("OrderTopic", "TagA", body);
msg.putUserProperty("orderId", "10086");
msg.putUserProperty("amount", "150.00");
msg.putUserProperty("region", "BEIJING");
producer.send(msg);
// 消费者:SQL92 表达式
consumer.subscribe("OrderTopic", MessageSelector.bySql(
"(amount > 100 AND region = 'BEIJING') OR orderId = 'VIP001'"
));
SQL92 支持的操作符
| 类型 | 操作符 |
|---|---|
| 比较 | =, <>, >, >=, <, <= |
| 逻辑 | AND, OR, NOT |
| 空值 | IS NULL, IS NOT NULL |
| 集合 | IN (a, b, c) |
| 字符串匹配 | LIKE 'prefix%'(仅前缀匹配) |
不支持:函数调用(MAX, COUNT)、子查询、JOIN、BETWEEN。
过滤原理
SQL92 过滤比 Tag 过滤多了一步——需要解析消息属性。
Tag 过滤流程:
ConsumeQueue Tag Hash 匹配 → CommitLog 读消息 → 返回
SQL92 过滤流程:
ConsumeQueue Tag Hash 匹配 → CommitLog 读消息 →
解析消息属性(properties)→ 执行 SQL 表达式 → 匹配才返回
关键差异:SQL92 过滤必须先读 CommitLog——因为消息属性(properties)存储在 CommitLog 的消息体中,不在 ConsumeQueue 中。这意味着:
- SQL92 过滤一定会产生磁盘 IO(读取 CommitLog 获取 properties)。
- 不匹配的消息已经被从磁盘读出来了,浪费了 IO,只是省了网络传输。
// RocketMQ Broker 端 SQL92 过滤核心代码(简化)
public boolean isMatched(MessageExt msg, SubscriptionData subscriptionData) {
if (subscriptionData.isClassFilterMode()) {
return true;
}
String expr = subscriptionData.getSubString();
// 解析消息 properties → 传入 SQL 表达式引擎
Map<String, String> properties = msg.getProperties();
// 使用 Apache Commons Jexl / 自研轻量引擎 执行
return sqlExpressionEngine.evaluate(expr, properties);
}
性能开销
| 维度 | Tag 过滤 | SQL92 过滤 |
|---|---|---|
| 是否需要读 CommitLog | 否(只读 ConsumeQueue) | 是 |
| 是否需要解析消息属性 | 否 | 是 |
| 过滤计算复杂度 | O(1),整数比较 | O(表达式复杂度),字符串解析 |
| 额外磁盘 IO | 无 | 需要读不匹配的消息的 properties |
| CPU 开销 | 极低 | 中等(需解析多个属性值 + 执行表达式) |
建议:能用 Tag 过滤就不用 SQL92。只有当需要基于业务字段过滤(如地区、金额、用户等级)且这些字段不能提升为 Tag 时,才使用 SQL92。
🤔 思考穿插:为什么 SQL92 比 Tag 灵活,却建议「能用 Tag 就不用 SQL92」?—— 因为 Tag 的 hash 存在 ConsumeQueue 里、不读 CommitLog;SQL92 的 properties 存在消息体里,必须先读 CommitLog 解析属性,带来磁盘 IO 和 CPU 开销。所以灵活性是有代价的,能用简单标签就别上表达式引擎。
三、消费端二次过滤:最后的防线
即使 Broker 已经过滤了,消费者仍然应该做二次过滤。两个原因:
- 哈希冲突:不同 Tag 可能有相同 hashCode,Broker 会误投。
- Broker 版本的 SQL 实现偏差:不同 RocketMQ 版本的 SQL 引擎行为可能略有不同。
consumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) -> {
List<MessageExt> filtered = msgs.stream()
.filter(msg -> {
// Tag 哈希冲突兜底:检查完整 Tag 字符串
String expectedTag = "PAID";
return expectedTag.equals(msg.getTags());
// 或者 基于属性的二次过滤
// String region = msg.getProperty("region");
// return "BEIJING".equals(region);
})
.collect(Collectors.toList());
if (filtered.isEmpty()) {
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
}
processMessages(filtered);
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
});
过滤配置与限制
开启 SQL92 过滤
# broker.properties
enablePropertyFilter = true # 开启 SQL92 过滤(默认关闭)
注意:SQL92 过滤需要额外的 CPU 开销,默认关闭是为了保护 Broker。在启用前务必评估 Broker 的 CPU 负载。
过滤在集群和广播模式下的表现
| 消费模式 | Tag 过滤 | SQL92 过滤 |
|---|---|---|
| 集群(CLUSTERING) | ✅ Broker 端过滤 | ✅ Broker 端过滤 |
| 广播(BROADCASTING) | ✅ Broker 端过滤 | ✅ Broker 端过滤 |
两种模式下的过滤行为一致——都在 Broker 端执行。
订阅关系一致性
RocketMQ 要求同一个 Consumer Group 内所有 Consumer 的订阅关系(Topic + Tag/SQL)必须完全一致。如果不同:
Consumer1: subscribe("Topic", "TagA")
Consumer2: subscribe("Topic", "TagB")
结果:Rebalance 时 Consumer2 可能拿到 Consumer1 的 Queue,
但 TagB 过滤掉了本应属于 Consumer1 的消息 → 消息丢失!
> **🤔 思考穿插**:为什么同一 ConsumerGroup 内订阅关系必须一致,否则会丢消息?—— 因为 Rebalance 时队列会在实例间重新分配,若两个实例订阅的 Tag 不同,A 的队列分给 B 后会被 B 的过滤规则误滤掉本属于 A 的消息。所以订阅关系不一致是「隐性丢消息」的经典坑。
实践中常见的问题
1. Tag 太多了怎么办?
Tag 本质上是”消息一级分类”。如果分类太多(几十个 Tag),考虑:
// ❌ 不要这样——Tag 爆炸
msg.setTags("PAID_BEIJING_VIP");
msg.setTags("PAID_SHANGHAI_VIP");
msg.setTags("PAID_BEIJING_NORMAL");
// ... 50 个 Tag
// ✅ 改为 Tag + 属性组合
msg.setTags("PAID"); // Tag 只表示状态
msg.putUserProperty("region", "BEIJING");
msg.putUserProperty("level", "VIP");
// 消费者用 SQL92 过滤:region='BEIJING' AND level='VIP'
2. Tag 过滤不生效?
检查:MessageListener 中收到的消息是否 Tag 不匹配?大概率是哈希冲突。在消费端检查 msg.getTags() 并做过滤。
3. SQL92 过滤导致消费延迟增加?
检查 Broker CPU 和磁盘 IO。SQL92 需要读 CommitLog 获取 properties,高 QPS 下这会产生大量随机读。考虑:
- 能否将 SQL92 条件转换为 Tag?(如
status='PAID'→ Tag=“PAID”) - 能否降低 SQL 表达式的复杂度?(减少 AND/OR 嵌套)
- 把过滤逻辑移到消费端(虽然多传数据,但 Broker 压力小)
总结
RocketMQ 的消息过滤是三级漏斗:
全部消息(CommitLog)
│
▼ Tag 过滤(ConsumeQueue Tag Hash,O(1),不读磁盘)
│ 淘汰约 80%-90% 不相关消息
│
▼ SQL92 过滤(读 CommitLog properties,解析表达式)
│ 在剩余 10%-20% 中进一步筛选
│
▼ 消费端二次过滤(检查完整 Tag 字符串 / 业务字段)
│ 最终防线(处理哈希冲突和边界情况)
最佳实践:Tag 做主分类(状态、类型),SQL92 做辅助筛选(地区、金额),消费端做二次兜底。三层配合,既保证了性能,又保证了准确性。
章末提问
追问 1:Tag 过滤为什么能做到 O(1) 且不读 CommitLog?
结论先行:因为 Tag 的 hashCode 被写进 ConsumeQueue 的 20 字节固定条目里(其中 8 字节是 Tag Hash),Broker 只需整数比较就能判断匹配。
因为:ConsumeQueue 是纯索引、不含消息体,匹配只需读那 8 字节做 hash 比对,不匹配就跳过、完全不读 CommitLog;整数比较在一个 CPU 周期内完成,所以是 O(1) 且零磁盘 IO。
追问 2:SQL92 过滤和 Tag 过滤的性能差异,根因是什么?
结论先行:根因是 SQL92 的过滤字段(properties)存在消息体里,必须先读 CommitLog 才能拿到;而 Tag 的 hash 直接存在 ConsumeQueue 索引里。
因为:Tag 过滤只需读索引、做整数比较;SQL92 必须先读 CommitLog 取 properties 再执行表达式,不匹配的消息也已经被从磁盘读出来了、白费 IO 和 CPU,只是省了网络传输。所以 SQL92 更灵活但更贵。
追问 3:同一个 ConsumerGroup 内订阅关系不一致,会有什么后果?
结论先行:Rebalance 时队列重新分配,订阅关系不同的实例会用自己的过滤规则误滤掉本属于对方的队列,导致消息丢失。
因为:RocketMQ 要求同一 Group 内所有 Consumer 的订阅关系完全一致;否则队列落到订阅不同 Tag 的实例上时,消息会被错误过滤、无法被任何实例正确消费,这是「隐性丢消息」的经典场景,且很难排查。