Clojure 实时通信与 Sente

深入 Clojure 实时通信实践:WebSocket 与 SSE 的选型依据、Sente 的握手协议与双向通道模型、服务端 Ring 路由与 ClojureScript 客户端接入、房间与广播实现、断线重连与状态同步策略,以及多实例伸缩、心跳保活与生产排错要点。

裸用 WebSocket 写实时功能,第一周很爽,第三个月就会开始还债:连接注册表没人管、心跳缺失导致僵尸连接堆积、断线后客户端状态与服务端不一致、多实例部署时广播只打到一半用户。Clojure 生态里的 Sente 正是为了把这些「实时通信的脏活」收敛成一层协议而存在。本文从传输选型讲到多实例伸缩,覆盖 Sente 的握手与通道模型、房间广播、断线重连与状态同步,帮你把实时能力做成可运维的基础设施,而不是散落在业务代码里的长连接。

1. 传输选型:WebSocket、SSE 与长轮询

1.1 三种传输的本质差异

实时通信不是「有 WebSocket 就上 WebSocket」。三条路线解决的是不同问题:

传输方向协议基础浏览器支持典型场景
WebSocket双向独立协议(HTTP 升级)全支持聊天、协同编辑、行情
SSE服务端 → 客户端单向纯 HTTP全支持(含自动重连)通知、进度、订阅流
长轮询伪双向纯 HTTP全支持兜底、老旧代理环境

关键判断顺序是:先问「客户端要不要主动推消息」。如果只是服务端推送,SSE 的心智负担和运维成本都远低于 WebSocket——它走标准 HTTP,天然兼容现有网关、鉴权和限流,浏览器还会自动重连。

1.2 选型决策表

维度WebSocketSSE长轮询
服务端推送支持支持支持
客户端上行支持不支持(需另开请求)不支持
自动重连需自实现浏览器内建天然
消息序号/续传需自实现Last-Event-ID 内建需自实现
代理/负载均衡友好度中(需支持 Upgrade)高高
连接数成本每连接一个长连接同同

在 Clojure 侧,底层传输的实现细节(http-kit 的 WebSocket、SSE 响应头、连接注册表)可以参阅 Clojure 网络服务深入 。本文关注的是再往上一层:如何把「连上之后怎么组织消息」这件事标准化。

1.3 为什么还需要 Sente 这一层

Sente 的价值在于它定义了一套传输无关的双向 RPC 协议:

  • 客户端与服务端用统一的 event 结构通信,形如 [event-id payload ?reply-fn];
  • 底层可以是 WebSocket,也可以是 Ajax 长轮询,业务代码不感知;
  • 内建心跳、自动重连、消息缓冲与 :chsk/uid 身份标识;
  • 服务端可以按 uid 或 channel 定向推送。

换句话说,Sente 把「实时通道」抽象成了类似 RPC 的接口,而把「连接怎么活下来」变成框架职责。

它不解决什么:Sente 不是消息队列,没有持久化、没有消费位点、没有顺序保证。它只保证「尽力把消息送到当前活着的连接」。需要可靠投递、需要重放历史消息,那是 Kafka 或事件溯源该干的事。把 Sente 当成「带推送的 RPC」而不是「消息总线」,架构就不会走偏。

2. Sente 架构与握手协议

2.1 双向通道模型

Sente 在两端各维护一个 channel socket(chsk):

  • 服务端:每个连接对应一个 uid,make-channel-socket-server 返回的 :ring-ajax-post、:ring-ajax-get-or-ws-handshake 挂到 Ring 路由上;
  • 客户端:make-channel-socket-client! 返回 chsk-send!、chsk/state、事件回调。

消息只有三种形态:普通事件(fire-and-forget)、请求-应答(带 reply-fn)、服务端推送(服务端主动 chsk-send!)。

2.2 握手流程

握手是 Sente 最容易出错的地方,本质是一次「身份协商」:

;; 服务端:握手时决定是否接受连接,并返回 uid
(defn ring-handler [req]
  (let [{:keys [ring-ajax-get-or-ws-handshake]} @sente-chsk]
    (ring-ajax-get-or-ws-handshake req)))

