178 lines
7.5 KiB
Markdown
178 lines
7.5 KiB
Markdown
|
|
---
|
|||
|
|
tags: [MQ, 消息过滤, 消息路由, Tag, SQL92, RocketMQ, RabbitMQ]
|
|||
|
|
create time: 2026-05-24 19:52
|
|||
|
|
---
|
|||
|
|
|
|||
|
|
# 消息过滤与路由
|
|||
|
|
|
|||
|
|
## 概述
|
|||
|
|
|
|||
|
|
一个 Topic 下可能有多种类型的消息,但某个 Consumer 只关心其中一部分。消息过滤(Filtering)让 Consumer 只接收自己需要的消息,消息路由(Routing)决定消息应该流向哪些队列或消费者。两者配合,才能在大规模系统中实现精准、高效的消息分发。
|
|||
|
|
|
|||
|
|
## 正文
|
|||
|
|
|
|||
|
|
### 消息过滤的需求
|
|||
|
|
|
|||
|
|
举个例子:电商系统的 `Topic_Order` 下有"创建"、"支付"、"取消"三种类型的消息。物流服务只关心"支付"类型,风控服务只关心"创建"和"取消"类型。如果没有过滤机制,每个 Consumer 都要接收全量消息再自行判断,浪费网络带宽和 CPU。
|
|||
|
|
|
|||
|
|
### Broker 端过滤 vs Consumer 端过滤
|
|||
|
|
|
|||
|
|
过滤发生在哪里,直接影响系统效率和灵活性:
|
|||
|
|
|
|||
|
|
| 维度 | Broker 端过滤 | Consumer 端过滤 |
|
|||
|
|
|------|-------------|----------------|
|
|||
|
|
| 网络传输 | 只传输匹配的消息,节省带宽 | 传输全部消息,Consumer 自行过滤 |
|
|||
|
|
| Broker 负载 | 增加(需要解析消息属性) | 无影响 |
|
|||
|
|
| 灵活性 | 受限于 Broker 支持的过滤语法 | 任意逻辑,完全灵活 |
|
|||
|
|
| 实现难度 | 高(Broker 需要理解消息语义) | 低(纯客户端逻辑) |
|
|||
|
|
| 典型代表 | RocketMQ Tag/SQL 过滤 | Kafka Consumer 自行过滤 |
|
|||
|
|
|
|||
|
|
```mermaid
|
|||
|
|
graph TD
|
|||
|
|
Producer["Producer"] -->|"发送消息\n带 Tag/Properties"| Broker["Broker"]
|
|||
|
|
|
|||
|
|
subgraph "Broker 端过滤"
|
|||
|
|
Broker -->|"根据过滤规则匹配"| Filter["过滤引擎"]
|
|||
|
|
Filter -->|"匹配的消息"| C1["Consumer A"]
|
|||
|
|
Filter -->|"匹配的消息"| C2["Consumer B"]
|
|||
|
|
end
|
|||
|
|
|
|||
|
|
subgraph "Consumer 端过滤"
|
|||
|
|
Broker -->|"全部消息"| C3["Consumer C"]
|
|||
|
|
C3 -->|"客户端过滤"| Logic["业务逻辑过滤"]
|
|||
|
|
end
|
|||
|
|
|
|||
|
|
style Producer fill:#4A90D9,color:#fff
|
|||
|
|
style Broker fill:#F5A623,color:#fff
|
|||
|
|
style Filter fill:#D0021B,color:#fff
|
|||
|
|
style C1 fill:#6EC1E0,color:#fff
|
|||
|
|
style C2 fill:#6EC1E0,color:#fff
|
|||
|
|
style C3 fill:#6EC1E0,color:#fff
|
|||
|
|
style Logic fill:#6EC1E0,color:#fff
|
|||
|
|
```
|
|||
|
|
|
|||
|
|
> [!question] 思考
|
|||
|
|
> Broker 端过滤可以减少网络传输,但会增加 Broker 负载。如何权衡?
|
|||
|
|
|
|||
|
|
关键在于过滤的"性价比"——如果过滤能淘汰 90% 的消息,那 Broker 多花一点 CPU 做过滤完全值得,因为省下的网络 IO 和 Consumer 处理时间远大于过滤开销。反过来,如果过滤只能淘汰 10% 的消息,不如在 Consumer 端过滤,把 Broker 的 CPU 留给更重要的事(如存储、复制)。实际生产中,Tag 过滤的性价比通常很高,因为一个 Tag 就能精准划分消息类型。
|
|||
|
|
|
|||
|
|
### 过滤方式一:Tag 过滤
|
|||
|
|
|
|||
|
|
RocketMQ 原生支持的最简单过滤方式。Producer 发送消息时指定 Tag,Consumer 订阅时用 Tag 表达式过滤。
|
|||
|
|
|
|||
|
|
```go
|
|||
|
|
// Producer: 发送带 Tag 的消息
|
|||
|
|
msg := NewMessage("Topic_Order", []byte(orderJSON))
|
|||
|
|
msg.SetTags("PAY") // 设置 Tag 为 PAY
|
|||
|
|
producer.Send(msg)
|
|||
|
|
|
|||
|
|
// Consumer: 只订阅 PAY 和 CANCEL 标签
|
|||
|
|
consumer.Subscribe("Topic_Order", "PAY || CANCEL")
|
|||
|
|
```
|
|||
|
|
|
|||
|
|
Tag 过滤发生在 Broker 端(ConsumeQueue 中存储了 Tag 的 hash 值),匹配效率很高。缺点是过滤粒度粗——只能按 Tag 精确匹配,不支持 `>`, `<`, `IN` 等复杂条件。
|
|||
|
|
|
|||
|
|
Tag 的底层实现很巧妙:ConsumeQueue 每条记录有 8 字节存储 Tag 的 hashcode,Broker 过滤时直接比较 hashcode,命中后再精确匹配 Tag 字符串,避免了解析消息体的开销。
|
|||
|
|
|
|||
|
|
### 过滤方式二:SQL92 表达式过滤
|
|||
|
|
|
|||
|
|
RocketMQ 支持基于 SQL92 子集的表达式过滤,功能比 Tag 强大得多。消息通过 `UserProperty` 设置自定义属性,Consumer 用 SQL 表达式过滤。
|
|||
|
|
|
|||
|
|
```go
|
|||
|
|
// Producer: 设置自定义属性
|
|||
|
|
msg := NewMessage("Topic_Order", []byte(orderJSON))
|
|||
|
|
msg.SetTags("PAY")
|
|||
|
|
msg.PutProperty("amount", "99.9")
|
|||
|
|
msg.PutProperty("region", "CN")
|
|||
|
|
producer.Send(msg)
|
|||
|
|
|
|||
|
|
// Consumer: SQL92 表达式过滤
|
|||
|
|
// 支持 AND、OR、IN、BETWEEN、IS NULL、比较运算符
|
|||
|
|
consumer.Subscribe("Topic_Order", "amount > 50 AND region IN ('CN','US')")
|
|||
|
|
```
|
|||
|
|
|
|||
|
|
SQL92 过滤同样在 Broker 端执行,Broker 会编译 SQL 表达式为语法树,对每条消息的属性进行求值。需要注意的是,SQL 过滤比 Tag 过滤消耗更多 Broker CPU,高吞吐场景下要谨慎使用。
|
|||
|
|
|
|||
|
|
支持的运算符和函数:
|
|||
|
|
- 比较:`>`, `<`, `>=`, `<=`, `=`, `<>`, `BETWEEN`
|
|||
|
|
- 逻辑:`AND`, `OR`, `NOT`
|
|||
|
|
- 集合:`IN`
|
|||
|
|
- 空值:`IS NULL`, `IS NOT NULL`
|
|||
|
|
- 字符串:`LIKE`(仅支持 `%` 通配符)
|
|||
|
|
|
|||
|
|
### 过滤方式三:Header 属性过滤
|
|||
|
|
|
|||
|
|
RabbitMQ 的 Headers Exchange 通过消息 Header 属性进行路由匹配,不依赖 Routing Key。
|
|||
|
|
|
|||
|
|
```go
|
|||
|
|
// RabbitMQ Headers Exchange 示例
|
|||
|
|
// Producer: 发送带 Header 的消息
|
|||
|
|
ch.Publish("orders_exchange", "", false, false, amqp.Publishing{
|
|||
|
|
Headers: amqp.Table{
|
|||
|
|
"x-match": "all", // all = 全部匹配, any = 任一匹配
|
|||
|
|
"type": "payment",
|
|||
|
|
"region": "CN",
|
|||
|
|
"amount": 99,
|
|||
|
|
},
|
|||
|
|
Body: orderJSON,
|
|||
|
|
})
|
|||
|
|
|
|||
|
|
// Consumer: 绑定队列时指定 Header 匹配规则
|
|||
|
|
ch.QueueBind("payment_cn_queue", "", "orders_exchange", false, amqp.Table{
|
|||
|
|
"x-match": "all",
|
|||
|
|
"type": "payment",
|
|||
|
|
"region": "CN",
|
|||
|
|
})
|
|||
|
|
```
|
|||
|
|
|
|||
|
|
Headers Exchange 的优势是路由规则完全由消息的 Header 决定,不需要像 Topic Exchange 那样设计 Routing Key 的层级结构。缺点是性能比 Direct/Topic Exchange 略差,因为需要逐个比较 Header 字段。
|
|||
|
|
|
|||
|
|
### 消息路由策略
|
|||
|
|
|
|||
|
|
路由决定消息从 Producer 到 Consumer 的流转路径,常见策略包括:
|
|||
|
|
|
|||
|
|
**Topic 路由**:最基础的路由方式,消息按 Topic 分发,每个 Topic 内按 Queue/Partition 分配。大部分场景下 Topic 路由就够用了。
|
|||
|
|
|
|||
|
|
**自定义路由规则**:当 Topic 粒度不够时,可以通过消息属性 + 过滤规则实现更精细的路由。比如同一个 Topic 下,按 `region` 属性路由到不同地域的 Consumer。
|
|||
|
|
|
|||
|
|
**消息再投递**:当 Consumer 处理失败或需要将消息转发到另一个 Topic 时,消息再投递(Consume-RePublish)模式就派上用场了。常见于消息转换、错误重试、死信转发等场景。
|
|||
|
|
|
|||
|
|
```go
|
|||
|
|
// 消息再投递示例:消费失败时转发到重试 Topic
|
|||
|
|
func handle(msg *MessageExt) ConsumeResult {
|
|||
|
|
err := process(msg)
|
|||
|
|
if err != nil {
|
|||
|
|
// 计算重试次数
|
|||
|
|
retryCount := msg.GetReconsumeTimes()
|
|||
|
|
if retryCount >= 3 {
|
|||
|
|
// 超过重试次数,投递到死信队列
|
|||
|
|
producer.Send("DLQ_Order", msg.Body)
|
|||
|
|
return ConsumeSuccess
|
|||
|
|
}
|
|||
|
|
// 重试:消息会自动投递到 %RETRY% Topic
|
|||
|
|
return ReconsumeLater
|
|||
|
|
}
|
|||
|
|
return ConsumeSuccess
|
|||
|
|
}
|
|||
|
|
```
|
|||
|
|
|
|||
|
|
### 各过滤方式对比
|
|||
|
|
|
|||
|
|
| 维度 | Tag 过滤 | SQL92 过滤 | Header 过滤 |
|
|||
|
|
|------|---------|-----------|-------------|
|
|||
|
|
| 过滤位置 | Broker 端 | Broker 端 | Broker 端 |
|
|||
|
|
| 过滤粒度 | 精确匹配 | 条件表达式 | 键值对匹配 |
|
|||
|
|
| 性能 | 极高 | 中 | 中 |
|
|||
|
|
| 灵活性 | 低 | 高 | 中 |
|
|||
|
|
| 适用场景 | 简单消息分类 | 复杂业务规则 | AMQP 路由 |
|
|||
|
|
| 典型 MQ | RocketMQ | RocketMQ | RabbitMQ |
|
|||
|
|
|
|||
|
|
## 关联笔记
|
|||
|
|
|
|||
|
|
- [[06-高级特性/17-MQ-延迟消息与定时消息|MQ 延迟消息与定时消息]]
|
|||
|
|
- [[06-高级特性/18-MQ-事务消息|MQ 事务消息]]
|
|||
|
|
- [[02-消息模型/3-MQ-消息模型|MQ 消息模型]]
|
|||
|
|
- [[03-协议与标准/5-AMQP-协议|AMQP 协议]]
|
|||
|
|
- [[07-主流MQ对比/24-RocketMQ|RocketMQ]]
|
|||
|
|
- [[07-主流MQ对比/23-RabbitMQ|RabbitMQ]]
|