为什么选择 SseIm Framework?
基于 HTTP + SSE 架构设计思想研发的即时通讯框架。3 行 YAML 配置,零 Java 代码即可获得实时消息推送能力。
SseIm Framework 是即时通讯的基础设施,专注于消息的实时推送。它解决"如何把消息从服务端实时推送到客户端",提供可插拔的 MQ 桥接和 SSE 推送能力。
配置驱动
3 行 YAML 配置 = 端点注册 + Stream 绑定 + SSE 连接管理。零 Java 代码开箱即用。
三段式解耦
进(Publisher)→ 转(Repeater)→ 出(Subscriber),每层职责单一、可独立替换,架构清晰可扩展。
双通道投递
广播(topicSink,O(1))+ 定向(routingKeyIndex + directedSink,O(K)),统一两种投递模式。
三层清理
doFinally(95%)+ Reaper(4.9%)+ maxLifetime(0.1%),保证连接不泄露。
核心架构
框架由核心组件协同工作,实现三段式解耦的消息流水线:
EndpointRegistry
端点注册表。存储 path → topic + mode 映射,支持运行时动态增删。
EndpointRouter
动态路由器。RouterFunction 实现,按 mode 分发 PUSH / PULL 请求,织入 AccessFilter 链。
TopicRegistry
Topic 组件注册表。通过工厂模式按需创建 Publisher / Repeater / Subscriber。
Publisher
消息发布器(进)。接收 POST 请求,构建 TopicMessage 经 Repeater 发送到 MQ。
Repeater
消息转发器(转)。桥接 MQ 两端,send → StreamBridge,receive → 拦截器链 → SSE 推送。
Subscriber
消息订阅器(出)。构建带生命周期钩子的 Reactor Flux,交付 SseConnectionManager 推送。
SseConnectionManager
SSE 连接管理器。连接注册/注销、背压策略、心跳保活、定时清理、优雅关闭。
SseConnectionRegistry
连接存储与索引。管理双通道架构(topicSink + directedSink),维护 routingKeyIndex 反向索引。
TopicLifecycleManager
Topic 生命周期管理。状态机(CREATED → ACTIVE → DESTROYING → DESTROYED)、引用计数、自动清理空闲 Topic。
ConnectionReaper
连接收割器。三层清理的第二层,定时扫描僵尸连接(30s),异步强制关闭,防止资源泄露。
RoutingKeyComposer
路由键组合策略。双通道投递的核心,支持广播和按 routingKey 定向投递,可自定义路由规则。
可插拔指标
Micrometer 指标自动暴露,连接数 / 订阅速率 / 发布速率 / 定时清理,无缝接入 Prometheus。
POST /api/chat/push → EndpointRouter → Publisher.push() → Repeater.send() → StreamBridge → MQ → Repeater.receive() → 拦截器链 → SseConnectionManager.publish() → Sink → GET /api/chat/pull SSE 客户端
设计哲学
声明式优于编程式
开发者说"我要什么",框架决定"怎么做"。配置驱动,人类负责意图,机器负责执行。
可组合优于可配置
小组件通过组合产生无限可能,而非通过配置膨胀。Factory + Interceptor + Filter 自由组合。
渐进式复杂度披露
零代码 → 拦截器 → 工厂 → 完全自定义。入门门槛低,天花板高,按需展开。
快速开始
三步集成 SseIm Framework 到你的 Spring Boot 应用,实现实时消息推送。零 Java 代码。
JDK 25+、Spring Boot 4.x、RabbitMQ(或其他 Spring Cloud Stream 支持的 Binder)
添加 Maven 依赖
<dependency>
<groupId>cn.gmlee.tools</groupId>
<artifactId>tools-im</artifactId>
<version>5.6.0-SNAPSHOT</version>
</dependency>
<dependency>
<groupId>cn.gmlee.tools</groupId>
<artifactId>tools-base</artifactId>
<version>5.6.0-SNAPSHOT</version>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-webflux</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.cloud</groupId>
<artifactId>spring-cloud-stream-binder-rabbit</artifactId>
</dependency>
配置 application.yml
3 行配置定义端点,框架自动完成端点注册、Stream 绑定、SSE 连接管理:
im:
endpoints:
- path: /api/chat/push
topic: im.chat
mode: push # POST 接收消息 → 发送到 MQ
- path: /api/chat/pull
topic: im.chat
mode: pull # GET 订阅 SSE 流
spring:
rabbitmq:
host: localhost
port: 5672
零 Java 代码。框架自动创建 RabbitMQ binding、注册端点路由、管理 SSE 连接。
客户端使用
// 订阅消息流(SSE)— 通过 URL 参数声明租户和房间
const es = new EventSource('/api/chat/pull?tenant=acme&room=lobby');
es.onmessage = e => console.log(JSON.parse(e.data));
// 广播消息(无 URL 参数 → 所有连接收到)
fetch('/api/chat/push', {
method: 'POST',
headers: {'Content-Type': 'application/json'},
body: JSON.stringify({from: 'alice', content: 'hello'})
});
// 定向投递 — 使用与订阅方相同的参数格式指定目标房间
fetch('/api/chat/push?tenant=acme&room=lobby', {
method: 'POST',
headers: {'Content-Type': 'application/json'},
body: JSON.stringify({from: 'alice', content: 'msg to acme/lobby'})
});
数据流
起源背景
从真实生产故障中提炼根因,针对性地选择更合适的技术方案。
从生产问题出发
SseIm Framework 不是"为了用新技术而换架构",而是从两个真实的生产故障中提炼出了根因,然后针对性地选择了更合适的技术方案。
竞价板块:连接泄露
原架构 WebSocket + STOMP + RabbitMQ。客户端连接管理不当,导致连接泄露 — 连接数只增不减,直到内存溢出,引发生产故障。
站内信板块:消息卡顿/超时
原架构同样是 WebSocket + STOMP + RabbitMQ。消息发送非常卡顿、缓慢,甚至超时,用户体验极差。
两个板块的问题表象不同,但根因指向同一个方向:WebSocket 的连接模型不适合这些场景。
WebSocket 的两个经典陷阱
陷阱一:连接泄露(竞价板块的根因)
WebSocket + STOMP 的连接模型是有状态长连接,服务端需要自己管理连接的生命周期。一旦客户端异常断开(网络波动、页面崩溃、移动端切换后台),如果清理逻辑有任何遗漏,连接就会残留在内存里 — 只增不减,直到 OOM。
这不是个别问题,WebSocket 连接泄露是业界非常常见的生产事故根因。
陷阱二:协议过重(站内信板块的根因)
STOMP 的消息投递模型是双向的,但站内信场景本质上是服务端单向推送。用 WebSocket + STOMP 做单向推送,是在用最重的协议做最简单的事,不必要的协议开销和连接管理复杂度都会变成性能瓶颈。
核心洞察:IM 不需要"双向长连接"
虽然 SSE 是单向通道,但大部分即时通讯场景其实只需要单向通道。因为消息发送是用户实时触发的离散动作,只需要能发出去即可,HTTP 才是简单又可靠的方案。
大多数人对 IM 的直觉是"需要双向通道",所以第一反应就是 WebSocket。但仔细拆解用户行为:
| 方向 | 行为特征 | 本质 | 最优协议 |
|---|---|---|---|
| 客户端 → 服务端(发消息) | 用户点"发送"按钮 → 一次 HTTP POST → 结束 | 离散动作,用完即走 | HTTP POST |
| 服务端 → 客户端(收消息) | 持续等待服务端主动推送 | 持续等待,需要长连接 | SSE |
两个方向用各自最合适的协议,而不是用 WebSocket 把两个方向绑在同一条连接上。
WebSocket 的问题在于"过度对称"
WebSocket 提供的是对称的双向长连接,但 IM 的两个方向根本不是对称的:
- 发消息是一次性的 HTTP 请求,不需要长连接
- 收消息需要持续等待,需要长连接
用 WebSocket 做 IM,相当于为了"服务端能推送"这个需求,被迫把"客户端发消息"也绑进了一条长连接里。然后就要自己处理:
- 这条连接上的消息协议(STOMP 帧解析)
- 这条连接的状态管理(心跳、重连、断线恢复)
- 这条连接的生命周期(谁负责关闭、怎么清理)
复杂度全是你自己加的,而 HTTP + SSE 方案里,这些复杂度本来就不存在。
SseIm 的解法:按需分配连接
- 需要推送时 → 建立 SSE 连接(长连接)
- 需要发消息时 → 发一个 HTTP 请求(用完即释放)
- 不需要推送时 → 没有任何连接占用
- 永远占着一条双向连接
- 发消息、收消息、心跳全在同一连接上
- 连接管理复杂度全部自己承担
原方案 vs 新方案对比
| 维度 | WebSocket + STOMP | HTTP + SSE |
|---|---|---|
| 连接模型 | 有状态双向长连接,服务端维护连接状态机 | HTTP 短连接 + SSE 单向长连接,HTTP 层天然感知断连 |
| 连接泄露风险 | 高 — 依赖应用层心跳探测,清理逻辑遗漏即泄露 | 低 — 三层清理机制兜底(doFinally → Reaper → maxLifetime) |
| 协议开销 | 重 — STOMP 帧解析、双向握手、连接状态维护 | 轻 — HTTP + SSE 原生支持,零额外协议开销 |
| 消息延迟 | 高 — 协议栈层次多,处理链路长 | 低 — 最短消息路径,直接 SSE 推送 |
| 客户端依赖 | 需要 STOMP 客户端 SDK | 浏览器原生 EventSource,零依赖 |
| 扩展性 | 直接绑定 RabbitMQ | Spring Cloud Stream 抽象,可换任意 MQ |
不是在"选 SSE 还是 WebSocket",而是在想清楚消息推送的本质是单向的之后,做出的必然选择。用 Spring Cloud Stream 而不是直接绑死某个 MQ,在解决当前问题的同时保留了架构的灵活性 — 以后换 MQ 不需要改业务代码。
设计路径
SseIm Framework 的诞生遵循一条清晰的路径:
生产故障暴露根因
竞价板块连接泄露(OOM)、站内信板块消息卡顿/超时。两个板块的问题都指向 WebSocket + STOMP 的架构缺陷。
架构思想形成
IM 的本质是单向推送,不需要双向长连接。HTTP + SSE 才是正确的协议组合。先有对问题的本质理解和解决思路。
实验验证可行性
通过原型实验验证 SSE 推送、Spring Cloud Stream 桥接、连接管理等核心技术方案的可行性。
框架设计与实现
可行性确认后,将经过验证的架构思想沉淀为可复用的开发框架 — 配置驱动、进-转-出三段式、SPI 扩展、渐进式复杂度披露。
框架本身就是把架构认知产品化 — 让开发者不需要理解全部思考过程,只需要写 3 行 YAML,就能按照经过验证的架构模式来工作。这正是框架的价值所在。
架构设计
配置驱动 · 三段式解耦 · 零侵入扩展
设计哲学
🎯 声明式优于编程式
开发者表达"意图",框架负责"执行"。3 行 YAML 配置即可完成端点注册、Stream 绑定、连接管理。
🔄 三段式解耦
进(Publisher)- 转(Repeater)- 出(Subscriber),每层职责单一,可独立替换,组合出无限可能。
🧩 渐进式扩展
零代码开箱 → 拦截器扩展 → 工厂自定义 → 完全替换链路。入门门槛低,天花板高。
全景架构
从开发者配置到客户端接收的完整链路:
im.endpoints: [{path, topic, mode}]
GET /api/chat/pull (SSE)
关键设计决策
| 决策 | 原因 | 收益 |
|---|---|---|
| 路径与 Topic 解耦 | URL 是对外契约,Topic 是内部实现。两者独立演进,互不影响。 | URL 可以 RESTful 化、版本化;Topic 可以按业务语义命名 |
| 进-转-出三段式 | 职责分离,每层可独立替换。Publisher 专注接入,Repeater 专注转发,Subscriber 专注分发。 | 高内聚低耦合,支持细粒度自定义 |
| 响应式全链路 | WebFlux 环境下,同步阻塞会阻塞事件循环。全链路响应式保证高并发下的吞吐量。 | 非阻塞、高吞吐、资源利用率高 |
| SPI 扩展点 | 不同场景有不同的定制需求。Factory / Interceptor / Filter / Listener 覆盖主要扩展场景。 | 零侵入扩展,按需渐进式自定义 |
配置驱动原理
框架通过 im.endpoints 配置自动完成所有基础设施创建:
im:
endpoints:
- path: /api/chat/pull # URL 路径(对外契约)
topic: im.chat # 内部 Topic(业务语义)
mode: pull # PUSH 或 PULL
URL 路径是对外契约,可以随意设计(RESTful、版本化);Topic 是内部实现,可以随意命名。两者互不影响,独立演进。
自动创建链路
启动时,EndpointAutoConfiguration 自动完成以下工作:
| 阶段 | 组件 | 动作 |
|---|---|---|
| 1 | EndpointRegistry |
注册端点配置(path → topic + mode) |
| 2 | TopicFactory |
按需创建 Stream binding(输出或输入)+ Consumer Bean |
| 3 | TopicRegistry |
按需创建 Publisher / Repeater / Subscriber 组件 |
| 4 | EndpointRouter |
构建 RouterFunction,动态路由请求到 PUSH / PULL 处理器 |
性能分析
用数据说话,用架构解释 — 每个指标都有设计支撑
消息延迟
1-5ms
端到端延迟(发布到接收)
IM 场景用户无感知
内存效率
~240 bytes
单连接内存占用
100 万连接仅 ~240 MB
消息吞吐
10 万 msg/s
单线程发布吞吐
Reactor Sinks 高效扇出
连接扇出
10 万连接/ms
单条消息扇出速率
Reactor 内部并行分发
延迟性能
| 指标 | 数值 | 说明 |
|---|---|---|
| 端到端延迟 | 1-5ms | 发布到接收全链路,IM 场景用户完全无感知 |
| 广播扇出延迟 | 10 万连接/ms | 单条消息扇出速率,Reactor 内部并行分发 |
| 定向投递延迟 | ~0.01ms × K | K = 目标连接数,反向索引 O(1) 查找 |
| MQ 传输延迟 | 0.5-2ms | 消息中间件网络传输 |
| 跨节点延迟 | 1-5ms | 集群环境下的节点间通信 |
设计支撑:消息路径优化
~0.5-1ms
~0.6-2.5ms
~0.001ms
SSE 基于 HTTP,看似"重",但实际上:① 无需双向握手;② 无需心跳保活(HTTP 层天然感知断连);③ 无需帧解析。协议层面的差异被架构简化抵消。
延迟主要来源:① MQ 网络传输(0.5-2ms);② 跨节点延迟(1-5ms)。SSE 协议本身开销极小(~0.001ms),架构简化是关键。
内存效率
| 指标 | 数值 | 说明 |
|---|---|---|
| 单连接内存 | ~240 bytes | 包含连接对象、原子变量、索引开销 |
| 100 万连接内存 | ~240 MB | 框架开销,不含 JVM 和 GC 预留 |
| 对比 WebSocket | 节省 2-3 倍 | WebSocket 方案 ~500-600 bytes/连接 |
| Topic 元数据 | ~40 bytes/Topic | 保留计数器条目,避免竞态 |
| 广播额外开销 | 0 bytes | 零拷贝广播,消息引用传递 |
设计支撑:位域优化 + 零拷贝
| 优化技术 | 节省内存 | 实现方式 |
|---|---|---|
| 位域双标志 | ~16 bytes/连接 | 单 AtomicInteger 存储 2 个 AtomicBoolean |
| 原子长整型时间戳 | ~8 bytes/连接 | AtomicLong 替代 Instant 对象 |
| 零拷贝广播 | 0 bytes 额外开销 | Reactor topicSink 共享,消息引用传递 |
| 无副本迭代 | 0 bytes 额外开销 | Reaper 直接迭代 ConcurrentHashMap.values() |
| 空 Topic 延迟移除 | ~40 bytes/Topic | 保留计数器条目,避免与 computeIfAbsent 竞态 |
内存主要占用:① 连接对象(~240 bytes);② Topic 元数据(~40 bytes/Topic);③ JVM 对象头(~16 bytes/对象)。位域优化已接近极限,进一步提升需要自定义对象布局。
并发吞吐
| 指标 | 单线程 | 多线程(10 线程) | 说明 |
|---|---|---|---|
| 消息发布 | 10 万 msg/s | 3-5 万 msg/s | Reactor Sinks CAS 序列化 |
| 消息扇出 | 10 万连接/ms | Reactor 内部并行 | 零拷贝广播,O(1) 复杂度 |
| 定向投递 | O(K),K=目标数 | 受 routingKeyIndex 限制 | 反向索引 O(1) 查找 |
| 连接建立 | 100 连接/ms | CAS 无锁并发 | 原子计数器递增 |
| Reaper 扫描 | 10 万连接 ~0.5ms | 单线程执行 | 零内存分配迭代 |
设计支撑:CAS 无锁 + Reactor Sinks
| 并发组件 | 数据结构 | 并发特性 | 性能 |
|---|---|---|---|
| 连接存储 | ConcurrentHashMap<String, SseConnection> |
分段锁(bin lock) | O(1) 查找,高并发下 bin 竞争可控 |
| 路由键索引 | 嵌套 ConcurrentHashMap |
两层 O(1) 查找 | 定向投递 O(K),K = 目标连接数 |
| 连接计数器 | AtomicLong + AtomicInteger |
CAS 无锁 | 极高并发下可能自旋,但连接创建频率远低于消息投递 |
| 消息分发 | Reactor Sinks.Many |
CAS 序列化 + Drain 循环 | 单线程 10 万 msg/s,多线程 3-5 万 msg/s |
多线程并发 tryEmitNext() 时,CAS 竞争可能导致消息丢失(返回 FAIL_NON_SERIALIZED)。框架提供内置重试策略,通过策略模式支持四种重试机制:
- no-retry(默认):单次尝试,性能最优,适合单线程发布场景
- busy-loop:忙等待重试直到成功或超时,可靠性最高
- bounded:固定次数重试,平衡可靠性和资源消耗
- exponential-backoff:指数退避重试,适合高并发竞争场景
配置示例:
im:
sse:
emit-retry:
default-strategy: busy-loop # 默认策略
busy-loop-timeout: 100ms # 超时时间
topic-overrides:
im.chat: bounded # 按 Topic 覆盖策略
可扩展性
| 指标 | 单节点 | 集群(100 节点) |
|---|---|---|
| 最大连接数 | 100 万 | 1 亿(理论无上限) |
| 消息吞吐 | 10 万 msg/s | 1000 万 msg/s |
| 广播延迟 | 1ms | 1-5ms(跨节点) |
| 内存占用 | ~500 MB | ~50 GB |
- JVM 堆内存:
-Xmx2g(100 万连接框架开销 ~240MB,预留 GC 和缓冲区) - 文件描述符:
ulimit -n 1000000 - GC 策略:G1GC 或 ZGC(降低 GC 停顿)
设计支撑:MQ 解耦 + 三层架构
可水平扩展
通过负载均衡
支持分区和复制
解耦发布和订阅
通过消费组负载均衡
可水平扩展
- MQ 吞吐量:RabbitMQ 单节点 ~2-5 万 msg/s,Kafka 单节点 ~10-50 万 msg/s
- 网络带宽:10 万连接 × 10 msg/s × 500 bytes = 50 MB/s
- Kafka 分区数:分区数 = 最大并行消费者数
连接管理可靠性
设计支撑:三层清理 + 位域双守卫
第一层:doFinally
响应式钩子,连接正常关闭时触发。覆盖率 ~95%,零额外开销。
第二层:ConnectionReaper
每 30s 扫描空闲连接,异步强制关闭。10 万连接扫描耗时 ~0.5ms。
第三层:maxConnectionLifetime
连接存活达到 24h 后自动终止。防止长期连接无限占用资源。
位域双守卫:三层清理可能并发触发,使用单个 AtomicInteger 的两个位保证清理仅执行一次:
| 位 | 名称 | 作用 |
|---|---|---|
bit 0 | FLAG_CLEANUP | 竞争清理执行权。三层机制谁先设置谁执行清理,其他线程跳过 |
bit 1 | FLAG_COUNTERS_DECREMENT | 独立保护计数器递减。避免 Reaper 与 doFinally 双重递减导致下溢 |
三层清理覆盖 95%+ 场景,位域双守卫保证并发安全。连接泄露风险降至最低,即使异常断开也能在 24h 内自动回收。
竞品对比
| 框架 | 协议 | 连接管理 | 消息延迟 | 内存效率 | 可扩展性 |
|---|---|---|---|---|---|
| SseIm | HTTP + SSE | 三层清理 + 位域守卫 | 1-5ms | ~240 bytes/连接 | MQ 解耦,水平扩展 |
| WebSocket + STOMP | WebSocket | 应用层心跳 | 1-3ms | ~500 bytes/连接 | 需要额外配置 |
| Socket.IO | WebSocket + 降级 | 心跳 + 重连 | 2-5ms | ~600 bytes/连接 | 需要 Redis Adapter |
| gRPC Streaming | HTTP/2 | 连接池 | <1ms | ~200 bytes/连接 | 负载均衡器支持 |
- 内存效率提升 2-3 倍:~240 bytes/连接 vs ~500-600 bytes(位域优化 + 无额外协议开销)
- 连接管理更可靠:三层清理 vs 应用层心跳,连接泄露风险显著降低
- 架构简单:配置驱动,零代码开箱,SPI 扩展点支持渐进式自定义
- 天然支持水平扩展:MQ 解耦允许 Publisher 和 Subscriber 独立扩缩容
配置参考
所有配置项通过 application.yml 配置。
端点配置(核心)
im:
endpoints:
- path: /api/chat/push
topic: im.chat
mode: push # PUSH: POST → MQ
- path: /api/chat/pull
topic: im.chat
mode: pull # PULL: GET → SSE
SSE 连接配置
| 配置路径 | 类型 | 默认值 | 说明 |
|---|---|---|---|
im.sse.max-total-connections | int | 100000 | 全局最大连接数 |
im.sse.max-connections-per-topic | int | 10000 | 单 Topic 最大连接数 |
im.sse.max-connection-lifetime | Duration | 24h | 连接最大存活时间(设 0 禁用) |
im.sse.heartbeat.enabled | boolean | true | 是否启用心跳 |
im.sse.heartbeat.interval | Duration | 30s | 心跳间隔 |
背压策略
| 配置路径 | 类型 | 默认值 | 说明 |
|---|---|---|---|
im.sse.backpressure.default-strategy | String | drop-oldest | 默认策略:buffer / drop-oldest / error |
im.sse.backpressure.default-buffer-size | int | 1024 | 默认缓冲区大小 |
im.sse.backpressure.topic-overrides | Map | {} | 按 Topic 覆盖策略 |
Reaper(僵尸连接清理)
| 配置路径 | 类型 | 默认值 | 说明 |
|---|---|---|---|
im.sse.reaper.enabled | boolean | true | 是否启用 |
im.sse.reaper.interval | Duration | 30s | 扫描间隔 |
im.sse.reaper.idle-timeout | Duration | 3600s | 空闲超时 |
SPI 扩展点
零侵入扩展,按需渐进式自定义。
1. Factory 模式(组件级自定义)
通过工厂模式替换默认实现:
// 1. 继承骨架类
public class ChatPublisher extends ImPublisher {
public ChatPublisher(String topic, Supplier<Repeater> repeaterSupplier) {
super(topic, repeaterSupplier);
}
@Override
public Serializable push(MultiValueMap<String, String> urlParams, Msg msg) {
log.info("[ChatPublisher] 推送聊天消息");
return super.push(urlParams, msg);
}
}
// 2. 实现工厂
@Component
public class ChatPublisherFactory implements PublisherFactory {
@Override
public Publisher create(String topic, Supplier<Repeater> repeaterSupplier) {
if ("im.chat".equals(topic)) {
return new ChatPublisher(topic, repeaterSupplier);
}
return null; // 其他 Topic 使用默认实现
}
}
2. RepeaterInterceptor(消息拦截器)
拦截消息流的关键节点,用于审计、持久化、历史回放:
@Component
public class RedisReplayInterceptor implements RepeaterInterceptor {
@Autowired
private RedisTemplate<String, byte> redisTemplate;
@Override
public Mono<Boolean> beforeSend(TopicMessage<Msg> message) {
// 发送前持久化到 Redis
redisTemplate.opsForList().rightPush(key(message.getTopic()), serialize(message));
return Mono.just(true); // false 拦截消息
}
@Override
public Flux<Msg> transformSubscribeStream(String topic, Flux<Msg> stream, ...) {
// 订阅时先回放历史消息
Flux<Msg> history = loadFromRedis(topic);
return Flux.concat(history, stream);
}
}
3. AccessFilter(端点访问控制)
Pipeline-Filter 模式的安全扩展:
@Component
@Order(10)
public class JwtAuthFilter implements AccessFilter {
@Autowired
private JwtService jwtService;
@Override
public void doFilter(AccessContext context, AccessFilterChain chain) {
String token = context.getHeader("Authorization")
.filter(h -> h.startsWith("Bearer "))
.map(h -> h.substring(7))
.orElseThrow(() -> new AccessDeniedException("Missing token", HttpStatus.UNAUTHORIZED));
UserDetails user = jwtService.validate(token);
context.setPrincipal(user);
chain.doFilter(context);
}
}
高阶用法
运行时动态注册、自定义消息类型、完全替换消息链路。
运行时动态注册端点
@RestController
public class AdminController {
@Autowired
private EndpointRegistry endpointRegistry;
@PostMapping("/admin/endpoint")
public R<?> addEndpoint(@RequestBody EndpointProperties config) {
endpointRegistry.register(config); // 即时生效
return R.ok();
}
}
自定义消息类型
@Data
public class ChatMsg implements Msg {
private String from;
private String content;
private long timestamp;
// 无需实现 build(),框架提供默认实现
}
完全自定义 Repeater(绕过 MQ)
public class DirectRepeater implements Repeater {
@Override
public Serializable send(TopicMessage<Msg> message) {
// 直接推送到 Redis Pub/Sub,绕过 MQ
redisTemplate.convertAndSend(message.getTopic(), serialize(message));
return message.getId();
}
@Override
public void receive(TopicMessage<Msg> message) {
// 从 Redis 接收后推送到 SSE
sseConnectionManager.publish(message.getTopic(), message);
}
@Override
public Flux<Msg> subscribe(MultiValueMap<String, String> urlParams) {
// 从 Redis 订阅实时流
return redisReactiveTemplate.listenToChannel(topic()).map(this::deserialize);
}
}
定向投递
基于地址模型的精准投递 — 同一套寻址语义,双通道互斥架构,O(K) 精准投递。
快速体验(零 Java 代码)
以聊天室为例。假设已配置 /api/chat/push(PUSH)和 /api/chat/pull(PULL),全局 routing-keys: ["room"]:
alice 和 bob 加入 lobby 房间
// 浏览器原生 EventSource — URL 参数即地址
const es = new EventSource('/api/chat/pull?room=lobby');
es.onmessage = e => console.log('收到:', JSON.parse(e.data));
两人的连接地址均为 "room=lobby"。不同房间(?room=main)的连接互不干扰。
发送消息到 lobby 房间 — 仅 alice 和 bob 收到
# 相同参数格式 → 相同 routingKey → 精准命中
POST /api/chat/push?room=lobby
Content-Type: application/json
{"from": "charlie", "content": "大家好"}
alice 和 bob 同时收到。/api/chat/push 不带参数时则为广播,所有连接收到。
订阅方和发布方使用相同的 URL 参数格式,框架自动完成地址建立、索引构建、精准投递。无需手写路由逻辑,无需额外依赖。下面逐步展开其工作原理。
核心概念
定向投递建立在对称寻址模型上 — 订阅方声明地址,发布方指定目标,框架负责匹配:
| 概念 | 载体 | 语义 | 类比 |
|---|---|---|---|
| 地址(Address) | ConnectionMetadata.routingKey | 订阅方的通信地址 —「消息送到哪里」 | 公寓房号 |
| 目标(Target) | TopicMessage.routingKeys | 发布方的投递目标集合 —「发给哪些地址」 | 快递单上的收件地址列表 |
匹配规则:目标的某个元素等于连接的地址时,消息投递到该连接。routingKeys 为空 → 广播(所有连接收到);非空 → 定向(仅匹配的连接收到)。
同一套 RoutingKeyComposer 同时用于订阅方的地址建立和发布方的目标提取 — 不是两套独立逻辑的拼接,而是一个模型的两侧投影。这从结构上保证了双方寻址语义的一致性。
寻址 API
两种路径,殊途同归 — 最终都生成 TopicMessage.routingKeys:
路径 A — HTTP 端点(零 Java 代码)
使用与订阅方相同的 URL 参数格式,框架自动提取 routingKey:
// 单键定向 — 投递到 lobby 房间
// ?room=lobby → targets = {"room=lobby"}
fetch('/api/chat/push?room=lobby', {
method: 'POST',
headers: {'Content-Type': 'application/json'},
body: JSON.stringify({from: 'bob', content: 'hi lobby'})
});
// 单键多值 — 批量投递到多个房间
// ?room=lobby&room=main → targets = {"room=lobby", "room=main"}
fetch('/api/chat/push?room=lobby&room=main', {
method: 'POST',
headers: {'Content-Type': 'application/json'},
body: JSON.stringify({from: 'charlie', content: '群聊'})
});
// 多键位置配对 — 多租户场景(需端点配置 routing-keys: ["tenant", "room"])
fetch('/api/room/push?tenant=acme&room=lobby', {
method: 'POST',
headers: {'Content-Type': 'application/json'},
body: JSON.stringify({from: 'system', content: '房间通知'})
});
// 无参数 = 广播
fetch('/api/chat/push', { method: 'POST', ... });
路径 B — Java API(服务端内部调用)
注入 RoutingKeyComposer 构建 routingKey,保证与 HTTP 路径格式一致:
@Autowired
private TopicRegistry topicRegistry;
@Autowired
private RoutingKeyComposer composer;
// 定向投递 — Composer 构建 routingKey,与 HTTP 路径格式一致
public void sendToRoom(String room, ChatMsg msg) {
MultiValueMap<String, String> params = new LinkedMultiValueMap<>();
params.add("room", room);
Set<String> targets = composer.extractRoutingTargets(null, params);
TopicMessage<ChatMsg> envelope = TopicMessage.<ChatMsg>builder()
.topic("im.chat")
.msg(msg)
.routingKeys(targets)
.build();
topicRegistry.ensureRepeater("im.chat").send(envelope);
}
// 批量投递
public void sendToRooms(Set<String> rooms, ChatMsg msg) {
MultiValueMap<String, String> params = new LinkedMultiValueMap<>();
rooms.forEach(room -> params.add("room", room));
Set<String> targets = composer.extractRoutingTargets(null, params);
TopicMessage<ChatMsg> envelope = TopicMessage.<ChatMsg>builder()
.topic("im.chat")
.msg(msg)
.routingKeys(targets)
.build();
topicRegistry.ensureRepeater("im.chat").send(envelope);
}
HTTP 路径和 Java API 共用同一个 RoutingKeyComposer — 相同的输入必然产生相同的 routingKey。两条路径只是入口不同,底层寻址逻辑完全一致。
路由键机制
理解了用法,现在看框架如何建立和管理路由键。
地址建立(routingKey 提取)
每个 SSE 连接都携带一个地址标识(routingKey)。框架按四层优先级自动提取:
| 优先级 | 来源 | 适用场景 | 示例 |
|---|---|---|---|
| ① | AccessFilter 设置的 principal(String) |
生产环境 — JWT / OAuth | JWT Filter → setPrincipal("user-123") |
| ② | PrincipalRoutingKeyConverter 从复杂 principal 提取 |
principal 为对象(如 UserDetails) |
UserDetails.getUsername() |
| ③ | X-Me 请求头 |
服务端调用、测试工具 | curl -H 'X-Me: user-123' ... |
| ④ | URL 参数按 routing-keys 配置组合 |
浏览器 EventSource(无法自定义请求头) | EventSource('/pull?room=lobby') |
③④ 层(URL 参数、X-Me 头)均为纯声明,无认证能力,仅用于开发调试。生产环境应通过 ① 层 AccessFilter 实现 JWT 认证,确保地址不可伪造。详见「生产环境建议」。
地址格式
routingKey 格式为规范化查询字符串:key 按字母排序,key=value 以 & 拼接。保证参数顺序无关性:
?room=lobby → "room=lobby"(单键) ?tenant=acme&room=lobby → "room=lobby&tenant=acme"(key 字母排序) ?room=lobby&tenant=acme → "room=lobby&tenant=acme"(相同结果,顺序无关) 无参数 → null(广播连接,不参与定向索引)
routingKey 生命周期
routingKey 贯穿投递全链路:
- 订阅时:Composer 从 URL 参数组合 routingKey → 存入
ConnectionMetadata→ 注册到反向索引(routingKeyIndex) - 发布时:同一 Composer 从 URL 参数提取目标集合 → 通过反向索引 O(K) 定位目标连接
- MQ 桥接:
routingKeys随TopicMessage作为不可变载荷完整保留,支持跨服务定向路由
同一 Topic 的 PUSH/PULL 端点必须使用相同的 routing-keys 配置,确保双方生成的 routingKey 格式完全一致。可全局配置 im.sse.routing-keys,也可通过 EndpointProperties.routingKeys 按端点覆盖。
双通道架构
消息经 MQ 到达 Repeater.receive() 后,按 routingKeys 是否为空分流到两条结构性互斥的通道:
routingKeys 为空?
topicSink
O(1) emit
routingKeyIndex
O(K) 精准投递
- 广播通道(
topicSink):共享的 per-topic Sink,O(1) emit - 定向通道(
directedSink):per-connection Sink,通过反向索引routingKeyIndex查找目标连接,O(K) 精准投递 - 两条通道结构性互斥:连接在创建时就决定走哪条通道,消息到达时无需 filter,不可能重复投递
广播 O(1)、定向 O(K)(K 为目标连接数,K << N)。纯广播连接零额外内存开销,仅有 routingKey 的连接额外 ~128 字节(directedSink + 索引条目)。
当定向消息的 routingKeys 无任何匹配的活跃连接时,消息静默丢弃(不排队等待)。Micrometer 指标 directed-publish.no-targets 会记录此情况,可用于监控告警。
完整数据流
EndpointRouter
Composer 提取 routingKeys
msg={content:"hi"}
routingKeys={"room=lobby"}
routingKeys 随消息
完整保留
routingKeyIndex 查找
O(K) 精准投递
routingKey="room=lobby"
directedSink
routingKey="room=main"
未命中
配置
两级配置:全局默认 + 端点级覆盖。
im:
sse:
routing-keys: [room] # 全局默认(按房间路由)
endpoints:
- path: /api/chat/pull
topic: im.chat
mode: pull # 继承全局 ["room"]
- path: /api/chat/push
topic: im.chat
mode: push # 继承全局 ["room"]
- path: /api/room/pull
topic: im.room
mode: pull
routing-keys: [tenant, room] # 端点覆盖:多租户聊天室
- path: /api/room/push
topic: im.room
mode: push
routing-keys: [tenant, room] # 必须与对应 pull 端点一致
routing-keys 语义:
null/ 未配置 /[]— 全部 URL 参数参与(默认)["*"]— 显式全部参数(等价于 null)["room"]— 指定字段["tenant", "room"]— 指定多字段
扩展点
自定义路由键组合器(RoutingKeyComposer)
默认实现 DefaultRoutingKeyComposer 按规范化查询字符串格式组合。如需自定义(如 HMAC 签名、自定义分隔符),实现 RoutingKeyComposer 接口并注册为 Bean,框架自动替换默认实现:
@Component
public class SignedRoutingKeyComposer implements RoutingKeyComposer {
@Override
public String composeRoutingKey(List<String> routingKeys,
MultiValueMap<String, String> params) {
String raw = buildRawKey(routingKeys, params);
return raw + "&sig=" + hmacSign(raw);
}
@Override
public Set<String> extractRoutingTargets(List<String> routingKeys,
MultiValueMap<String, String> params) {
// 需与 composeRoutingKey 保持寻址语义一致
}
}
自定义 Principal 转换器(PrincipalRoutingKeyConverter)
当 principal 为非 String 类型时,框架通过 PrincipalRoutingKeyConverter 将其转换为 routingKey:
@Component
public class UserDetailsRoutingKeyConverter
implements PrincipalRoutingKeyConverter {
@Override
public boolean supports(Class<?> principalType) {
return UserDetails.class.isAssignableFrom(principalType);
}
@Override
public String convert(Object principal, AccessContext context) {
return ((UserDetails) principal).getUsername();
}
}
多个转换器按 getOrder() 排序依次尝试,首个 supports() 返回 true 且转换结果非空的获胜。
生产环境建议
URL 参数和 X-Me 请求头均为纯声明,任何人可伪造。生产环境必须通过 AccessFilter 实现认证,确保地址不可篡改。
最小 JWT 认证示例:
@Component
@Order(10)
public class JwtAuthFilter implements AccessFilter {
@Autowired
private JwtService jwtService;
@Override
public void doFilter(AccessContext context, AccessFilterChain chain) {
String token = context.getHeader("Authorization")
.filter(h -> h.startsWith("Bearer "))
.map(h -> h.substring(7))
.orElseThrow(() -> new AccessDeniedException(
"Missing token", HttpStatus.UNAUTHORIZED));
String userId = jwtService.parseUserId(token);
context.setPrincipal(userId);
chain.doFilter(context);
}
}
客户端改为通过请求头传 token:
import { fetchEventSource } from '@microsoft/fetch-event-source';
await fetchEventSource('/api/chat/pull', {
headers: { 'Authorization': 'Bearer eyJhbG...' },
onmessage(ev) { console.log(JSON.parse(ev.data)); }
});
认证方式决定路由策略:JWT principal 为 userId 时按用户身份路由("user-123");若业务需要按房间/租户路由("room=lobby&tenant=acme"),可在 JWT Filter 中从 token claims 提取 tenant/room 并通过 Converter 组合,或继续使用 URL 参数方式(④ 层)。
监控运维
运行时查询 API、Micrometer 指标、日志级别。
运行时查询 API
框架提供 ImAdminController,暴露以下接口:
| 接口 | 说明 |
|---|---|
GET /im/admin/endpoints | 查询所有已注册端点 |
GET /im/admin/topics | 查询所有活跃 Topic 及连接数 |
GET /im/admin/connections/{topic} | 查询指定 Topic 的连接详情 |
GET /im/admin/stats | 查询全局统计信息 |
生产环境应通过 AccessFilter 为 /im/admin/** 添加访问控制(认证、授权、IP 白名单)。
Micrometer 指标
自动暴露以下指标(前缀 im.sse),无缝接入 Prometheus / Grafana:
| 指标名 | 类型 | 说明 |
|---|---|---|
connections.total | Gauge | 当前总连接数 |
connections.active | Gauge | 按 Topic 的活跃连接数 |
subscribe.rate | Counter | 订阅速率 |
publish.rate | Counter | 发布速率 |
reaper.zombies | Counter | 清理的僵尸连接数 |
故障排查
常见问题诊断与解决方案。
常见问题
| 现象 | 可能原因 | 解决方案 |
|---|---|---|
| SSE 连接无法建立 | 端点未配置 | 检查 im.endpoints 配置,确认 mode: pull 的端点存在 |
| 消息未收到 | RabbitMQ binding 未创建 | 查看启动日志中 [TopicFactory] 相关日志,确认 binding 已创建 |
| 连接频繁断开 | 反向代理空闲超时 | 确认 Nginx proxy_read_timeout 大于心跳间隔(默认 30s) |
| 连接数持续增长 | 客户端未正确关闭连接 | 确认页面卸载时调用 eventSource.close()。Reaper 会清理空闲连接 |
日志级别
logging:
level:
cn.gmlee.tools.im: DEBUG # 框架核心日志
cn.gmlee.tools.im.sse: DEBUG # SSE 连接管理日志