Skip to content
Go back

RocketMQ 消息过滤:Tag 过滤、SQL92 表达式与过滤原理

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

性能特点

局限性

  1. 每个消息只能有一个 Tag:不能做复杂的多维度过滤。
  2. 哈希冲突:不同 Tag 的 hashCode 相同 → 消费者收到不该收到的消息 → 消费端需要二次过滤。

🤔 思考穿插:为什么 Tag 用 hashCode 过滤就快,却要消费端二次过滤兜底?—— 因为 hashCode 是 8 字节整数,整数比较是 O(1)、不读磁盘;但不同 Tag 可能哈希冲突导致误投,所以消费者要用完整 Tag 字符串再校验一次。所以「快」和「准」之间存在冲突,兜底才能两全。

  1. 无法基于消息内容过滤: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 中。这意味着:

  1. SQL92 过滤一定会产生磁盘 IO(读取 CommitLog 获取 properties)。
  2. 不匹配的消息已经被从磁盘读出来了,浪费了 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 已经过滤了,消费者仍然应该做二次过滤。两个原因:

  1. 哈希冲突:不同 Tag 可能有相同 hashCode,Broker 会误投。
  2. 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 下这会产生大量随机读。考虑:


总结

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 的实例上时,消息会被错误过滤、无法被任何实例正确消费,这是「隐性丢消息」的经典场景,且很难排查。


Share this post on:

Previous Post
RocketMQ 消息领域模型:Topic、Tag、MessageQueue、Group 与 Offset 的层次设计
Next Post
RocketMQ 消息积压处理:监控、流控、扩容与 DLQ 的全链路方案