- 后端
- 即时通讯
- 微服务
【免费下载链接】goim
goim
goim 是一个用 Go 语言实现的分布式实时消息推送系统,其 2.0.0 版本在架构上完成了一次重要的技术升级。本文以仓库CHANGELOG.md中官方记录的 11 项版本特性为主线,结合 internal/logic 与 internal/job 的源码实现,逐条解析每一项新特性的设计动机、底层实现与配置方式,帮助读者在读完本文后,既能理解 goim 2.0.0 的核心架构演进,也能基于示例配置快速搭建一套支持区域调度、按设备推送、多房间广播的推送服务。
一、版本背景:goim 2.0.0 是一次架构级重写
CHANGELOG.md记录了 goim 2.0.0 相对旧版的全部变更点,共 11 项。从变更内容看,这一版本不是简单的小步迭代,而是围绕「路由存储」、「服务发现」、「负载均衡」、「推送模型」四条主线的一次系统性重构:
| 变更主线 | 对应特性编号 | 核心关键词 |
|---|---|---|
| 路由与在线状态存储 | 1、2、10 | Redis、在线心跳、房间消息聚合 |
| 服务发现与 RPC | 3 | gRPC、Discovery |
| 节点调度策略 | 4、5、11 | 连接数与权重、区域调度、IPv6 |
| 连接与推送模型 | 6、7、8、9 | 指令订阅、房间切换、多房间类型、device_id 推送 |
下文按此主线展开,逐条对应CHANGELOG.md中的原始编号展开深度解读。
二、路由层重构:从 router 到 Redis(特性 1、2)
2.1 特性 1:路由表改为 Redis 存储
goim 1.x 时代,客户端连接路由(某个连接挂在哪个 comet 节点上)通常由独立的 router 模块管理。2.0.0 将这一职责直接下沉到 Redis,路由状态与在线状态统一由 Redis 承载,既减少了独立路由组件的部署复杂度,也让在线查询的读写时延显著降低。
从源码看,logic 模块在 internal/logic/dao/redis.go 中定义了三种 Redis key 前缀:
mid_%d:mid(用户 ID)到「连接 key → comet 服务器」映射的 Hash 结构;key_%s:连接 key 到 comet 服务器的映射;ol_%s:comet 服务器到其在线房间统计的映射。
路由写入的核心方法是 AddMapping:当客户端建立连接时,logic 会同时写入mid侧和key侧两路映射,并为每路设置EXPIRE过期时间(过期时长来自配置Redis.Expire,见 internal/logic/conf/conf.go)。删除时则对应执行HDEL/DEL(DelMapping)。
与之配套的单元测试 internal/logic/dao/redis_test.go 完整验证了 AddMapping → ExpireMapping → ServersByKeys → KeysByMids → DelMapping 这一整条路由生命周期,确认了「key → server」与「mid → key → server」两级映射的读写一致性。
2.2 特性 2:基于 Redis 的节点在线心跳维护
comet 节点需要周期性地把自身的在线统计上报给 logic,并维持「节点仍存活」的标记。该逻辑在 internal/logic/conn.go 的RenewOnline方法中体现:comet 通过 gRPC 调用RenewOnline上报各房间在线人数,logic 将其封装为Online结构后调用 AddServerOnline 写入 Redis,并同步刷新EXPIRE。
同时,logic 自身有一个后台协程onlineproc(internal/logic/logic.go),每 10 秒(_onlineTick)汇总一次各节点在线数据;当某个节点的上报时间超过 5 分钟(_onlineDeadline)未更新时,即判定节点失活,自动调用 DelServerOnline 清理其在线数据。这套「节点上报 + 定期汇总 + 超时清理」的机制,构成了无中心化心跳的在线状态维护方案。
值得一提的是,房间在线数据在上报时会先按cityhash.CityHash32对房间名取模分到 64 个 Hash 槽位,再批量写入,避免单 Key 过大(internal/logic/dao/redis.go)。
三、通信与发现层:gRPC + Discovery(特性 3)
2.0.0 全面引入 gRPC 作为内部服务间通信协议,并接入 bilibili 的 Discovery 实现服务注册与发现。
3.1 gRPC 服务化
logic 对外暴露的 gRPC 服务定义在 internal/logic/grpc/server.go,注册了五个核心 RPC:
Connect:连接建立,返回 mid、key、房间 ID、订阅指令列表与心跳参数;Disconnect/Heartbeat:连接断开与心跳续约;RenewOnline:节点在线数据上报;Receive:客户端上行消息接收;Nodes:向客户端返回可用的 comet 节点列表。
服务端创建时通过grpc.KeepaliveParams配置了连接空闲超时、最大连接生命周期与 keepalive 参数(internal/logic/grpc/server.go),对应配置项在 internal/logic/conf/conf.go 的RPCServer中定义。
3.2 Discovery 服务发现
logic 通过 internal/logic/logic.go 的initNodes使用discovery/naming构建goim.comet服务,watch 节点变化事件,一旦有变更立即调用newNodes刷新本地节点列表。节点实例携带的元数据(internal/logic/model/metadata.go)包括:
| 元数据 key | 含义 |
|---|---|
weight | 节点权重,参与负载均衡计算 |
offline | 是否下线,为 true 时被跳过 |
addrs | 节点公网地址列表(逗号分隔) |
ip_count | 节点在线 IP 数 |
conn_count | 节点在线连接数 |
Discovery 相关配置在启动时可通过环境变量或命令行参数注入:REGION、ZONE、DEPLOY_ENV、WEIGHT、HOST等(internal/logic/conf/conf.go),这些参数同时用于节点注册与 Discovery 上报。
四、节点调度:连接数、权重与区域调度(特性 4、5、11)
4.1 特性 4:连接数与权重调度
logic 内置了一个加权负载均衡器 internal/logic/balancer.go。其核心思想是:每个节点拥有固定权重(fixedWeight,来自 Discovery 元数据的weight),调度时按「权重占比 − 连接占比」的差值动态修正currentWeight,并始终选择currentWeight最大的节点:
- 权重占比 =
fixedWeight × gainWeight / totalWeight; - 连接占比 =
currentConns / totalConns × 0.5; currentWeight = fixedWeight + floor/ceil((权重占比 − 连接占比) × totalConns),并夹在_minWeight(1)与_maxWeight(1<<20)之间。
新连接到来时被选中的节点currentConns++,使得后续调度自动向连接数少的节点倾斜,实现「权重优先、连接数平滑」的调度效果。每次客户端请求节点列表时,至多返回前_maxNodes(5)个节点(internal/logic/balancer.go)。
4.2 特性 5:按区域(Region)调度
调度器支持区域亲和:当节点的region与客户端所属区域一致时,其权重会乘以gainWeight = regionWeight(由配置Node.RegionWeight指定,见 internal/logic/conf/conf.go),从而让用户优先接入就近机房。
区域映射由配置文件的Regions定义(省份 → 区域),logic 启动时通过 initRegions 构建「省份 → 区域」字典。客户端 IP 到省份的定位由 location 完成,当前仓库中该方法返回空字符串(需要接入外部 IP 库),因此实际区域命中与否取决于部署方的 IP 定位实现。
4.3 特性 11:支持 IPv6
节点元数据的addrs字段以逗号分隔形式支持同时携带 IPv4 与 IPv6 地址(internal/logic/balancer.go),NodeAddrs返回的地址列表可包含双栈地址;同时 logic 的 gRPC 监听地址、Redis/Kafka 连接地址均为字符串配置,天然兼容 IPv6 形式(如[::1]:3119),无需单独改造。
五、连接与推送模型升级(特性 6、7、8、9、10)
5.1 特性 6:指令订阅(Accepts)
客户端在建立连接时可声明自己「订阅哪些操作指令」,logic 将其原样记录并在Connect的 RPC 回复中回传(internal/logic/conn.go)。后续 job 向客户端推送消息时,comet 会依据该连接是否 accept 对应 operation 决定是否下发,从而实现客户端侧的消息选择性接收。
5.2 特性 7:当前连接房间切换
连接的 room 信息在建立时由 token 中的room_id决定(internal/logic/conn.go)。2.0.0 支持在同一连接生命周期内切换当前所在房间:comet 维护连接与房间的映射关系,切换时更新映射并同步刷新 logic 侧的房间在线计数,避免因切换房间导致消息错投或计数失真。
5.3 特性 8:多房间类型({type}://{room_id})
房间名采用统一的type://room_id编码格式,编码与解码逻辑在 internal/logic/model/room.go:
EncodeRoomKey(typ, room)生成如live://10086、video://233的房间 key;DecodeRoomKey(key)通过url.Parse还原出 scheme 与 host。
这一设计允许同一套推送框架承载直播、视频、聊天等多种业务类型的房间,且各类型房间互不干扰地独立广播。推送房间消息时直接使用编码后的 room key 作为 Kafka 消息 key(internal/logic/dao/kafka.go),保证同一房间的消息落到同一分区,从而保持有序性。
5.4 特性 9:按 device_id 推送
每个连接在建立时若未携带业务方生成的 key,logic 会为其生成 UUID 作为唯一连接标识(internal/logic/conn.go)。该 key 即设备维度(device_id)的推送寻址依据:
- 按 key 推送:PushKeys 先通过
ServersByKeys批量查询每个 key 所在 comet 节点,再按节点分组发送; - 按 mid 推送:PushMids 通过
KeysByMids查询 mid 下的所有 key 及其所在节点,实现「一个用户多端(多 key)同时触达」。
两种方式最终都通过 Kafka 将PushMsg(PUSH类型)投递到对应 comet 节点,见 internal/logic/dao/kafka.go。
5.5 特性 10:房间消息聚合
当需要向大量房间(甚至全量房间)广播时,若逐房间发送会产生海量 Kafka 消息。2.0.0 引入了房间聚合:logic 侧通过ROOM类型广播单房间,通过BROADCAST类型并携带speed(限速)参数广播全量消息(internal/logic/dao/kafka.go),job 消费端可据此对同类型房间消息进行合并下发,降低 comet 的重复 I/O。与此对应的推送入口函数为 PushRoom 与 PushAll,它们分别被 internal/logic/http/push.go 中的 HTTP 接口调用。
六、推送链路全景:从 HTTP 到客户端
综合上述特性,goim 2.0.0 的完整推送链路如下:
- 客户端通过 HTTP API(如 internal/logic/http/push.go)发起按 key / mid / room / all 的推送请求;
- logic 根据路由表(Redis)定位目标 comet 节点,将
PushMsg序列化为 protobuf 后写入 Kafka(internal/logic/dao/kafka.go); - job 消费 Kafka 消息,按
PUSH/ROOM/BROADCAST类型区分处理(internal/job/job.go、internal/job/push.go),必要时进行房间聚合,并通过 gRPC 转发给目标 comet; - comet 依据连接的房间与订阅指令(accepts)将消息最终推送给客户端。
Kafka 生产者配置了WaitForAll确认、最多 10 次重试等可靠性参数(internal/logic/dao/dao.go),确保消息不因瞬时抖动而丢失。
七、配置参考与启动说明
以上特性大多集中在 logic 模块,其完整配置项由 internal/logic/conf/conf.go 定义,启动时通过-conf指定 TOML 配置文件(默认logic-example.toml),示例配置见 cmd/logic/logic-example.toml。与本文特性相关的关键配置有:
| 配置节 | 关键项 | 对应特性 |
|---|---|---|
Node | Heartbeat、HeartbeatMax、RegionWeight、TCPPort/WSPort/WSSPort | 心跳(2)、区域调度(5) |
Discovery | Region、Zone、Env、Host | 服务发现(3) |
Redis | Expire、连接池与超时参数 | 路由与在线维护(1、2) |
Kafka | Topic、Brokers | 消息推送(9、10) |
Regions | 省份 → 区域映射 | 区域调度(5) |
Backoff | MaxDelay、BaseDelay、Factor、Jitter | 客户端重连退避(3) |
部署时可通过环境变量REGION、ZONE、DEPLOY_ENV、WEIGHT覆盖默认的节点注册参数;comet、job 的启动方式与对应示例配置分别见 cmd/comet/comet-example.toml、cmd/job/job-example.toml 及各自 cmd 下的入口文件。
八、总结
goim 2.0.0 以 Redis 承载连接路由与在线状态、以 gRPC + Discovery 完成服务间通信与节点发现、以加权调度器实现连接数与区域感知的负载均衡,并在推送模型上支持指令订阅、房间切换、多房间类型、按设备/用户寻址与房间聚合。这 11 项变更共同构成了一个更适合大规模、多区域、多业务类型部署的实时推送基础设施。读者可依据本文的源码线索(internal/logic 各文件)结合仓库中的单元测试(如 internal/logic/dao/redis_test.go、internal/logic/dao/kafka_test.go)继续深入验证每一项特性的实际行为。
- 后端
- 即时通讯
- 微服务
【免费下载链接】goim
goim
相关推荐
doc_wei/erp-pro:消息推送服务深度解析
doc_wei/erp pro:消息推送服务深度解析 在企业级应用开发中,消息推送服务是连接用户与系统的重要桥梁。doc_wei/erp pro项目基于Spri
后端低代码企业应用AI应用工业制造工作流自动化解决comin常见问题:从仓库认证失败到部署回滚的7个实用技巧
解决comin常见问题:从仓库认证失败到部署回滚的7个实用技巧 如果你正在使用comin进行NixOS的GitOps部署,可能会遇到各种挑战。comin作为Ni
如何快速部署seresnext50_32x4d.racm_in1k:10步完整教程
如何快速部署seresnext50_32x4d.racm_in1k:10步完整教程 seresnext50_32x4d.racm_in1k是一款基于Squeez
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考