# 异步任务处理-RabbitMQ+WebSocket 负载均衡 **Repository Path**: hnu-turing/asynchronous-update ## Basic Information - **Project Name**: 异步任务处理-RabbitMQ+WebSocket 负载均衡 - **Description**: 异步任务处理 Demo:客户端 → RabbitMQ → 独立消费者集群 → WebSocket 实时通报 - **Primary Language**: Unknown - **License**: Not specified - **Default Branch**: master - **Homepage**: None - **GVP Project**: No ## Statistics - **Stars**: 5 - **Forks**: 0 - **Created**: 2026-08-10 - **Last Updated**: 2026-08-31 ## Categories & Tags **Categories**: Uncategorized **Tags**: None ## README # 异步任务处理 Demo:客户端 → RabbitMQ → 独立消费者集群 → WebSocket 实时通报 演示**工作队列(Work Queue)+ 负载均衡 + 异步回执**的经典模式: - 多个客户端通过 Web 应用提交任务(HTTP) - 任务进入 RabbitMQ **共享工作队列**,由**多个独立消费者进程**竞争消费——谁空闲谁拿(`prefetch=1` 公平分发) - 消费者可多实例**动态扩缩容**,处理过程中按进度回报 - 进度/结果经 RabbitMQ 通知扇出到 Web 应用各实例,由 WebSocket 推回对应客户端 - WebSocket 会话状态统一存 Redis(跨实例可见) ``` 客户端1 ─┐ ├─HTTP 提交任务→ Web应用 →┐ 客户端2 ─┘ ▼ mqws.task.queue(共享工作队列) │ 竞争消费(prefetch=1) ┌─────────────┼─────────────┐ ▼ ▼ ▼ consumer-1 consumer-2 consumer-3 ← 独立进程,可动态增减 │ 处理中按步骤回报进度 ▼ mqws.notify.topic(扇出) ▼ Web应用各实例消费 → 查本机 WebSocket 连接 → 推回客户端 ▲ 会话状态:Redis(谁在线/挂在哪个实例) ``` ## 模块结构(Maven 多模块) ``` mq-websocket-demo(聚合父工程) ├── common/ # 消息契约:TaskMessage / NotifyMessage / 常量(三方共用) ├── web-app/ # Web 应用:任务提交接口 + 通知推送(依赖 Redis + RabbitMQ) └── task-consumer/ # 独立消费者应用(只依赖 RabbitMQ),跑几个实例就是几个消费者 ``` ## 快速开始 1. 启动中间件:`docker-compose up -d`(RabbitMQ + Redis) 2. 启动 Web 应用:`mvn spring-boot:run -pl web-app` 3. 启动消费者(想开几个开几个,可中途增减): ```bash java -jar task-consumer/target/mqws-task-consumer.jar --app.instance-id=worker-a java -jar task-consumer/target/mqws-task-consumer.jar --app.instance-id=worker-b ``` 4. 打包(根目录执行一次即可):`mvn package` 5. 打开 `http://localhost:8080/`: - 每开一个标签页连一个 userId → 模拟多个客户端 - 用「提交任务」卡片连发多个任务(可设不同耗时)→ 观察进度通知来自哪个 worker - 消息里带 `由 worker-x 处理`,直观看到任务在消费者间的分配 ## 观察点 - **负载均衡**:RabbitMQ `prefetch=1` + 单线程消费 → 消息只发给空闲消费者;日志 `[任务-阶段3-领取]` 打在谁那里,活就是谁干的 - **动态扩展**:演示中途再启动 worker-c,新任务立刻会被它分到;停掉一个 worker,剩余任务自动由别的实例接手 - **负载均衡 ≠ 随机分配**:同耗时任务趋于轮询,不同耗时自动偏向先完成的实例 - **RabbitMQ 控制台**:`http://localhost:15672`(guest/guest) - `mqws.task.queue`:unacked/ready 数量变化可见任务排队与消费 - `spring.gen-*`:web 应用的通知扇出匿名队列 - **Redis**:会话状态三本账(见下) ## Redis 键结构 | 键 | 类型 | 内容 | | --- | --- | --- | | `mqws:online:users` | Set | 所有在线 userId | | `mqws:user:{userId}:sessions` | Set | 该用户的所有 sessionId | | `mqws:session:{sessionId}` | String | `userId\|instanceId` | ## 接口速查 | 方法 | 路径 | 说明 | | --- | --- | --- | | WebSocket | `/ws/notify/{userId}` | 客户端长连接 | | POST | `/api/tasks` | 提交任务 `{"userId":"user-1","taskName":"清洗数据","durationMs":8000}`(userId 为 null = 结果广播) | | GET | `/api/tasks/send?taskName=&durationMs=&userId=` | 浏览器地址栏快速提交 | | POST | `/api/messages` | 直接发通知消息 `{"userId":"user-1","content":"..."}` | | GET | `/api/messages/send?userId=&content=` | 快速发通知 | | GET | `/api/messages/online-count` | 在线用户数(Redis)+ 本实例 ID | ## 日志跟踪(按关键字过滤) | 日志标记 | 位置 | 含义 | | --- | --- | --- | | `[任务-阶段1-HTTP]` | web-app | 收到任务提交 → 202 返回 taskId | | `[任务-阶段2-MQ发送]` | web-app | 任务投入 `mqws.task.queue` | | `[任务-阶段3-领取]` | task-consumer | 某 worker 抢到任务 | | `[任务-阶段4-处理]` | task-consumer | 逐步处理(进度回报) | | `[任务-完成]` | task-consumer | 任务结束,发完成通知 | | `[阶段3-MQ消费]` `[阶段4-本机路由]` `[阶段5-WS送达]` | web-app | 通知扇出 → 查本机连接 → 推送 | | `[WS-连接]` `[Redis-会话注册/注销]` | web-app | WebSocket 会话生命周期 | ## 冒烟验证 中间件就绪后,起 web-app + 至少 2 个消费者,然后: ```bash node scripts/smoke-test.cjs # 默认打 http://localhost:8080 $env:SMOKE_HOST="localhost:18080" ; node scripts/smoke-test.cjs ``` 断言覆盖:定向通知、4 任务全部完成、**至少 2 个不同 worker 参与处理(负载均衡)**、定向回执不串线、在线数统计。 ## pom 说明 - Spring Boot 4.0.6(Jackson 3,`spring-boot-starter-webmvc`)+ JDK 17;BOM 直引、插件显式固定版本。 - 父工程声明了 id 为 `public` 的仓库(URL 为中央仓库):本机 `~/.m2` 缓存来源于同名仓库;离线构建用 `mvn -o package`。 - 消费者模块依赖 `spring-boot-starter-jackson`:`JacksonJsonMessageConverter` 是 Jackson 3 实现,纯 AMQP starter 不含 Jackson。 ## 配套文档 - [docs/演讲稿.md](docs/演讲稿.md):30 分钟演示讲稿(含架构图 / 流程图) - [docs/面试八股文.md](docs/面试八股文.md):29 题递进式面试 Q&A——**场景八股(10 个场景全围绕本项目演示动作展开)→ 传统八股(Q1-Q19 递进深水区)**,场景结尾带锚点主动引导面试官;另附**简历职责描述**(简历版 5 条 + 面试口述版 + 量化口径)