;; 关键:握手时服务端可以读取会话、决定 uid
(defn- :ring-ajax-get-or-ws-handshake [req]
  (let [uid (or (get-in req [:session :uid])
                (str (java.util.UUID/randomUUID)))]
    ;; 返回 200 并附带 uid,客户端后续所有消息都带这个 uid
    ...))

握手失败时 Sente 会返回 403 或 401,客户端进入 :failed 状态而不会无脑重试——这正是「启动失败快」的思路:认证不过就别反复握手。

2.3 传输回退

Sente 客户端默认策略是「能用 WebSocket 就用,否则退回 Ajax」。回退由 :type 参数控制:

(require '[taoensso.sente :as sente])

(def chsk-client
  (sente/make-channel-socket-client!
   "/chsk"                                  ; 握手与上行端点
   {:type :auto                             ; :auto | :ws | :ajax
    :packer :edn}))                         ; 序列化方式

:auto 会在 WebSocket 握手失败时自动降级。生产上要留意:降级是静默的,如果网关把 Upgrade 头吃掉了,你会一路退回长轮询而毫无察觉。因此要在客户端把 :type 上报到监控。

2.4 序列化与协议契约

Sente 的 :packer 决定消息如何序列化,常用 :edn 与 :json:

  • :edn:Clojure 原生,支持关键字、集合、ratio 等;但要求两端都是 Clojure,且 EDN 解析器对不可信输入有安全风险,必须配白名单读取器。
  • :json:跨语言友好,但关键字会变成字符串,需要显式转换。

协议契约上建议显式定义事件 schema:每个 event-id 的 payload 形状用 Malli 或 spec 描述,在服务端入口和客户端出口各校验一次。这样当两端版本不同步时,能在边界立刻报错,而不是在业务深处抛 NullPointerException。事件命名建议用带命名空间的关键字(:chat/message、:room/join),避免全局关键字污染。

3. 服务端接入:Ring 路由与身份

3.1 依赖与基本接线

(ns app.realtime
  (:require [taoensso.sente :as sente]
            [taoensso.sente.server-adapters.http-kit :refer [get-sch-adapter]]))

(defonce chsk-server
  (sente/make-channel-socket-server!
   (get-sch-adapter)
   {:user-id-fn (fn [ring-req] (get-in ring-req [:session :uid]))}))

(defn chsk-routes []
  (let [{:keys [ring-ajax-post ring-ajax-get-or-ws-handshake]} @chsk-server]
    ["/chsk" {:post ring-ajax-post
              :get  ring-ajax-get-or-ws-handshake}]))

用 reitit 挂载时,把这两个端点当作普通路由注册即可;中间件(会话、CSRF 白名单、限流)照常生效。

3.2 身份:user-id-fn 与握手状态

user-id-fn 是安全边界:它从 Ring 请求中提取 uid。不要信任客户端传来的 uid,必须来自服务端会话或令牌解析。握手中的 :client-id 是浏览器本地随机值,只用于区分同一用户的多个标签页,绝不能当作身份。

3.3 与中间件整合

两个坑:

  1. CSRF:Sente 的 POST 上行是跨站的表单提交,若用了 ring-anti-forgery,需要把 /chsk 的 POST 加入白名单,或改用 WebSocket 主通道。
  2. 会话读取:ring-ajax-get-or-ws-handshake 必须跑在 wrap-session 之内,否则 user-id-fn 拿不到会话,所有人都会变成匿名。

4. 客户端接入:ClojureScript 与 re-frame

4.1 建立连接

(ns app.ws
  (:require [taoensso.sente :as sente]
            [re-frame.core :as rf]))

(def chsk
  (sente/make-channel-socket-client!
   "/chsk" {:type :auto :packer :edn}))

(defn start! []
  (let [{:keys [chsk ch-recv send-fn]} chsk]
    ;; 把服务端推送的事件汇入 re-frame 事件循环
    (rf/dispatch [:ws/start ch-recv send-fn])
    chsk))

4.2 事件分发

Sente 客户端的 ch-recv 是一个 core.async 通道,里面流出 [event-id payload]。典型做法是起一个 go-loop 把消息转成 re-frame 事件:

(rf/reg-event-fx
 :ws/start
 (fn [{:keys [db]} [_ ch-recv send-fn]]
   {:db (assoc db :ws/send send-fn)
    :fx [[:ws/listen ch-recv]]}))

(defn- listen-loop [ch-recv]
  (go-loop []
    (when-let [[event-id payload] (<! ch-recv)]
      (case event-id
        :chsk/state  (rf/dispatch [:ws/state payload])
        :chsk/recv   (rf/dispatch [:ws/push payload])
        :chsk/handshake (rf/dispatch [:ws/ready payload])
        nil)
      (recur))))

把连接状态(:chsk/state)也纳入 re-frame,好处是 UI 可以直接订阅「是否已连接」,而不用到处 (.-open? chsk)。

4.3 上行请求与应答

Sente 支持「请求-应答」语义:chsk-send! 的第三个参数是回调,服务端用 ?reply-fn 回包。

;; 客户端:发送并等待应答(带超时)
(chsk-send! [:room/join {:room-id "r-42"}] 3000
  (fn [reply]
    (if (sente/cb-success? reply)
      (rf/dispatch [:room/joined reply])
      (rf/dispatch [:room/join-failed reply]))))

;; 服务端:处理并回包
(defn handle-event [{:keys [event ?reply-fn uid]}]
  (case (first event)
    :room/join (let [room-id (get-in event [1 :room-id])]
                 (join! room-id uid)
                 (when ?reply-fn (?reply-fn {:ok true :members (count (get @rooms room-id))})))))

超时参数很重要:没有它,一次丢包就会让 UI 永远停在「加入中」。所有带应答的上行都应当带超时,并在超时后走降级路径(提示重试或退回轮询)。

4.4 与 re-frame 单向数据流的契合

4.5 与 re-frame 单向数据流的契合

re-frame 的事件循环天生适合承载实时消息:服务端推送 → :ws/push → 更新 app-db → 订阅重渲染。不要在 ch-recv 的 go-loop 里直接改 DOM 或调用业务函数,否则就绕开了单一数据源。完整的前端数据流组织可参考 ClojureScript 全栈开发 。

5. 房间、广播与定向推送

5.1 广播模型

服务端的 chsk-send! 签名是 (chsk-send! uid event),因此「广播」本质上就是遍历 uid 集合逐个发送。Sente 自己维护了一个 :uid->ch 注册表(可从 chsk-server 的 atom 里读到 :connected-uids)。

5.2 房间实现

(defonce rooms (atom {}))   ; room-id -> #{uid}

(defn join! [room-id uid]
  (swap! rooms update room-id (fnil conj #{}) uid))

(defn leave! [room-id uid]
  (swap! rooms update room-id disj uid))

(defn broadcast! [room-id event]
  (let [{:keys [send-fn]} @chsk-server]
    (doseq [uid (get @rooms room-id)]
      (send-fn uid event))))

三个工程要点:

  • 房间成员要随连接断开自动清理:监听 :chsk/uidport-close 事件,把 uid 从所有房间摘掉,否则会留下幽灵成员,广播时白白遍历。
  • 成员集合是热点数据:高并发下 swap! 会成为瓶颈,房间数不多时可以接受;房间极多时应改用分片或直接用 Redis 维护。
  • 广播是 O(n) 的:一个万人房间每次广播就是一万次 send-fn,必须做批量或走发布订阅中间层。

5.3 定向推送与在线状态

connected-uids 提供了 :any、:ws、:ajax 三个集合,可用于:

  • 在线人数统计(:any 的 count);
  • 「只推给还在线的用户」,避免给已离线 uid 发消息;
  • 区分 WebSocket 与降级用户,排查「为什么只有部分人收不到」。

6. 断线重连与状态同步

6.1 自动重连

Sente 客户端内建重连,状态在 :chsk/state 里流转::open → :closed → :open。默认是指数退避,避免服务端重启时被全量客户端同时打崩(惊群)。

;; 可配置的重连策略(示意)
(sente/make-channel-socket-client!
 "/chsk"
 {:type :auto
  :ws-kalive-ms 10000      ; 心跳间隔
  :lp-timeout 60000})      ; 长轮询超时

6.2 消息缓冲:不丢,但可能重复

Sente 在断线期间会把待发消息缓冲在客户端,重连后补发。这带来一个必须接受的语义:至多一次会变成至少一次。因此:

  • 上行消息要带幂等键,服务端按 key 去重;
  • 不要用「发一条消息就期望恰好执行一次」的假设写业务。

6.3 重连后的状态对齐

这是实时系统最容易翻车的地方。断线期间服务端数据可能已经变了,客户端手里的还是旧快照。两种可靠策略:

  1. 版本号 + 增量:每条推送带 version,重连时客户端上报本地 version,服务端补发差异;差异过大则退回全量。
  2. 快照 + 订阅:重连成功后先拉一次快照,再重新订阅增量流。简单可靠,适合状态不大的场景。
(rf/reg-event-fx
 :ws/ready
 (fn [{:keys [db]} _]
   ;; 重连成功:先对齐版本,再恢复订阅
   {:fx [[:http/get-snapshot {:on-success [:ws/snapshot]}]
         [:ws/resubscribe (:room db)]]}))

关键心法:把「连接」和「状态」解耦。连接可以断可以重连,状态必须由版本或快照显式对齐,绝不能假设连接一恢复数据就是对的。

7. 生产化:心跳、伸缩与可观测

7.1 心跳与僵尸连接

TCP 长连接在 NAT、负载均衡空闲超时下会「假活」——对端已死,本端还以为是开的。心跳是唯一解药:Sente 默认 10s 发一次,超过阈值未收到就判定断开并重连。调参原则是心跳间隔小于链路上最短的空闲超时(常见是网关的 60s)。

7.2 多实例下的广播

单机 chsk-send! 只在本进程的 uid 表里找连接。一旦水平扩容,A 实例的用户收不到 B 实例广播的消息。解决方案是引入中间层:

;; 订阅 Redis 频道,把跨实例消息转发给本进程连接
(defn start-cross-instance! []
  (redis/subscribe! "realtime:broadcast"
    (fn [[room event]]
      (broadcast! room event))))

架构上就是「每个实例只推自己持有的连接,跨实例靠 pub/sub 转发」。Kafka 或 NATS 也可以承担这个角色,选型可参考 Clojure Kafka 事件流集成 。

7.3 指标与排错

必看指标:

指标含义异常信号
在线连接数:any 大小骤降 = 网关或握手问题
握手失败率401/403 比例升高 = 认证或会话配置错
平均重连间隔客户端上报过短 = 服务端不稳定
消息积压缓冲队列长度持续增长 = 消费端过慢
传输类型分布ws/ajax 占比ajax 占比高 = Upgrade 被拦

把传输类型分布暴露出来尤其重要,因为降级是静默的,只有指标能告诉你实时性其实已经退化成了长轮询。

7.4 上线检查清单

检查项期望不做的后果
握手鉴权uid 来自会话/令牌任何人都能冒充
心跳间隔小于网关空闲超时僵尸连接堆积
房间清理断开时摘除 uid广播越跑越慢
跨实例转发pub/sub 已接入扩容后消息丢一半
上行超时全部带 timeoutUI 永久卡在加载态
幂等键上行消息带去重键重连后重复执行
传输类型上报ws/ajax 分布可查静默降级无人知
背压慢客户端有丢弃策略内存被单个客户端拖垮

背压这一项常被忽略:如果某个客户端读得极慢,服务端的发送缓冲会无限增长。Sente 层可以设置发送队列上限,超出后主动断开该连接——牺牲一个慢客户端,保住整个进程。

8. 小结

Clojure 的实时通信能力,底层靠 http-kit 这类适配器,上层靠 Sente 把「连接生命周期」与「消息语义」标准化。真正决定成败的是几个工程决策:

  • 传输选型:先问是否需要客户端上行,单向推送优先 SSE;
  • 握手即认证:user-id-fn 是安全边界,别信客户端 uid;
  • 房间要能自动清理:连接断开必须摘除成员,否则广播会越跑越慢;
  • 重连必对齐状态:用版本号或快照显式同步,不要假设连接恢复数据就对;
  • 扩容要跨实例转发:单机 uid 表不跨进程,靠 pub/sub 补齐。

把这五件事做对,实时功能才从「能跑」变成「能运维」。消息流的上下游编排可继续参考 Clojure core.async 深入 。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「clojure」更多文章

  1. Clojure 桌面 UI:cljfx 与 JavaFX 实战
  2. JSON/EDN 序列化与数据格式互操作
  3. JVM 调优与容器化部署:GC、JFR 与 Docker