☰
goim 集群部署与推送协议实战指南:Go 语言实现的 IM 与实时推送服务
2026/9/28 2:45:19 网站建设 项目流程
  • 后端
  • 即时通讯
  • 微服务

【免费下载链接】goim

goim

项目地址:https://gitcode.com/gh_mirrors/go/goim
点击查看免费下载

goim(Terry-Mao/goim)是一个用纯 Go 语言编写、支持集群部署的 IM(即时通讯)与推送通知服务器,能够承载单用户推送、批量推送、房间推送和全量广播等场景。本文以仓库的 README_en.md 为主线,完整讲解其核心特性、依赖安装、集群部署、配置参数,并结合 推送 HTTP 协议 与 客户端长连接协议 给出可直接复用的验证与接入方法,读完即可动手搭建一套属于自己的推送集群。

项目概览与核心特性

goim 定位为轻量、高性能的 IM 及实时推送服务集群。官方英文文档将其能力归纳为以下要点,每条都对应仓库中实际的代码实现:

特性说明源码佐证
轻量级、高性能、纯 Golang服务端全部由 Go 编写,无重型运行时依赖全部代码位于 internal 与 cmd 目录
支持单推、多推、房间推、广播提供 4 类 HTTP 推送接口internal/logic/http/push.go
单 Key 多订阅者,可限制最大订阅数一个用户 ID(mid)可在多个设备/通道同时在线internal/logic/dao/redis.go 中的 mid→key 映射
心跳支持(应用心跳、TCP、KeepAlive)客户端与服务器保持长连接活性api/protocol/operation.go 中OpHeartbeat与OpHeartbeatReply
支持鉴权(未授权用户不能订阅)连接必须携带 token 完成认证internal/comet/server_websocket.go 的authWebsocket
多协议支持(WebSocket、TCP)客户端接入层同时监听两套协议internal/comet/server_tcp.go 与 internal/comet/server_websocket.go
可伸缩架构job、logic 模块可动态无限扩展internal/job/job.go 通过 Discovery 动态发现 comet 节点
基于 Kafka 的异步推送推送消息先入 Kafka,再由 job 异步消费分发internal/logic/dao/kafka.go

架构上,goim v2.0 由logic(逻辑层)、comet(长连接接入层)、job(推送任务层)三类无状态/可扩展节点构成,配合Kafka(异步消息队列)、Redis(在线状态与 key 映射)、Discovery(服务注册发现)三个基础设施组件协同工作:

安装依赖:从零准备运行环境

原文档给出的安装流程面向 Linux 服务器(以yum包管理为例),共分为三步,其中前两步是 goim 自身运行的前提:

1. 安装 Java 运行环境

goim 本身是纯 Go 服务,不需要 JVM;安装 Java 是因为Kafka 依赖 ZooKeeper/JVM才能运行:

$ yum -y install java-1.7.0-openjdk

若使用较新版本的 Kafka,建议按所选 Kafka 版本要求选择对应的 JDK 版本(例如 JDK 8),仓库也提供了 scripts/jdk8.sh 与 scripts/kafka.sh、scripts/zk.sh 等辅助脚本可供参考。

2. 安装并启动 Kafka

Kafka 是整个异步推送链路的“传输带”:logic 把推送消息序列化后写入 Kafka 的goim-push-topic主题,job 从该主题消费后转发给 comet。安装与启动请参照 Kafka 官方 QuickStart 文档,仓库中 internal/logic/dao/kafka.go 使用 sarama 客户端发送消息,internal/job/job.go 使用sarama-cluster以消费者组方式拉取消息。启动 Kafka 后默认监听127.0.0.1:9092,与各示例配置中的brokers保持一致。

3. 安装 Go 语言环境

按官方安装包安装 Go 后配置环境变量:

export GOROOT=/usr/local/go export PATH=$PATH:$GOROOT/bin export GOPATH=/data/apps/go source /etc/profile

当前仓库为 Go Module 工程(见 go.mod),因此并不强制依赖 GOPATH 目录结构,可以直接在任意目录 clone 后构建。

