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 "情报发布成功,请查看控制台日志!";
    }
}

  1. 启动 Spring Boot 项目。
  2. 浏览器访问: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. 生产环境避坑指南

  1. 消息丢失防范
    • 发布端:开启 publisher-confirms,确保消息成功到达 Exchange。
    • 消费端:关闭自动 ACK,采用手动 ACKAcknowledgeMode.MANUAL)。推送成功后再 basicAck,如果推送服务挂了则 basicNack 重新入队。
  2. 避免重复推送(幂等性)
    • 如果一篇文章同时带了 [Java][后端] 标签,而某个用户同时订阅了这两个标签。那么该用户会被查出来两次,收到两条同样的推送。
    • 解决方案:在 Consumer 端,可以借助 Redis 做一个去重校验。Key 为 push:{intelligenceId}:{userId},如果存在则说明已推送过,直接跳过。
  3. 大 V 效应(海量订阅用户如何处理)
    • 如果某个标签(如 [日常公告])有 100 万人订阅,Consumer 查库和循环推送会耗时很久。
    • 解决方案:Consumer 只做“分片”工作。比如把 100 万用户分成 1000 个批次,每批次 1000 人,将这 1000 个批量推送任务再次作为消息发入 batch.push.queue 中,利用多台机器并发推送。

5. 总结 (Conclusion)

一句话总结这套架构的优势:异步解耦提升了接口响应速度,RabbitMQ 路由机制让标签匹配变得极其优雅,整体架构具备良好的横向扩展能力。


Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