基于webSocket协议单连接约束下的高并发:企微智能机器人的选举、主备切换与流量分发


前言

企业微信智能机器人(长连接模式)有几条很”硬”的约束,它们直接决定了后端架构长什么样:

  • 连接唯一:一个机器人同一时刻只能保持一条有效长连接,新连接建立时会直接踢掉旧连接;
  • 窗口有限:流式消息从首次发送算起有明确的生命周期上限,超时后消息自动结束,无法继续刷新;
  • 频率受限:对同一个会话的发送频率有硬性上限,被动回复与主动推送共用一个额度。

官方对高可用的建议也很明确:用主备切换,而不是同时多连接。

这几条约束叠在一起,架构的出发点就定死了:接入侧必须选出唯一的节点,而且这个节点还不能”重”——AI 问答任务动辄要跑几分钟,让唯一入口承担计算,等于把整个系统的天花板钉死在一台机器上。

于是整套机制被拆成三件事,本文就按这三块来讲:

要回答的问题对应的机制
谁来持有这条唯一连接?选举
持有者挂掉怎么办?主备切换
上千人并发怎么扛?流量分发

一、整体结构:一个镜像,三种角色

这套架构最核心的设计是:所有实例跑的是同一份代码,角色由”锁”决定,而不是由部署决定。

                    企业微信(唯一一条有效长连接)
                                 │
                                 ▼
        ┌────────────────────────────────────────────────┐
        │              接入节点(Leader)                 │
        │   独占长连接:接收消息 / 校验 / 落库 / 派发任务  │
        │   轻量推送:轮询任务进展,把结果回写给用户        │
        └───────────────────┬────────────────────────────┘
                            │ 派发计算任务(服务发现寻址)
          ┌─────────────────┼─────────────────┐
          ▼                 ▼                 ▼
   ┌────────────┐    ┌────────────┐    ┌────────────┐
   │  同构实例  │    │  同构实例  │    │  同构实例  │
   │ 备用 + 计算│    │ 备用 + 计算│    │ 备用 + 计算│  ← 可水平扩容
   └─────┬──────┘    └─────┬──────┘    └─────┬──────┘
         └─────────────────┼─────────────────┘
                           ▼
        ┌────────────────────────────────────────────────┐
        │    共享存储:选主用的锁  +  任务与消息的状态机    │
        │                ↑ 全系统唯一的"真相"              │
        └────────────────────────────────────────────────┘
角色由什么决定职责
接入节点持有分布式锁独占长连接;接收消息与事件;排重、鉴权、路由、落库;把重活派发出去;负责结果回写
备用节点没有抢到锁周期性参与竞争;同时作为计算节点承接任务
计算节点被接入节点选中执行 AI 任务,把过程与结果持续写回共享存储

这个”同构”设计的收益很直接:不需要为接入角色单独准备一套部署,不需要人工介入切主,扩容计算节点时顺带壮大了备用池。锁在谁手上,谁就是接入层。


二、选举:把分布式共识降维成一个原子操作

选举要解决的问题很朴素:在任意时刻,有且只有一个实例认为自己是接入节点。 实现上只用了一把带过期时间的分布式锁。

2.1 三件事:抢、续、放

  • 抢锁:以”不存在才写入”的语义一次性完成抢锁并附带过期时间,写入的值是本实例的唯一标识。这一步本身是原子的,不需要额外加锁。
  • 续约:持锁期间定期重置过期时间,让锁一直归属于自己。
  • 释放:正常退出时主动删掉锁,让备用节点立刻可以上位。

关键在于续约和释放都必须校验锁的归属:只有当锁里存的还是自己的标识时才允许操作。否则就会出现”读出来还是自己的、动手时已经被别人抢走”的时间窗——最后的结果是A 把 B 的锁删了,两个节点同时自认为可以建连。这是分布式锁最经典的翻车点,也是选举机制里最不能省的一步校验。

2.2 时间参数怎么定

参数取值设计意图
锁的过期时间20s崩溃场景下的最大接管时间
续约间隔7s约为过期时间的 1/3,留出两次容错机会
竞争间隔5s备用节点的抢锁频率,决定优雅下线的接管速度

