如何用 RabbitMQ 写一个高性能的“标签订阅与推送”系统?
1. 场景引入
主播之前再做威胁情报的交易平台,在很多资讯类、威胁情报类或社区类产品中,我们经常遇到这样的需求:
- 用户 A 关注了标签
[Web安全]和[Java]。 - 用户 B 关注了标签
[Python]和[AI]。 - 当有运营人员或创作者发布了一篇带有
[Java]标签的情报/文章时,系统需要立刻且准确地给用户 A 发送站内信或 App 推送。
传统做法的痛点:如果每次发布文章都去数据库里 SELECT 查出所有关注该标签的用户,然后在一个线程里循环发通知,一旦用户量上来,接口响应会极慢,甚至导致请求超时。
破局方案:引入 RabbitMQ,利用其强大的路由(Routing)和发布/订阅(Pub/Sub)机制,将“发布情报”与“推送通知”彻底解耦!
2. 核心架构设计
在这里,我们采用 RabbitMQ 的 Topic Exchange(主题交换机) 来实现。
为什么选 Topic? 因为情报可能包含多个标签,且后续可能扩展出更复杂的订阅规则(比如不仅按标签,还按情报级别划分:info.level.high.tag.java),Topic 的通配符匹配最灵活。

为什么不给每个用户建一个 Queue?
理论上,可以让用户 A 拥有一个独立的 Queue,并绑定到他关注的标签 RoutingKey 上。这在 RabbitMQ 层面是最完美的 Pub/Sub。但是,在互联网 ToC 场景下(动辄百万用户),创建百万个 Queue 会导致 MQ 内存溢出。因此,我们采用**“统一处理队列 + 消费端查库动态分发”**的方案,兼顾了性能与可维护性。
3. Demo实验
环境准备
- JDK 8 或以上
- 本地安装并启动了 RabbitMQ(默认端口 5672)
- Spring Boot 2.x 或 3.x
第一步:引入 Maven 依赖 pom.xml
只需要最基础的 Web 和 RabbitMQ 依赖即可。
<dependencies>
<!-- Web 用于提供测试接口 -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
<!-- RabbitMQ 核心依赖 -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>
</dependencies>
第二步:配置文件 application.yml
配置 RabbitMQ 的连接信息。
spring:
rabbitmq:
host: 127.0.0.1
port: 5672
username: guest
password: guest
server:
port: 8080
第三步:RabbitMQ 核心配置类 RabbitMQConfig.java
在这里定义 Topic 交换机和队列,并完成绑定。
import org.springframework.amqp.core.*;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@Configuration
public class RabbitMQConfig {
// 交换机名称
public static final String EXCHANGE_NAME = "intel.topic.exchange";
// 推送队列名称
public static final String PUSH_QUEUE_NAME = "intel.push.queue";
// 路由键前缀匹配
public static final String ROUTING_KEY_PATTERN = "tag.#";
@Bean
public TopicExchange intelExchange() {
// 创建 Topic 交换机
return new TopicExchange(EXCHANGE_NAME);
}
@Bean
public Queue pushQueue() {
// 创建队列,true 表示持久化
return new Queue(PUSH_QUEUE_NAME, true);
}
@Bean
public Binding bindingPushQueue(Queue pushQueue, TopicExchange intelExchange) {
// 将队列绑定到交换机,监听所有以 "tag." 开头的路由键
return BindingBuilder.bind(pushQueue).to(intelExchange).with(ROUTING_KEY_PATTERN);
}
}
第四步:模拟订阅关系与实体类 MockDataService.java
为了让 Demo 能跑起来,我们模拟一个数据库,里面存了用户的订阅标签。
import org.springframework.stereotype.Service;
import java.util.*;
import java.util.stream.Collectors;
@Service
public class MockDataService {
// 模拟数据库:存储 用户ID 和 他订阅的标签列表
private static final Map<Long, List<String>> USER_SUBSCRIPTIONS = new HashMap<>();
static {
// 用户 1 订阅了 Java 和 后端
USER_SUBSCRIPTIONS.put(1L, Arrays.asList("Java", "后端"));
// 用户 2 订阅了 Python 和 AI
USER_SUBSCRIPTIONS.put(2L, Arrays.asList("Python", "AI"));
// 用户 3 仅订阅了 Java
USER_SUBSCRIPTIONS.put(3L, Collections.singletonList("Java"));
}
/**
* 模拟:根据标签查询订阅了该标签的用户ID集合
*/
public List<Long> findUsersByTag(String tag) {
return USER_SUBSCRIPTIONS.entrySet().stream()
.filter(entry -> entry.getValue().contains(tag))
.map(Map.Entry::getKey)
.collect(Collectors.toList());
}
}
第五步:情报发布者(生产者) IntelPublisher.java
发布情报时,拆分标签,并发送 MQ 消息。
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
@Service
public class IntelPublisher {
@Autowired
private RabbitTemplate rabbitTemplate;
public void publish(String intelId, String title, List<String> tags) {
System.out.println("=== 运营发布新情报: [" + title + "] ===");
// 1. 正常业务逻辑:保存情报到数据库...
// 2. 遍历标签,触发 MQ 消息
for (String tag : tags) {
// 拼装 RoutingKey,例如:tag.java
String routingKey = "tag." + tag.toLowerCase();
// 构建消息体
Map<String, Object> message = new HashMap<>();
message.put("intelId", intelId);
message.put("title", title);
message.put("triggerTag", tag);
// 发送消息到 Topic Exchange
rabbitTemplate.convertAndSend(RabbitMQConfig.EXCHANGE_NAME, routingKey, message);
System.out.println("[发布者] 已发送 MQ 消息,路由键: " + routingKey);
}
}
}
第六步:推送消费者(防重复推送) PushConsumer.java
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
@Component
public class PushConsumer {
@Autowired
private MockDataService mockDataService;
// 本地缓存,模拟 Redis 记录:防止同一篇文章向同一个用户推送多次
// Key 格式: {intelId}_{userId}
private final Set<String> pushedRecords = ConcurrentHashMap.newKeySet();
@RabbitListener(queues = RabbitMQConfig.PUSH_QUEUE_NAME)
public void receivePushTask(Map<String, Object> message) {
String intelId = (String) message.get("intelId");
String title = (String) message.get("title");
String triggerTag = (String) message.get("triggerTag");
System.out.println("\n[消费者] 收到推送任务,触发标签: " + triggerTag);
// 1. 查库:获取订阅了该标签的用户
List<Long> userIds = mockDataService.findUsersByTag(triggerTag);
if (userIds.isEmpty()) {
System.out.println(" -> 没有用户订阅该标签,忽略推送");
return;
}
// 2. 遍历用户,执行推送
for (Long userId : userIds) {
// == 核心防重逻辑 ==
// 如果一篇文章有[Java]和[后端],用户A都订阅了,防止给用户A发两次
String deduplicationKey = intelId + "_" + userId;
// 如果 add 成功,说明之前没推送过;如果 add 返回 false,说明已推送过
if (pushedRecords.add(deduplicationKey)) {
System.out.println(" [推送成功] 给用户 " + userId + " 发送 App 通知: 您关注的 [" + triggerTag + "] 有更新: " + title);
} else {
System.out.println(" 拦截 [重复推送] 给用户 " + userId + " (原因:已通过其他标签推送过该文章)");
}
}
}
}
第七步:测试接口 TestController.java
提供一个 HTTP 接口,方便你在浏览器或者 Postman 中触发测试。
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RestController;
import java.util.Arrays;
import java.util.UUID;
@RestController
public class TestController {
@Autowired
private IntelPublisher intelPublisher;
@GetMapping("/test/publish")
public String testPublish() {
String intelId = UUID.randomUUID().toString();
// 模拟发布一篇情报,包含两个标签:"Java" 和 "后端"
// 根据我们的 Mock 数据,用户1 恰好同时订阅了这两个标签!可以完美测试防重逻辑。
intelPublisher.publish(
intelId,
"Spring Boot 整合 RabbitMQ 最佳实践",
Arrays.asList("Java", "后端")
);
return "情报发布成功,请查看控制台日志!";
}
}
- 启动 Spring Boot 项目。
- 浏览器访问:
http://localhost:8080/test/publish
[发布者] 已发送 MQ 消息,路由键: tag.java
[发布者] 已发送 MQ 消息,路由键: tag.后端
[消费者] 收到推送任务,触发标签: Java
[推送成功] 给用户 1 发送 App 通知: 您关注的 [Java] 有更新: Spring Boot 整合 RabbitMQ 最佳实践
[推送成功] 给用户 3 发送 App 通知: 您关注的 [Java] 有更新: Spring Boot 整合 RabbitMQ 最佳实践
[消费者] 收到推送任务,触发标签: 后端
拦截 [重复推送] 给用户 1 (原因:已通过其他标签推送过该文章)
4. 生产环境避坑指南
- 消息丢失防范:
- 发布端:开启
publisher-confirms,确保消息成功到达 Exchange。 - 消费端:关闭自动 ACK,采用手动 ACK(
AcknowledgeMode.MANUAL)。推送成功后再basicAck,如果推送服务挂了则basicNack重新入队。
- 发布端:开启
- 避免重复推送(幂等性):
- 如果一篇文章同时带了
[Java]和[后端]标签,而某个用户同时订阅了这两个标签。那么该用户会被查出来两次,收到两条同样的推送。 - 解决方案:在 Consumer 端,可以借助 Redis 做一个去重校验。Key 为
push:{intelligenceId}:{userId},如果存在则说明已推送过,直接跳过。
- 如果一篇文章同时带了
- 大 V 效应(海量订阅用户如何处理):
- 如果某个标签(如
[日常公告])有 100 万人订阅,Consumer 查库和循环推送会耗时很久。 - 解决方案:Consumer 只做“分片”工作。比如把 100 万用户分成 1000 个批次,每批次 1000 人,将这 1000 个批量推送任务再次作为消息发入
batch.push.queue中,利用多台机器并发推送。
- 如果某个标签(如
5. 总结 (Conclusion)
一句话总结这套架构的优势:异步解耦提升了接口响应速度,RabbitMQ 路由机制让标签匹配变得极其优雅,整体架构具备良好的横向扩展能力。
更多推荐

所有评论(0)