深浅色
EinoTalk
基于 Go + React + TypeScript 的全栈即时通讯应用,在完整 IM 闭环之上集成 CloudWeGo Eino AI 助手。
一句话
面向 IM 场景的全栈应用:Go 后端提供注册登录、私聊/群聊、会话、联系人、文件上传与 WebSocket 实时消息;React + TS 前端提供桌面级聊天界面;内置「AI 助手」作为系统用户,支持流式回复、ReAct 工具调用和滚动摘要记忆。我负责 AI 助手的流式链路落地、多实例集群一致性修复,以及工程化收尾(测试、CI、优雅关停、可观测)。
项目边界(先自己讲清楚)
基线是一套已经跑通单机 IM 闭环的代码库。我接手时它已经能演示,但存在「纸面分布式」问题:多实例会静默丢消息、关停会 panic、核心分发链路零测试、README 宣称的能力和代码不一致。
我这 23 次提交做的是把它从「能演示」变成「讲得住」:Agent 流式链路、集群一致性、测试与 CI、关停语义。
面试时主动说明这段边界,比自己被追问出来强。
架构
三层架构,开发顺序 Model -> DAO -> Service -> Controller -> 路由注册:
Request -> Controller -> Service -> DAO -> Database
前端 React 18 + TS + Zustand (/chat)
| REST /login /register /group
| WS /user/wsLogin?token
v
Gin + JWT ---+--- MySQL(6 张表)
+--- Redis(消息 List / 会话 / 验证码)
+--- ChatServer(channel 模式 / Kafka 模式)
|
+--- Eino Agent(ReAct + Tools + Stream)
|
+--- LLM(OpenAI 兼容 / Mock)
每一层为什么存在:
- Controller 只做参数绑定与鉴权,不写业务
- Service 承载业务逻辑,按 gorm / redis / email / chat / kafka 分域
- DAO 只暴露数据库连接与事务
- chat/dispatch.go 是消息分发的公共核心,把「落库 -> 推送 -> 缓存 -> 触发 Agent」抽成一条 pipeline,channel 模式和 Kafka 模式共用
- agent/ 与 IM 主链路解耦:Agent 在独立协程里跑,25s 超时 + 背压保护,LLM 挂了不影响消息收发
我的模块
1. AI 流式链路(internal/agent/、internal/service/chat/、frontend/src/)
- internal/agent/service.go:接入 Eino ReAct Agent(MaxStep=8),构建上下文 [system(含摘要) + 历史 + 当前问题]
- 三级降级:StreamWithTools(带工具的流式)-> Stream(纯流式)-> Generate(一次性);Eino 初始化失败自动回退 Mock,保证无 API Key 也能演示
- 增量聚合:LLM 的 token 流按 80ms / 32 字符双阈值批量聚合后再推,避免逐 token 推帧把 WebSocket 打爆
- 流式帧协议:{"stream":"delta","stream_id":"Mxxx","content":"增量文本"},结束时发 done 帧;done 帧才落库 message 表并写 Redis List
- 前端:applyStreamMessage 拼接增量 + pendingMessages 离线队列,渲染打字机效果
- 可观测:middleware.IncAgentCall(costMs, isErr) 统计调用次数/错误率/平均耗时,GET /metrics 可见
2. 多实例集群一致性(internal/service/cluster/、internal/service/chat/kafka_server.go)
- 新增集群在线路由表 + 实例间消息总线(cluster/bus.go、registry.go、routing.go),让多实例部署时知道「这个用户在哪个实例上」
- KafkaServer 接入跨实例定向投递:消息先查路由表,命中本机直接推,命中别的实例走总线转发
- Kafka 分区键从「消息维度」改为按会话维度,保证同一会话的消息进同一分区、消费有序
- Producer acks 改为 RequireOne,Topic 改为经 controller 创建(原来依赖自动创建,生产环境容易踩坑)
- 修复 Renew 在 pipeline 中报 NOSCRIPT(Lua 脚本未预加载)
3. 稳定性与关停(internal/service/chat/server.go、client.go)
- 关停改用 done 信号 + context 控制,消除「多发送者 close channel」导致的 send on closed channel panic 和 goroutine 泄漏
- 关停不再清空整个 Redis,改成前缀白名单删除(原来一个实例退出会把别的实例的缓存也清掉)
4. 测试与 CI
- dispatch.go 从零测试补到 18 个用例,并借此发现 2 个真实缺陷
- 新增 6 个双实例端到端集成测试(跑真 Redis)、8 个关停语义测试、5 个分区键用例
- CI 增加 Redis 服务容器,让集群逻辑能在 CI 里跑真实集成测试;CI 失败时输出 annotation 摘要
- 前端 CI 增加构建校验
5. 工程化收尾
- 前端后端地址改为构建期注入,WS 地址自动派生,干掉硬编码 127.0.0.1
- Go 版本三处对齐到 1.25(原来 README 写 1.25、go.mod 写 1.26.1,是虚构版本)
- applyEnvOverrides 原是空壳函数(假装支持环境变量覆盖),真正实现并补测试
- 新增 docker-compose.cluster.yml + deploy/gateway.conf,多实例一键起
难点与取舍
1. LLM 逐 token 推帧 vs 批量聚合
- 问题:Eino 的 Stream 每产生一个 token 就回调一次,直接推 WebSocket 会导致帧数量爆炸(一句话几百帧),前端频繁 re-render,压测时 P99 直接劣化
- 方案 A:逐 token 推,前端自己节流 -> 前端复杂度上升,且网络开销没省
- 方案 B:服务端聚合,80ms 或 32 字符任一触发就推一帧 -> 延迟可控(最长多等 80ms),帧数降一个数量级
- 选 B。代价:首字延迟最多增加 80ms,需要把聚合窗口写进文档避免后人误改
2. 多实例怎么保证消息不丢
- 问题:单机时 Clients map 里有连接就能推;扩到多实例后,用户在 A 实例、消息从 B 实例发出,B 查不到连接就静默丢弃
- 方案 A:广播到所有实例,各自判断 -> 简单但 N 倍放大流量
- 方案 B:维护在线路由表 + 实例间定向总线 -> 精准投递,代价是要处理路由表一致性与实例上下线
- 选 B。路由表先用 Redis 实现(已有依赖,不引入 ZooKeeper/etcd),路由失效时降级为广播
- 代价:路由表有过期窗口,实例非正常退出时可能短暂投递失败
3. 关停时的 channel 语义
- 问题:原实现 Close() 里直接 close(channel),而多个 goroutine 都可能往这个 channel 写,触发 send on closed channel panic
- 方案 A:加锁 + 标志位判断 -> 能缓解,但仍有竞态窗口
- 方案 B:不用 close 传信号,改用 done channel + context.Cancel,写方通过 select { case <-done: return } 退出,由 sync.WaitGroup 等待收敛
- 选 B。这是 Go 并发里「用关闭表达结束」的经典坑,代价是代码里多一层 select
4. 群聊 Agent 触发词
- 问题:@AI助手 的匹配写在通用分支里,私聊/系统消息也会命中,导致 Agent 重复回复,还会生成伪群 ID
- 方案:把触发词判断限定在群消息分支内,触发词剥离后再传给 LLM
- 代价:无,纯缺陷修复。但值得讲的是「为什么写测试才发现」——dispatch.go 的 18 个用例就是这个背景下补的
5. 关停要不要清 Redis
- 问题:原实现关停时 Scan + Del 清空,看似「清理现场」,实际上多实例部署时会把其他实例的会话缓存一起删掉
- 方案:改成前缀白名单 + 默认关闭清理
- 取舍:本地开发时残留数据需要手动清,但生产正确性优先
量化结果
压测环境:4 核 8G,channel 模式,单实例,不含 LLM 调用(bench/k6-ws.js,发送间隔 1s):
| VU | 总消息 | P95 | P99 | Transmit 峰值 | 现象 |
|---|---|---|---|---|---|
| 100 | 6k/min | 45ms | 90ms | 12 | 正常 |
| 200 | 12k/min | 68ms | 130ms | 35 | 正常 |
| 500 | 30k/min | 110ms | 220ms | 98 | 接近满,偶发 429 |
| 1000 | 60k/min | 180ms | 350ms | 100 | 限流生效,需调大 CHANNEL_SIZE 或切 Kafka |
其他数字:
- Agent 与消息主链路解耦:LLM 调用超时 25s,走独立协程,不阻塞消息收发
- 流式聚合:帧数从「每 token 一帧」降到 80ms / 32 字符一帧
- 测试:dispatch.go 18 个用例(发现 2 个真实缺陷)、双实例集成测试 6 个、关停语义测试 8 个
如果重做
- 先写 dispatch 测试再改集群。我当时先实现路由表再补测试,结果 dispatch.go 的 18 个用例一跑就翻出 2 个既有缺陷;如果先补测试,集群改造会少走弯路。
- 路由表不放在 Redis 单点。现在路由表依赖 Redis,Redis 挂了集群就退化成广播。更稳的做法是引入 etcd/成员协议(gossip),或至少给路由表加本地缓存 + 心跳续期。
- Agent 记忆不该停在摘要。现在是「短期窗口 + 滚动摘要」,摘要是有损的。下一步应上向量检索做双层召回(摘要兜底 + 语义召回),代价是引入 embedding 成本和一致性维护。
- 压测要带 LLM。现在基准明确排除了 LLM,面试官很容易追问「加上 LLM 之后呢」。应该补一组 Mock LLM 的压测,把 Agent 链路的容量也说清楚。