续约间隔取过期时间的 1/3 是这套机制里最重要的一个比例。它意味着一两次续约失败不会立刻丢锁,还有机会在过期前补上。如果按 1/2 甚至更激进地设置,一次网络抖动就足以让锁过期,两个节点同时认为自己可以建连。

过期时间本身是”可用性”与”安全性”的取舍:设长了,崩溃后无人接管的时间就久;设短了,续约失败的误判概率就高。20s 配上 7s 的节奏,刚好给出两次容错。

而优雅退出走的是另一条路——不等过期,直接让锁,于是接管时间从 20s 量级缩短到一个竞争周期。

2.3 防脑裂:续约失败就”自杀”

选主类系统最怕的不是节点挂掉,而是两个节点都以为自己是主。落到长连接上,现象会非常糟:新连接踢旧连接,旧连接又重连再踢回来,形成”连接乒乓”,用户感知到的就是消息时好时坏、断断续续。

这里的处理原则只有一句话:没有锁,就没有身份。 一旦续约发现自己已经不再是锁的持有者,立即放弃接入节点身份,并主动断开现有连接,回到竞争队列重新排队——而不是心存侥幸地”再撑一会儿”。

这个”主动断”的动作还有一层工程意义:断开连接会顺带打断等待连接的阻塞流程,让实例能干脆利落地退出接入态,而不是傻等连接自然断开。

2.4 隔离:让多机器人、多环境互不干扰

锁按机器人维度隔离,因此同一套 Redis 上跑多个机器人、多个环境的实例时彼此不干扰。

此外还有一道环境闸门:当多套测试环境共用同一个机器人凭证时,只放行其中一套接入。原因很现实——这类”配置层面的脑裂”必须在代码层面堵死,因为它的现象(连接反复断开)和网络故障几乎一模一样,排查成本极高。


三、主备切换:三种故障,接管时间都是可预期的

接入节点的主流程就是一个”竞争循环”:抢到锁就进入接入态,连接结束或被抢占就退出来重新竞争。对应到真实故障,收尾路径有三条:

故障类型旧接入节点动作新节点接管耗时
优雅下线(发布、缩容)主动让锁秒级(一个竞争周期)
进程崩溃 / 被强制终止无(已经死了)≤ 过期时间(20s 量级)
连接被踢 / 网络抖动停止续约 → 让锁 → 回竞争循环立即重抢,通常仍由本节点夺回

这三条路径合起来,把”故障接管”从一件需要人工判断和干预的事故,变成了一个时间上限已知的自动过程:最坏情况也在 20s 内自愈,发布扩容期间的接管则在秒级完成。

切换过程中有两个容易被忽略、但必须处理的细节:

  • 每次上位都要重建连接、重注册事件,并且先清理上一轮的监听。否则每次切主都会叠加一层回调,一条消息被处理多次——这种问题在切主不频繁的环境里可以潜伏很久。
  • 等待连接结束要靠在事件上挂状态,而不是依赖”连接调用会阻塞”的直觉。事件驱动的客户端往往建立连接就返回了,真正需要等待的是”断开”这个事件。

四、流量分发:接入只做接入,计算可以水平扩

入口只有一个,但计算不能只有一台机器。分发被拆成两条互不干扰的链路。

4.1 第一跳:从接入节点到计算节点

接入节点收到消息后并不自己跑 AI,而是做一遍轻量处理再派发:消息排重、身份与权限校验、指令路由、会话级并发拦截(同一个会话同时只允许一个任务在跑,避免上下文互相污染)、先给用户一个”已收到”的回执,最后把任务和原始消息上下文一并落库,再把计算请求投出去。

原始消息上下文也要落库,这一点很关键:后续推送阶段需要用它来重建”往哪个会话回复、回复到哪条消息”的上下文。也就是说,接入节点承接的是状态,而不是连接。

派发的重点在寻址与降级,一共三层:

层级策略失败后果
1服务发现,任取一个健康实例回退到第 2 层
2配置的固定地址请求失败则回退到第 3 层
3接入节点本地执行业务不中断,只是占用接入节点资源

