消息过滤
消费者订阅某个主题后,Apache RocketMQ 会将该主题中的所有消息投递给消费者。但是,如果您希望消费者仅接收与业务相关的消息,可以在 Apache RocketMQ Broker 上设置过滤器。本主题介绍了消息过滤功能及其工作原理,同时介绍了消息的分类方式,并提供了不同过滤方法的使用示例。
适用场景
Apache RocketMQ 遵循发布-订阅模式。作为一种面向消息的中间件,Apache RocketMQ 被广泛用于促进分布式上游和下游应用程序之间的通信。在实际场景中,应用程序可能会使用不同的方法来消费消息。这些应用程序都可以订阅同一个 Apache RocketMQ 主题,并可以通过设置过滤器,只接收与其相关的消息。
通过使用 Apache RocketMQ 的消息过滤功能,您可以有效地管理发送给不同消费者的消息。这可以防止您的系统因大量非关键任务消息而负担过重。
Apache RocketMQ 的消息过滤功能在主题级别生效,允许您管理分布在多个服务中的同一业务消息。如果您需要管理不同业务的消息,可以订阅不同的主题。
功能概述
定义
消息过滤功能根据消费者配置的条件对消息进行过滤,并将符合条件的消息发送给消费者。
首先,在 Apache RocketMQ 生产者和消费者上定义消息属性和标签。然后,在消费者端设置过滤条件,Apache RocketMQ Broker 会根据这些条件过滤消息,并将过滤后的消息发送给消费者。
工作机制 
消息过滤包含以下步骤:
生产者:生产者在初始化消息之前,会为消息添加属性和标签。这些属性和标签用于匹配消费者设置的过滤条件。
消费者:消费者在消息初始化和消费过程中,调用订阅注册操作,将订阅的主题和消息(或过滤条件)告知 Broker。
Broker:在收到消费者的消息请求时,Apache RocketMQ Broker 会根据消费者提交的过滤条件表达式动态过滤消息,并将匹配过滤条件的消息发送给消费者。
分类
Apache RocketMQ 支持基于标签(Tag)的过滤和基于属性的 SQL 过滤。下表对比了这两种方法。
| 项目 | 基于标签的过滤 | 基于属性的 SQL 过滤 |
|---|---|---|
| 过滤目标 | 消息标签(Tags)。 | 消息属性,包括自定义属性和系统属性。消息标签属于系统属性 (TAGS)。 |
| 过滤能力 | 精确匹配。 | 基于 SQL 语法的匹配。 |
| 适用场景 | 基于标签的简单过滤。 | 涉及标签与属性之间关系的复杂过滤。 |
有关如何使用这些过滤方法的更多信息,请参阅 基于标签的过滤 和 基于属性的 SQL 过滤。
订阅一致性
过滤表达式是订阅的一部分。根据 Apache RocketMQ 的发布-订阅模式,消费者组内各消费者的订阅必须保持一致,包括过滤表达式,以避免某些消息无法被消费。有关更多信息,请参阅 订阅。
基于标签的过滤
基于标签的过滤是 Apache RocketMQ 提供的基础消息过滤功能。该功能基于生产者设置的标签对消息进行过滤。消费者使用这些标签指定要消费的消息。
适用场景
下图展示了电子商务交易场景中的一个示例。从下单到收货的过程中会产生一系列消息,例如:
订单消息
支付消息
物流消息
这些消息被发送到名为 Trade_Topic 的主题,该主题有多个系统作为订阅者,包括:
支付系统:仅订阅支付消息。
物流系统:仅订阅物流消息。
交易成功率分析系统:订阅订单消息和支付消息。
实时计算系统:订阅所有消息。