部署 goim 集群

原文档描述了早期版本基于$GOPATH/src/goim与router模块的部署方式;当前仓库已演进为 goim v2.0(见 README.md 顶部),模块精简为 logic、comet、job 三个,并以 Discovery 取代了有状态的 router 节点。下面先完整保留原文档的经典部署流程,再给出当前仓库的实际部署方式。

方式一:经典流程(对应原英文文档)

$ yum install git $ cd $GOPATH/src $ git clone https://github.com/Terry-Mao/goim.git $ cd $GOPATH/src/goim $ go get ./...

然后逐个编译并安装 router、logic、comet、job 模块(配置文件需按实际机器环境调整):

$ cd $GOPATH/src/goim/router $ go install $ cp router-example.conf $GOPATH/bin/router.conf $ cp router-log.xml $GOPATH/bin/ $ cd ../logic/ $ go install $ cp logic-example.conf $GOPATH/bin/logic.conf $ cp logic-log.xml $GOPATH/bin/ $ cd ../comet/ $ go install $ cp comet-example.conf $GOPATH/bin/comet.conf $ cp comet-log.xml $GOPATH/bin/ $ cd ../logic/job/ $ go install $ cp job-example.conf $GOPATH/bin/job.conf $ cp job-log.xml $GOPATH/bin/

注意:旧版router目录与*-example.conf、*-log.xml文件在 v2.0 仓库中已不存在,此流程适合从早期版本升级/参考的场景;v2.0 请使用下面的方式二。

方式二:v2.0 仓库实际构建方式(推荐)

当前仓库提供了统一的 Makefile 构建与启动入口(见 Makefile):

make build # 编译三个模块并生成 target/ 目录