三层降级保证了”服务发现挂了”和”计算节点全挂了”都不会让用户收不到回复。

但有一条红线:当计算节点返回的是”这个会话已有任务在跑”这类业务拒绝时,绝不允许降级到本地重跑。因为本地重跑等于绕开会话互斥、主动制造并发,会直接污染会话上下文。这属于”降级降错方向”,比不降级更危险。

4.2 第二跳:从计算节点到用户

计算节点做两件事:把模型输出流写进缓存,同时把累积内容实时写进共享存储。它完全不碰长连接。

回写给用户的工作,统一由接入节点上的定时任务负责,按固定周期扫描共享存储里的任务状态:

  • 并发推送:一轮扫描到的记录并发发送,而不是顺序发送。上百条记录顺序推送的耗时是线性累加,并发之后约等于最慢一条的耗时,量级上的差别。
  • 增量去重:记录每条消息上次推送的内容,内容没变就不推——这是对会话频率限制最直接的尊重。
  • 心跳文案:长时间没有新内容时,按固定间隔推一次”已进行 X 分 XX 秒,仍在处理中”的提示,让用户知道后台还在跑,同时把频率控制在限速的安全边际内。
  • 超时截断 + 事后补推:接近流式消息窗口上限时,先把当前内容收尾,等后台任务真正完成后再主动推送完整结果(带重试)。用户不会看到”消息卡死在半截”。
  • 状态由推送侧标记:所有”已回复”的状态扭转都发生在推送成功之后,保证”系统认为已回复”和”用户确实收到”这两件事一致,否则就会出现用户什么都没收到、系统却以为处理完毕的黑洞。

4.3 为什么”状态全落库”是这套分发能成立的前提

把两条链路连起来看,会得到一个不显眼但很关键的性质:接入节点不持有任何不可恢复的内存状态。 收到的消息、AI 的输出、推送的进度、会话的互斥锁,全部沉淀在共享存储里。

于是当接入节点挂掉、新节点上位时,在途任务不会丢:新节点接手长连接后,第一轮扫描就能读到所有未推送完成的任务,继续把结果推给用户。

这就是为什么”接入节点”可以随时被替换——它只是一副可以随时换掉的”嘴和耳朵”,记忆并不在它身上。


五、这套机制给出了什么能力

能力靠什么实现达到的效果
稳定通信单连接 + 选举 + 续约校验从机制上消除多主抢连导致的”连接乒乓”
故障自愈锁的过期 + 主动让锁崩溃 20s 量级、优雅下线秒级自动接管
发布不中断下线时主动让锁,新实例随即参与竞争滚动发布期间消息不丢
高并发承载从”每个用户一个等待协程”改为”定时任务统一推送”支撑上千人同时问答
水平扩容接入与计算解耦,计算节点无状态加实例即加算力
长任务不丢结果窗口到期前截断 + 完成后主动补推突破流式消息的生命周期限制
失败可观测双维度状态机 + 失败原因留痕失败可定位、可人工修正恢复
多机器人 / 多环境隔离锁按机器人维度隔离 + 环境闸门环境之间零互相干扰

小结

如果把企业微信的长连接约束翻译成一句话,就是:接入路径必须是唯一的,但计算路径不能是唯一的。

这套机制的三个部分,各自回答一个问题:

  • 选举回答”谁来持有这条唯一连接”—— 用一把带过期时间的锁,把分布式共识问题降维成一个原子操作;
  • 主备切换回答”它挂了怎么办”—— 用过期、心跳、主动让锁三条路径,把接管时间变成可预期的秒级 / 20s 量级;
  • 流量分发回答”上千人并发怎么扛”—— 把唯一入口做成轻量派发器,把状态全部沉淀到共享存储,让计算侧可以随便扩、随便重启。

其中我觉得最值得带走的一点是:接入节点不应该持有状态。 接入、计算、推送三者分离之后,“主备切换”就不再是一件需要小心翼翼的事故处理,而只是一次普通的进程重启。