标签设置
生产者在发送消息之前,每条消息只能添加一个标签。
标签是一个字符串。建议字符串的最大长度为 128 个字符。
过滤规则
基于标签的过滤实现的是基于字符串的精确过滤。您可以设置以下过滤规则:
单标签匹配:您可以将过滤表达式设置为单个标签,以仅接收携带该标签的消息。
多标签匹配:您可以在过滤表达式中设置多个标签,以接收携带其中任何一个标签的消息。使用两个竖线 (||) 分隔标签。||例如:Tag1||Tag2||Tag3 表示所有带有 Tag1、Tag2 或 Tag3 的消息都会被发送给消费者。
全部匹配:您可以使用星号 (*)*来匹配所有标签,这意味着主题中的所有消息都会被发送给消费者。
示例
设置标签并发送消息
Message message = messageBuilder.setTopic("topic")
// Specify the message index key so that the system can use a keyword to accurately locate the message.
.setKeys("messageKey")
// Specify the message tag so that consumers can use the tag to filter the message.
// This example indicates that the tag of the message is set to "TagA".
.setTag("TagA")
// Message body.
.setBody("messageBody".getBytes())
.build();
指定一个标签并订阅消息
String topic = "Your Topic";
// Subscribe to messages that carry tag "TagA".
FilterExpression filterExpression = new FilterExpression("TagA", FilterExpressionType.TAG);
pushConsumer.subscribe(topic, filterExpression);
指定多个标签并订阅消息
String topic = "Your Topic";
// Subscribe to messages that carry tag TagA, TagB, or TagC.
FilterExpression filterExpression = new FilterExpression("TagA||TagB||TagC", FilterExpressionType.TAG);
pushConsumer.subscribe(topic, filterExpression);
订阅主题中的所有消息
String topic = "Your Topic";
// Subscribe to all messages.
FilterExpression filterExpression = new FilterExpression("*", FilterExpressionType.TAG);
pushConsumer.subscribe(topic, filterExpression);
基于属性的 SQL 过滤
基于属性的 SQL 过滤是 Apache RocketMQ 提供的进阶消息过滤方法。它根据生产者为消息配置的属性和属性值(也称为键和值)对消息进行过滤。生产者可以为一条消息设置多个属性。消费者可以在 SQL 表达式中指定属性,以接收特定的消息。
由于标签也是一种系统属性,因此基于标签的过滤实际上是基于属性的 SQL 过滤的一种特殊形式。在 SQL 语法中,标签属性表示为 TAGS。
适用场景
下图展示了电子商务交易场景中的一个示例。从下单到收货的过程中产生一系列消息,这些消息被分为订单消息和物流消息。物流消息配置了 region(地区)属性,其值为 Hangzhou(杭州)和 Shanghai(上海)。
订单消息
物流消息
region 属性值为 Hangzhou 的物流消息。
region 属性值为 Shanghai 的物流消息。
这些消息发送到名为 Trade_Topic 的主题,该主题包含以下订阅系统:
物流系统 1:仅订阅 region 属性值为 Hangzhou 的物流消息。
物流系统 2:订阅所有物流消息。
订单跟踪系统:仅订阅订单消息。
实时计算系统:订阅所有消息。