make build会创建target/目录,将 cmd/comet/comet-example.toml、cmd/logic/logic-example.toml、cmd/job/job-example.toml 复制为target/*.toml,并编译出target/comet、target/logic、target/job三个可执行文件。

启动时既可以用make run一键拉起,也可以手动指定 flag 参数(region/zone/deploy.env用于 Discovery 注册,weight用于负载均衡权重,comet 还需addrs指定对外 IP):

nohup target/logic -conf=target/logic.toml -region=sh -zone=sh001 -deploy.env=dev -weight=10 2>&1 > target/logic.log & nohup target/comet -conf=target/comet.toml -region=sh -zone=sh001 -deploy.env=dev -weight=10 -addrs=127.0.0.1 2>&1 > target/comet.log & nohup target/job -conf=target/job.toml -region=sh -zone=sh001 -deploy.env=dev 2>&1 > target/job.log &

停止服务使用make stop。如果启动失败,请查看对应进程的日志文件(如panic-*.log)定位问题。

环境变量与命令行参数

每个模块的配置解析都位于各自的conf/conf.go中,支持“环境变量优先、命令行 flag 覆盖”的取值方式:

参数环境变量说明
-regionREGION可用区域,如sh
-zoneZONE可用分区,如sh001/sh002
-deploy.envDEPLOY_ENV部署环境,如dev/fat1/uat/pre/prod
-host机器 hostname机器主机名,默认取os.Hostname()
-weightWEIGHT负载均衡权重
-addrs(仅 comet)ADDRS服务器对外公网 IP,多个用逗号分隔
-offline(仅 comet)OFFLINE节点下线标记
-debug(仅 comet)DEBUG调试日志开关

例如 supervisord 托管时可在配置中写入environment=REGION=sh,ZONE=sh001,DEPLOY_ENV=dev。详见 internal/comet/conf/conf.go、internal/logic/conf/conf.go、internal/job/conf/conf.go。

配置文件逐项详解

原英文文档中Configurations章节标注为 TODO,本节基于仓库真实配置文件补齐完整参数说明,这也是把 goim 部署到生产环境的必修课。

comet 接入层:cmd/comet/comet-example.toml

[discovery] nodes = ["127.0.0.1:7171"] # Discovery 注册中心地址 [rpcServer] addr = ":3109" # comet 对外 gRPC 服务地址(job 推送用) timeout = "1s" [rpcClient] dial = "1s" # 连接 logic 的拨号超时 timeout = "1s" [tcp] bind = [":3101"] # TCP 长连接监听地址(可配多个) sndbuf = 4096 # 发送缓冲区大小(字节) rcvbuf = 4096 # 接收缓冲区大小(字节) keepalive = false # 是否开启 TCP KeepAlive reader = 32 # Reader 协程数(配合 round 池) readBuf = 1024 readBufSize = 8192 writer = 32 # Writer 协程数 writeBuf = 1024 writeBufSize = 8192 [websocket] bind = [":3102"] # WebSocket 监听地址 tlsOpen = false # 是否开启 TLS tlsBind = [":3103"] # WSS 监听地址 certFile = "../../cert.pem" # 证书路径 privateFile = "../../private.pem" # 私钥路径 [protocol] timer = 32 # 定时器分片数 timerSize = 2048 # 每个分片的定时器容量 svrProto = 10 # 服务器侧 proto 环形队列容量 cliProto = 5 # 客户端侧 proto 环形队列容量 handshakeTimeout = "8s" # 握手超时时间 [whitelist] Whitelist = [123] # 白名单用户 mid,用于调试 WhiteLog = "/tmp/white_list.log" # 白名单日志输出文件 [bucket] size = 32 # bucket 数量(影响并发分片与在线分布) channel = 1024 # 每个 bucket 的 channel 容量 room = 1024 # 每个 bucket 的房间容量 routineAmount = 32 # 广播推送协程数 routineSize = 1024 # 广播推送协程缓冲队列大小
  • bucket是 comet 的核心数据分片:每个连接(channel)按 key 哈希落入某个 bucket,广播推送时并行向所有 bucket 分发,见 internal/comet/bucket.go 与 internal/comet/grpc/server.go 中PushMsg/Broadcast/BroadcastRoom的实现。
  • svrProto/cliProto对应 internal/comet/channel.go 中客户端上行与服务器下行两个环形缓冲队列的深度,直接影响高并发下的吞吐与背压。
  • 示例证书 examples/cert.pem 与 examples/private.pem 可配合tlsOpen = true做本地 WSS 验证。

logic 逻辑层:cmd/logic/logic-example.toml

[discovery] nodes = ["127.0.0.1:7171"] [regions] # 区域 -> 省份映射,用于节点调度权重 "bj" = ["北京","天津","河北","山东","山西","内蒙古","辽宁","吉林","黑龙江","甘肃","宁夏","新疆"] "sh" = ["上海","江苏","浙江","安徽","江西","湖北","重庆","陕西","青海","河南","台湾"] "gz" = ["广东","福建","广西","海南","湖南","四川","贵州","云南","西藏","香港","澳门"] [node] defaultDomain = "conn.goim.io" # 默认接入域名 hostDomain = ".goim.io" # 接入域名后缀 heartbeat = "4m" # 应用层心跳间隔 heartbeatMax = 2 # 心跳最大容忍次数(超过则判定离线) tcpPort = 3101 # 对外 TCP 端口(comet 暴露) wsPort = 3102 # 对外 WebSocket 端口 wssPort = 3103 # 对外 WSS 端口 regionWeight = 1.6 # 同区域权重 [backoff] # 节点重连退避策略 maxDelay = 300 baseDelay = 3 factor = 1.8 jitter = 0.3 [rpcServer] network = "tcp" addr = ":3119" # logic 的 gRPC 服务地址(comet 连接) timeout = "1s" [rpcClient] dial = "1s" timeout = "1s" [httpServer] network = "tcp" addr = ":3111" # HTTP 推送接口监听端口 readTimeout = "1s" writeTimeout = "1s" [kafka] topic = "goim-push-topic" # 推送消息主题 brokers = ["127.0.0.1:9092"] # Kafka 集群地址 [redis] network = "tcp" addr = "127.0.0.1:6379" # Redis 地址 active = 60000 # 最大活跃连接数 idle = 1024 # 最大空闲连接数 dialTimeout = "200ms" readTimeout = "500ms" writeTimeout = "500ms" idleTimeout = "120s" expire = "30m" # mid/key 映射过期时间

node段中的tcpPort/wsPort/wssPort必须与 comet 侧 [tcp]/[websocket] 的监听端口一一对应,否则客户端通过域名接入时会连错端口。redis的expire = 30m决定了在线映射的有效期,它需要大于node.heartbeat * heartbeatMax(即 4m × 2 = 8m),否则心跳周期稍长就会被判定离线。

job 任务层:cmd/job/job-example.toml

[discovery] nodes = ["127.0.0.1:7171"] [kafka] topic = "goim-push-topic" # 与 logic 保持同一个主题 group = "goim-push-group-job" # 消费者组 ID brokers = ["127.0.0.1:9092"]

job 以消费者组方式消费 Kafka 中的推送消息,因此多个 job 节点使用同一个group即可自动分区负载均衡;同时 job 会通过 Discovery 监听goim.comet服务的实例变化,动态建立/销毁到各 comet 节点的 gRPC 连接(见 internal/job/job.go 的watchComet/newAddress)。

启动服务与日志排查

参照原文档,启动后各进程的崩溃日志可分别落到独立文件,方便按模块排查:

$ cd /$GOPATH/bin $ nohup $GOPATH/bin/router -c $GOPATH/bin/router.conf 2>&1 > /data/logs/goim/panic-router.log & $ nohup $GOPATH/bin/logic -c $GOPATH/bin/logic.conf 2>&1 > /data/logs/goim/panic-logic.log & $ nohup $GOPATH/bin/comet -c $GOPATH/bin/comet.conf 2>&1 > /data/logs/goim/panic-comet.log & $ nohup $GOPATH/bin/job -c $GOPATH/bin/job.conf 2>&1 > /data/logs/goim/panic-job.log &

启动顺序建议为:Kafka → Redis/Discovery → logic → comet → job。如果失败,优先检查对应panic-*.log中的连接错误(Kafka broker 不可达、Discovery 未注册、端口被占用等)。

验证推送:HTTP 推送接口实战

原文档将 推送 HTTP 协议 作为部署完成后的第一道验收步骤。logic 的 HTTP 服务监听:3111,路由注册见 internal/logic/http/server.go(实际路径前缀为/goim,如/goim/push/room,与旧文档路径有所差异)。

接口一览

名称URLHTTP 方法
single push/1/pushPOST
multiple push/1/pushsPOST
room push/1/push/roomPOST
broadcasting/1/push/allPOST

公共响应体

响应码说明
1成功
65535内部错误
{ "ret": 1 }

单用户推送(single push)

uid即推送目标用户 id,消息体为任意 JSON:

curl -d "{\"test\":1}" "http://127.0.0.1:7172/1/push?uid=0"
{"ret": 1}

批量推送(multiple push)

u为 uid 数组,m为消息体,可一次覆盖多个在线用户:

curl -d "{\"u\":[1,2,3,4,5],\"m\":{\"test\":1}}" "http://127.0.0.1:7172/1/pushs"
{"ret": 1}

房间推送(room push)

rid为房间 id,消息将发给该房间内的所有在线连接:

curl -d "{\"test\": 1}" "http://127.0.0.1:7172/1/push/room?rid=1"
{"ret": 1}

全量广播(broadcasting)

curl -d "{\"test\": 1}" "http://127.0.0.1:7172/1/push/all"
{"ret": 1}

推送链路源码级解读

在 internal/logic/http/push.go 中,四个接口分别调用PushKeys/PushMids/PushRoom/PushAll(见 internal/logic/push.go):

  • 按 key 推送:PushKeys先从 Redis 用MGET批量查出每个 key 所在的 comet 服务器(internal/logic/dao/redis.go 的ServersByKeys),按服务器分组后写入 Kafka;
  • 按 mid 推送:PushMids通过KeysByMids用HGETALL取出每个 mid 下的全部在线 key,实现“单用户多设备同时收到消息”;
  • 房间/广播:PushRoom与PushAll直接构造pb.PushMsg(类型分别为ROOM、BROADCAST)写入 Kafka,由 job 消费后调用 comet 的 gRPC 接口(见 internal/comet/grpc/server.go)完成BroadcastRoom或带Speed限速的Broadcast分发。

整个链路为logic → Kafka → job → comet(grpc) → channel/room → 客户端,推送请求在 logic 处即返回,实际下发是异步完成的,这正是其高吞吐的核心设计。

客户端长连接协议

原文档链接的 Comet 客户端协议 说明了客户端如何接入长连接层。comet 同时支持 WebSocket 与 TCP 两套协议,接入后即可收到上述推送。

WebSocket 协议

请求地址:ws://DOMAIN/sub

以 JSON Frame 通信,请求与响应同构:

{ "ver": 102, "op": 10, "seq": 10, "body": {"data": "xxx"} }
参数是否必填类型说明
vertrueint协议版本号
optrueint操作码
seqtrueint序列号(服务器返回时对应客户端发送的 seq)
body-json推送的 JSON 消息体

在服务端,comet 要求 WebSocket 的请求路径必须为/sub(internal/comet/server_websocket.go 中校验req.RequestURI != "/sub"即断开),握手完成后进入鉴权流程authWebsocket:客户端必须发送OpAuth操作码完成认证,服务器回OpAuthReply,否则无法订阅。

TCP 协议

请求地址:tcp://DOMAIN

采用二进制帧,响应结构同请求:

参数是否必填类型说明
package lengthtrueint32 大端包总长度
header lengthtrueint16 大端头部长度
vertrueint16 大端协议版本
operationtrueint32 大端操作码
seqtrueint32 大端序列号(对应 jsonp 回调)
bodyfalsebinary消息体,长度为package length - header length

操作码(Operations)

原文档表格给出了 4 个基础操作码,仓库中 api/protocol/operation.go 定义了完整的操作码集合:

操作码含义源码常量
2客户端发送心跳OpHeartbeat
3服务器回复心跳OpHeartbeatReply
7鉴权请求OpAuth
8鉴权响应OpAuthReply
0/1握手请求/响应OpHandshake/OpHandshakeReply
4/5发送消息/响应OpSendMsg/OpSendMsgReply
9原始消息OpRaw
10/11proto 就绪/结束OpProtoReady/OpProtoFinish
12/13切换房间/响应OpChangeRoom/OpChangeRoomReply
14/15订阅/响应OpSub/OpSubReply
16/17取消订阅/响应OpUnsub/OpUnsubReply

心跳处理逻辑可在 internal/comet/server_websocket.go 中看到:客户端发来OpHeartbeat后,服务端改写为OpHeartbeatReply并重置定时器,同时周期性通过 gRPC 向 logic 续报在线状态。

客户端示例与参考文档

仓库自带一个基于静态文件的 WebSocket 客户端 Demo(examples/javascript 目录):通过main.go在:1999端口提供静态服务,浏览器打开 examples/javascript/index.html 即可连接 comet 并实时观察推送消息(消息推送逻辑见 examples/javascript/client.js):

$ cd examples/javascript $ go run main.go

更多权威参考:

  • 推送 HTTP 协议中文文档 与 英文文档
  • Comet 客户端协议中文文档 与 英文文档
  • 协议结构图 与 握手流程说明
  • 中文总览见 README_cn.md,英文总览见 README_en.md

至此,你已经掌握了 goim 的依赖准备、集群部署、配置调优、推送接口验证与客户端接入的完整闭环,可以在自己的服务器上搭建一套可横向扩展的实时推送系统。

  • 后端
  • 即时通讯
  • 微服务

【免费下载链接】goim

goim

项目地址:https://gitcode.com/gh_mirrors/go/goim
点击查看免费下载
上一篇:Newton中的触觉反馈:接触力计算与应用全解析
下一篇:react-native-code-push性能分析报告:更新流程各阶段耗时优化

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询