消息属性设置
生产者可以在发送消息前为消息设置自定义属性。每个属性都是一个自定义的键值对。
一条消息可以设置多个属性。
过滤规则
编写过滤表达式时必须遵循 SQL92 语法。具体说明如下:
| 语法 | 描述 | 示例 |
|---|---|---|
| IS NULL | 指定属性不存在。 | a IS NULL:属性 a 不存在。 |
| IS NOT NULL | 指定属性存在。 | a IS NOT NULL:属性 a 存在。 |
| > >= < <= | 用于比较数值。此语法不能用于比较字符串。如果用于比较字符串,消费者启动时会报错。注意:可转换为数值的字符串也被视为数值。 | a IS NOT NULL AND a > 100:属性 a 存在且其值大于 100。 a IS NOT NULL AND a > 'abc':错误示例。abc 是一个字符串,因此无法将 a 与 abc 进行数值比较。 |
| BETWEEN xxx AND xxx | 用于比较数值。此语法不能用于比较字符串。如果用于比较字符串,消费者启动时会报错。该语法等同于>>= xxx AND <= xxx。表示属性值在两个数值之间或等于这两个数值中的任意一个。 | a IS NOT NULL AND (a BETWEEN 10 AND 100):属性 a 存在且其值大于等于 10 且小于等于 100。 |
| NOT BETWEEN xxx AND xxx | 用于比较数值。此语法不能用于比较字符串。如果用于比较字符串,消费者启动时会报错。该语法等同于 < xxx OR>> xxx。表示属性值小于左侧数值或大于右侧数值。 | a IS NOT NULL AND (a NOT BETWEEN 10 AND 100):属性 a 存在且其值小于 10 或大于 100。 |
| IN (xxx, xxx) | 表示属性值包含在集合中。集合中的元素只能是字符串。 | a IS NOT NULL AND (a IN ('abc', 'def')):属性 a 存在且其值为 abc 或 def。 |
| = <> | = 和 <> 是等于和不等于运算符。它们可用于比较数值和字符串。 | a IS NOT NULL AND (a = 'abc' OR a<>'def'):属性 a 存在且其值为 abc,或者其值不为 def。 |
| AND OR | 逻辑 AND 和逻辑 OR 运算符。它们用于组合简单的逻辑函数,且每个逻辑函数必须放在括号内。 | a IS NOT NULL AND (a > 100) OR (b IS NULL):属性 a 存在且值大于 100,或者属性 b 不存在。 |
SQL 属性过滤通过配置自定义消息属性并定义 SQL 过滤表达式来实现。过滤表达式可能无法生成有效结果。Apache RocketMQ Broker 处理消息的逻辑如下:
异常处理:如果在评估过滤表达式时报告异常,Broker 默认会将收到的消息过滤掉,而不将其投递给消费者。例如,当比较数值和非数值时会发生异常。
Null 值处理:如果过滤表达式的计算结果为 NULL 或非布尔值,Broker 默认会将收到的消息过滤掉,而不将其投递给消费者。布尔值表示真值(true 或 false)。假设您没有为生产者发送的消息配置某个自定义属性,但该属性被用作 SQL 表达式中的过滤条件,此时过滤表达式的计算结果即为 NULL。
数值不一致处理:如果自定义消息属性的值为浮点数,但过滤表达式中使用的属性值为整数,Broker 默认会将收到的消息过滤掉,而不将其投递给消费者。
示例
为消息设置标签和属性并发送消息
Message message = messageBuilder.setTopic("topic")
// Specify the message index key so that the system can use a keyword to accurately locate the message.
.setKeys("messageKey")
// Specify the message tag so that consumers can use the tag to filter the message.
// This example indicates that the message tag is set to "messageTag".
.setTag("messageTag")
// You can also set custom attributes for the messages, such as environment, region, and logical branch.
// In this example, the custom attribute is region and the attribute value is Hangzhou.
.addProperty("Region", "Hangzhou")
// Message body.
.setBody("messageBody".getBytes())
.build();
基于自定义属性订阅和过滤消息
String topic = "topic";
// Subscribe only to messages whose value of the region attribute is Hangzhou.
FilterExpression filterExpression = new FilterExpression("Region IS NOT NULL AND Region='Hangzhou'", FilterExpressionType.SQL92);
simpleConsumer.subscribe(topic, filterExpression);
基于多个自定义属性订阅和过滤消息
String topic = "topic";
// Subscribe to messages whose value of the region attribute is Hangzhou and value of the price attribute is greater than 30.
FilterExpression filterExpression = new FilterExpression("Region IS NOT NULL AND price IS NOT NULL AND Region = 'Hangzhou' AND price > 30", FilterExpressionType.SQL92);
simpleConsumer.subscribe(topic, filterExpression);
订阅主题中的所有消息
String topic = "topic";
// Subscribe to all the messages.
FilterExpression filterExpression = new FilterExpression("True", FilterExpressionType.SQL92);
simpleConsumer.subscribe(topic, filterExpression);
使用说明
合理设置消息的主题和标签。
您可以使用主题、标签和属性来划分消息。在划分消息时,请注意以下事项:
消息类型:不同类型的消息(如顺序消息和普通消息)必须使用不同的主题进行划分。不要使用标签来划分消息类型。
业务领域:不同的业务领域和部门必须使用不同的主题。例如,物流消息和支付消息的主题必须不同。物流消息可以使用标签进一步细分为普通消息和紧急消息。
消息数量和重要性:数量或链路重要性不同的消息应拆分到不同的主题中。