即时通讯 IM 服务

2023-04-14T14:11:02+08:00 | 19分钟阅读 | 更新于 2026-04-14T14:11:02+08:00

@

学习目标

学完本章,你应该能够:

  1. 讲清楚 IM 系统为什么比普通 Web 服务难:长连接维护、分布式下消息投递这两座大山。
  2. 对比 XMPP / WebSocket / WebRTC 三种协议,说清各自适用场景与选型理由。
  3. 讲透 WebSocket 原理:握手升级、报文结构、长连接保活,并能动手起一个 WebSocket 服务。
  4. 设计分布式 IM 的跨节点转发:说清注册机制 vs 广播机制的取舍,以及"网关 + 后端 + Kafka"为什么是生产标配。
  5. WebSocket + Kafka 搭一个最简群聊,理解 ACK、可靠投递、顺序性、离线消息这些面试深度题。

前置知识(如果下面任一点生疏,先回看对应章):

  • 第07章 Kafka:知道 Topic / 分区 / 消费组(群聊广播要用到"每节点独立消费组")。
  • 第12-14章 服务治理:知道注册发现、负载均衡(理解分布式转发的前提)。
  • 基本 TCP/HTTP 知识:懂"握手"“长连接"的字面意思即可。
  • Go 并发基础:goroutinesync.Map/sync.Mutex(Hub 维护连接要用)。

本章你会动手做的事

  • gorilla/websocket 起一个会定时推送时间戳的 WebSocket 服务,用 wscat 或浏览器连上去看效果。
  • 把单节点广播的 Hub 写出来,体会"用户上线注册、下线注销、发消息广播”。
  • 用 Kafka 把"网关接收 → 广播给所有节点"的群聊链路跑通,验证每个网关节点都是独立消费组。

一、IM 简介

IM(Instant Message,即时通讯)是一种通过网络实时传递文本、多媒体内容和文件的通信方式,使用户能够在几乎瞬间与他人交流,消除时间和地理限制。

1.1 IM 的核心特征

  • 即时性:消息实时传递,依赖实时通信协议(XMPP / WebSocket / WebRTC)。
  • 多媒体支持:文本、图片、音频、视频、文件,不同类型消息在展示和存储上要求不同,是实践中难处理的点。
  • 跨平台:移动端(iOS/Android)、桌面端(Windows/macOS/Linux)、Web 端多端同步。

白话类比:普通 Web 服务像"你寄信、对方某天去信箱翻";IM 像"你一开口,对方手机立刻响"。难点不在"存消息",而在"怎么让消息秒到、且对方在哪一目了然"。这就是长连接和实时推送要解决的问题。

1.2 常见 IM 工具

工具特点
WhatsApp端到端加密、全球覆盖、轻量
微信多功能平台,融合社交/支付/小程序
Telegram开放 API、隐私保护、支持机器人
Slack团队协作、频道、第三方集成
Microsoft Teams与 Office 365 深度集成

1.3 IM 系统的完整需求

一个完整 IM(如微信)要解决的问题非常复杂:

  • 实时聊天(文本、图片、多媒体消息)
  • 群组功能(创建群、邀请、群管理)
  • 权限控制(能否加群、能否看资源)
  • 用户模块(注册、登录、资料管理)
  • 用户关系(加好友、拉黑、屏蔽)
  • 消息记录 / 历史消息
  • 内容审核(国内合规要求)
  • 搜索(文件、聊天记录、群、用户)
  • 多端信息同步

本课程聚焦最核心的实时聊天功能。其余功能本质是增删改查,但组合在一起做成 IM,难度不亚于做一个微博或小红书。

flowchart TD
    RT[实时聊天
本课重点] --> G[群组/权限] RT --> U[用户/关系] RT --> H[历史/审核] RT --> S[搜索/多端同步]

这张图在讲:IM 的"冰山"——露在水面的是实时聊天,水面下是群组、关系、历史、审核等一大堆系统工程。本课只抠最顶上的实时聊天。


二、IM 核心协议对比

2.1 XMPP

XMPP(Extensible Messaging and Presence Protocol)是基于 XML 的开放通信协议。

  • 特点:开放标准、支持扩展、多终端同步、即时消息与在线状态管理。
  • 定位:最初专为 IM 设计,但如今影响力下降。
  • 适用:需要和其他 IM 互通的场景。官网 https://xmpp.org/

2.2 WebSocket

当前主流 IM 基本都构建在 WebSocket 上。

  • 双向通信:全双工,服务器可主动推送。
  • 低延迟:实时性高,适合对延迟敏感的应用。
  • 跨域支持:允许跨域建立连接。
  • 应用范围:IM、在线游戏、实时协作、Web 应用实时数据传输。

2.3 WebRTC

WebRTC 专注于浏览器间实时音视频通信。

  • 实时音视频:直接在浏览器中进行音视频通话。
  • 点对点通信:更加直接高效。
  • 端到端加密:安全性好。
  • 适用:Web 会议、在线教育、视频聊天。官网 https://webrtc.org/

2.4 三者对比

维度XMPPWebSocketWebRTC
通信类型文本消息+在线状态全双工通用通信实时音视频点对点
灵活性高(开放标准,可扩展)中(简洁但不够灵活)低(专用)
实时性较高(受轮询影响)高(全双工低延迟)高(取决于网络)
安全性TLS/SSL加密通信端到端加密
flowchart LR
    A[XMPP
文本/状态 互通] --> IM[IM 系统] B[WebSocket
全双工主流] --> IM C[WebRTC
音视频] --> IM

这张图在讲:三种协议各有主场——要互通选 XMPP,普通消息主流用 WebSocket,音视频用 WebRTC。本课程主角是 WebSocket。


三、WebSocket 原理详解

3.1 为什么用 WebSocket 而不是轮询

在没有 WebSocket 之前,前端等待后端结果只能用轮询(如打赏结果查询)。轮询缺点明显:

  • 流量放大:大量无效请求加重后端负担。
  • 雪崩风险:服务端压力大时无法正常响应,前端进一步轮询,形成恶性循环。

WebSocket 通过持久连接 + 全双工通信彻底解决这个问题。

注意:持久连接是双刃剑。一个服务端很难撑住大量 WebSocket 连接,如果每个连接通信频繁就更难。Go 在压测下表现并非最优,需要调优。

⚠️ 新手必踩的坑:长连接很"贵"。一个 HTTP 请求完了就释放,连接数可以很大;但一个 WebSocket 连接是"常驻"的,占内存、占文件描述符。几千个连接还好,几十万上百万就对单机是巨大考验。所以 IM 网关要独立部署、要横向扩展,不是随便一个业务服务就能扛的。

3.2 WebSocket 初始化过程

WebSocket 的握手是在 HTTP 基础上协商升级协议:

  1. 客户端发送 HTTP 升级请求,询问能否用 WebSocket 通信。
  2. 服务端答复可以,协议升级为 WebSocket。
  3. 两端使用 WebSocket 持续通信。
  4. 任意一端发起关闭即可中止通信。

客户端升级请求关键头部

GET /ws HTTP/1.1
Upgrade: websocket
Connection: Upgrade
Sec-WebSocket-Key: dGhlIHNhbXBsZSBub25jZQ==    <!-- 握手合法性验证 -->
Sec-WebSocket-Protocol: chat                     <!-- 子协议 -->
Sec-WebSocket-Version: 13                        <!-- 协议版本 -->
Origin: https://example.com                      <!-- 跨域控制 -->

服务端响应关键头部

HTTP/1.1 101 Switching Protocols
Upgrade: websocket
Connection: Upgrade
Sec-WebSocket-Accept: s3pPLMBiTxaQ9kYGzzhZRbK+xOo=   <!-- 由 Key 计算得出 -->
Sec-WebSocket-Protocol: chat                           <!-- 选定子协议 -->
sequenceDiagram
    participant C as 客户端
    participant S as 服务端
    C->>S: GET /ws + Upgrade: websocket
    S-->>C: 101 Switching Protocols
    Note over C,S: 之后走 WebSocket 全双工帧
    C->>S: 文本/二进制帧
    S-->>C: 主动推送帧

这张图在讲:WebSocket 先用一次 HTTP 握手(101 切换协议),之后就脱离 HTTP,双方都能随时发帧,服务端也能主动推——这是它和请求/响应模型最大的不同。

3.3 WebSocket 报文(帧)结构

字段长度含义
FIN1 bit是否消息最后一个片段
RSV3 bits保留字段
Opcode4 bits操作码(文本帧/二进制帧等)
Mask1 bitPayload 是否经过掩码处理
Payload Length7 / 7+16 / 7+64 bits负载长度(动态扩展)
Masking Key4 bytes仅 Mask=1 时存在
Payload Data任意长度实际消息内容

3.4 长连接的两种含义

“长连接"是含糊的说法,可能指:

  1. TCP 长连接:TCP 协议栈自身的 Keep-Alive 机制,定时发送保活报文,发现对端无响应就关闭连接。
  2. 应用层长连接:如 HTTP Connection: Keep-Alive,RPC 心跳请求等,是应用层约定复用 TCP 连接。

WebSocket 的长连接兼具两层含义:底层 TCP Keep-Alive + 应用层心跳。


四、WebSocket API 实战(Go)

4.1 最简 WebSocket 服务

使用 github.com/gorilla/websocket 库:

package ws

import (
    "fmt"
    "net/http"
    "time"

    "github.com/gorilla/websocket"
)

// Upgrader 用于把 HTTP 连接升级为 WebSocket
var upgrader = websocket.Upgrader{
    ReadBufferSize:  1024, // 读缓冲区大小
    WriteBufferSize: 1024, // 写缓冲区大小
    // CheckOrigin: 检查 Origin,防范跨域攻击。生产环境要严格校验
    CheckOrigin: func(r *http.Request) bool {
        return true // 演示用,生产环境应该校验 Origin 白名单
    },
    EnableCompression: false, // 是否启用压缩
}

// HandleWS 处理 WebSocket 连接
func HandleWS(w http.ResponseWriter, r *http.Request) {
    // 步骤 1:升级 HTTP 为 WebSocket
    conn, err := upgrader.Upgrade(w, r, nil)
    if err != nil {
        fmt.Printf("upgrade failed: %v\n", err)
        return
    }
    defer conn.Close()

    // 步骤 2:启动一个 goroutine 定时往客户端推送数据
    go func() {
        for {
            err := conn.WriteMessage(websocket.TextMessage,
                []byte(fmt.Sprintf(`{"ts":%d}`, time.Now().Unix())))
            if err != nil {
                return
            }
            time.Sleep(time.Second)
        }
    }()

    // 步骤 3:主循环不断读取客户端发来的数据
    for {
        _, msg, err := conn.ReadMessage()
        if err != nil {
            fmt.Printf("read failed: %v\n", err)
            return
        }
        fmt.Printf("received: %s\n", msg)
        // 实际业务中转交给 handler 处理
    }
}

⚠️ 新手必踩的坑:CheckOrigin 别直接 return true。演示为了方便关掉了跨域校验,但生产上这等于任何人都能从任意网站连你的 WebSocket(CSWSH 攻击)。务必校验 Origin 白名单。另外,读和写要用不同的 goroutineconn.ReadMessageconn.WriteMessage 对同一个连接并发调用是安全的,但多个 goroutine 同时写同一个 conn 会出问题——写要串行化或用单独写 goroutine。

4.2 BufferPool 优化

WriteBufferPool 是一种对象池,避免每次连接都新分配 Buffer 内存。这是常见的性能优化措施:

var upgrader = websocket.Upgrader{
    WriteBufferPool: &sync.Pool{
        New: func() interface{} {
            return make([]byte, 1024)
        },
    },
}

4.3 多客户端协调(单节点广播)

最简单的 IM 场景:A 发消息,服务端要转发给同样连上来的 B 和 C。

package ws

import "sync"

// Hub 维护所有在线连接,实现广播
type Hub struct {
    mu    sync.RWMutex
    conns map[int64]*websocket.Conn // uid -> conn
}

func NewHub() *Hub {
    return &Hub{conns: make(map[int64]*websocket.Conn)}
}

// Register 用户上线时注册连接
func (h *Hub) Register(uid int64, conn *websocket.Conn) {
    h.mu.Lock()
    defer h.mu.Unlock()
    h.conns[uid] = conn
}

// Unregister 用户下线时移除
func (h *Hub) Unregister(uid int64) {
    h.mu.Lock()
    defer h.mu.Unlock()
    // 注意:原示例 delete(h, uid) 是笔误,应为 delete(h.conns, uid)
    delete(h.conns, uid)
}

// Broadcast 广播消息,skipUid 用于不转发给自己
func (h *Hub) Broadcast(msg []byte, skipUid int64) {
    h.mu.RLock()
    defer h.mu.RUnlock()
    for uid, conn := range h.conns {
        if uid == skipUid {
            continue // 不转发给自己
        }
        // 注意:实际生产中写入要异步+超时控制,避免慢消费者拖垮整个 hub
        _ = conn.WriteMessage(websocket.TextMessage, msg)
    }
}
flowchart TD
    A[用户 A 上线] --> R[Hub.Register A]
    B[用户 B 上线] --> R2[Hub.Register B]
    A -- 发消息 --> BC[Hub.Broadcast]
    BC --> W[写 B 连接]
    BC -. 跳过 A 自己 .-> X[不写 A]

这张图在讲:单节点下,Hub 用 map[uid]conn 记下所有在线连接,广播就是"遍历 map 挨个写”,并跳过发送者自己。

4.4 sync.Map 与 syncx.Map

Go 内置 map 不是线程安全的,并发场景要用 sync.Map。核心操作:

  • Store:写入
  • Load:加载
  • LoadOrStore:有则加载,无则写入
  • CompareAndSwap:CAS 修改

课程中使用 ekit 库的 syncx.Map,是对 sync.Map 的泛型封装,类型更安全。


五、IM 收发消息在分布式下的难点

5.1 核心问题

IM 后端是分布式系统,部署多个节点。A 连上节点 1,B 连上节点 2,A 给 B 发消息时,节点 1 怎么知道 B 连在哪个节点上?群聊场景下这个问题更复杂。

白话类比:A 在 1 号营业厅办业务,B 在 2 号营业厅。A 说"帮我告诉 B 一声",1 号厅得先知道 B 在哪家厅,才能把话传过去。这就是分布式 IM 的核心难题——连接状态分散在多个节点,消息要跨节点投递

5.2 方案一:注册机制

B 连上节点后,把自己的信息注册到注册中心。节点 1 收到 A 发给 B 的消息后,查询注册中心找到 B 所在节点,转发过去。

A → 节点1 → 注册中心(查 B 在哪) → 节点2 → B

缺点:注册中心容易成为瓶颈。每个用户上线都要注册,多端登录时压力更大。

5.3 方案二:广播机制

节点 1 不知道 B 在哪,索性给所有节点广播。节点 2 收到后发现 B 连着自己,就转发给 B。

广播的实现方式:

  • RPC 广播调用
  • 消息队列(各节点订阅同一 topic)
  • Redis 发布订阅

缺点:消息队列或 Redis 这类中间件容易成为瓶颈。

flowchart TD
    subgraph 注册机制
        A1[A→节点1] --> REG[(注册中心)]
        REG --> N2[节点2→B]
    end
    subgraph 广播机制
        A2[A→节点1] --> B1[广播到所有节点]
        B1 --> N3[节点2 命中 B]
        B1 --> N4[节点3 无 B 丢弃]
    end

这张图在讲:注册机制"精准找人但压注册中心",广播机制"广撒网但压消息中间件"。两者各有瓶颈。

5.4 实践方案:网关 + 后端服务

生产环境 IM 通常分两层:

  • 网关(Gateway):负责维护 WebSocket 连接,稳定少变更,轻量易扩展。
  • 后端服务:处理业务逻辑,微服务架构下有大量服务。

分离的原因:

  • 网关稳定 → 用户连接不会因后端变更而频繁迁移。
  • 网关轻量 → 容易横向扩展支撑更多连接。

六、基于 WebSocket + Kafka 的最简群聊 IM

6.1 整体架构

A(客户端) ──ws──> 网关1 ──> Kafka(topic=msg) ──> 网关1/2/3(各自消费)
                                                  │
                                                  ├── B(连网关1)
                                                  ├── C(连网关2)
                                                  └── D(连网关3)
flowchart LR
    A[客户端 A] -->|ws| GW1[网关1]
    GW1 -->|发到 Kafka| K[(Kafka topic=msg)]
    K --> GW1
    K --> GW2[网关2]
    K --> GW3[网关3]
    GW1 --> B[用户 B]
    GW2 --> C[用户 C]
    GW3 --> D[用户 D]

这张图在讲:所有网关都往同一个 Kafka topic 发消息,也都消费这个 topic。于是任意网关收到的消息,会被"广播"到所有网关,再由各网关下发给连在自己身上的用户。Kafka 在这里充当了"广播总线"。

6.2 Message 格式定义

package im

// Message IM 消息格式(最简版)
type Message struct {
    Seq     int64  `json:"seq"`     // 前端生成的序列号,本次连接内唯一即可
    Type    string `json:"type"`    // 消息类型:text/image/system...
    Content string `json:"content"` // 具体内容,根据 Type 解析
    Cid     int64  `json:"cid"`     // channel id:聊天 ID(单聊/群聊/系统通知)
}

字段说明

  • seq:前端生成,用于前后端关联(如 ACK、重试)。后端有自己的 ID。
  • Type:标识不同消息类型,支持多媒体和系统消息。
  • Cid:聊天 ID。单聊和群聊都是聊天,系统通知也有类似 id(如反诈频道)。

6.3 Gateway 接收消息并转发

package im

import (
    "encoding/json"
    "errors"

    "github.com/gorilla/websocket"
)

// Gateway WebSocket 网关
type Gateway struct {
    conn    *websocket.Conn
    uid     int64
    backend BackendService // 转发到后端服务
}

// ReceiveAndForward 接收客户端消息并转发到后端
func (g *Gateway) ReceiveAndForward() error {
    for {
        // 步骤 1:读消息
        _, data, err := g.conn.ReadMessage()
        if err != nil {
            return err
        }

        // 步骤 2:反序列化为 Message
        var msg Message
        if err := json.Unmarshal(data, &msg); err != nil {
            // 通知前端格式错误
            g.sendError(0, "invalid message format")
            continue
        }

        // 步骤 3:转发到后端服务处理(存储、审核、分发)
        if err := g.backend.SendMessage(g.uid, msg); err != nil {
            // 通知前端发送失败
            g.sendError(msg.Seq, "send failed")
            continue
        }

        // 步骤 4:返回 ACK 给前端
        g.sendAck(msg.Seq)
    }
}

func (g *Gateway) sendAck(seq int64) {
    ack, _ := json.Marshal(map[string]any{"type": "ack", "seq": seq})
    _ = g.conn.WriteMessage(websocket.TextMessage, ack)
}

func (g *Gateway) sendError(seq int64, reason string) {
    e, _ := json.Marshal(map[string]any{"type": "error", "seq": seq, "reason": reason})
    _ = g.conn.WriteMessage(websocket.TextMessage, e)
}

6.4 后端服务:找成员 + 投递 Kafka

package im

import "context"

// BackendService 后端服务接口
type BackendService interface {
    SendMessage(ctx context.Context, sender int64, msg Message) error
}

// KafkaBackend 基于 Kafka 的后端实现
type KafkaBackend struct {
    groupSvc GroupService // 查询群成员
    producer KafkaProducer
}

func (b *KafkaBackend) SendMessage(sender int64, msg Message) error {
    // 步骤 1:找到当前聊天的所有成员
    members, err := b.groupSvc.GetMembers(context.Background(), msg.Cid)
    if err != nil {
        return err
    }

    // 步骤 2:投递到 Kafka,每个网关节点都会消费
    //    生产中这里还会包含:存储消息记录、内容审核、同步搜索等
    return b.producer.Send("msg", encodeMsg(sender, msg, members))
}

6.5 Gateway 订阅 Kafka 并下发

package im

import "context"

// ConsumeLoop 每个网关节点都作为独立消费组,消费全部消息
func (g *Gateway) ConsumeLoop(ctx context.Context) error {
    // 步骤 1:以"节点 ID"为消费组名,保证每个网关都独立消费全量消息
    consumer := g.kafka.Consumer("msg-group-" + g.nodeID) // 每个节点独立消费组
    for {
        select {
        case <-ctx.Done():
            return ctx.Err()
        case msg := <-consumer.Messages():
            var event MsgEvent
            json.Unmarshal(msg.Value, &event)
            // 步骤 2:只把连在自己节点上的目标用户真正下发
            for _, uid := range event.TargetUids {
                if conn, ok := g.hub.Get(uid); ok {
                    _ = conn.WriteMessage(websocket.TextMessage, event.Payload)
                }
            }
        }
    }
}

关键点:每个网关节点必须是独立的消费组,否则一个消息只会被一个节点消费,无法广播到所有节点上的用户。

⚠️ 新手必踩的坑:群聊消费组不能共享。如果所有网关共用一个消费组名,Kafka 会把消息"分摊"给组内成员——一条消息只会被其中一个网关消费,连在别的网关上的用户就收不到。所以每个网关节点必须用不同的消费组(用 nodeID 区分),才能各自拿到全量消息再本地过滤下发。


七、IM 核心难点补充

PDF 中重点讲了架构,这里补充 IM 工程上的几个核心难点,面试常考:

7.1 长连接维护

  • 心跳机制:客户端定时发心跳,服务端检测超时断开,避免半开连接。
  • 重连策略:断线后指数退避重连,避免雪崩。
  • 连接迁移:网关扩缩容时如何平滑迁移连接(优雅关闭 + 客户端重连)。

7.2 消息可靠投递

  • ACK 机制:客户端发消息 → 服务端 ACK → 客户端确认。未收到 ACK 则重试。
  • 去重:基于客户端 seq 或服务端消息 ID 去重,避免重发导致重复。
  • 不丢消息:服务端持久化成功后再 ACK;推送时先持久化再推送。
sequenceDiagram
    participant C as 客户端
    participant S as 服务端
    C->>S: 发消息 seq=1
    S->>S: 持久化落库
    S-->>C: ACK(seq=1)
    Note over C: 未收到 ACK 则按 seq 重发

这张图在讲:可靠投递的标准姿势——服务端"先落库、再 ACK",客户端"没 ACK 就重发",配合 seq 去重,做到不丢不重。

7.3 消息顺序性

  • 单聊:同一会话内消息按发送顺序到达,通常用服务端单调递增 ID。
  • 群聊:群内消息全局有序,需要集中的 ID 分配器或逻辑时钟。

7.4 离线消息

  • 用户离线时消息持久化到 DB。
  • 用户上线后拉取未读消息,拉取后标记已读。
  • 离线消息存储一般有时效(如 7 天/30 天)。
flowchart TD
    M[消息到达] --> P[持久化 DB]
    P --> ON{用户在线?}
    ON -- 是 --> D[实时下发]
    ON -- 否 --> OFF[存离线]
    OFF --> U[用户上线拉取未读]
    U --> R[标记已读]

这张图在讲:离线消息的本质就是"在线就推、不在线先存",上线再补拉。再加已读回执计算未读数。

7.5 已读未读

  • 客户端上报已读位置(lastReadSeq)。
  • 服务端比对消息 ID 与已读位置,计算未读数。
  • 群聊已读未读更复杂,需为每个用户维护独立位置。

八、OpenIM 实战

8.1 OpenIM 是什么

OpenIM 是一款开源 IM 解决方案,提供完整的工具和库来构建实时通讯应用。它极大简化了 IM 应用开发流程。

优势

  • 开源:开放源代码,活跃社区支持。
  • 灵活可定制:灵活架构,支持自定义消息类型。
  • 跨平台:Web/iOS/Android 多端通用。

8.2 OpenIM 整体架构

┌─────────────────────────────────────────┐
│  SDK 层(灰色):多语言/多平台 SDK 接入  │
├─────────────────────────────────────────┤
│  access layer:访问层,少量业务逻辑     │
├─────────────────────────────────────────┤
│  service layer:核心模块                │
│   push / auth / user / msg / friend /   │
│   group ...                             │
├─────────────────────────────────────────┤
│  消息队列 → msgtransfer → 本地缓存      │
├─────────────────────────────────────────┤
│  存储层:数据存储 + 消息存储            │
└─────────────────────────────────────────┘
flowchart TD
    SDK[SDK 层] --> ACC[access layer]
    ACC --> SVC[service layer
push/auth/user/msg/group] SVC --> MQ[消息队列 → msgtransfer → 本地缓存] MQ --> STORE[存储层
数据 + 消息]

这张图在讲:OpenIM 的分层和我们自己设计的一致——接入层管连接、service 层管业务、底下是消息总线与存储。

8.3 接入 OpenIM 的基本架构

OpenIM 部署涉及三部分:

  • OpenIMServer:IM 服务端,承担所有 IM 后端责任。
  • APP Client:应用客户端,通过接入 OpenIMSDK 搭建。
  • APP Server:应用服务端,需要和 OpenIM 交互。

8.4 Docker Compose 部署

# 克隆部署仓库
git clone https://github.com/openimsdk/openim-docker openim-docker
cd openim-docker
make init  # 初始化部署参数

# 配置 OPENIM_IP(用环境变量)
export OPENIM_IP=<你的服务器 IP>

# 启动(建议先清理 Docker 中同名容器,避免冲突)
docker-compose up -d

注意:启动过程缓慢,需要科学上网拉镜像。最好清理 Docker 中之前部署的同名容器,防止端口冲突。

8.5 试用 OpenIM

  1. OPENIM_IP(不要用 localhost,小部分功能会有问题)访问 Web 客户端。
  2. 用任意手机号注册,演示环境验证码统一为 666666
  3. 在一个账号中添加另一个账号为好友,对方同意后才能开始对话。
  4. 管理后台:http://OPENIM_IP:11002,默认账号密码都是 chatAdmin

8.6 OpenIM 监控

OpenIM 自带监控接入:

关键监控指标

  • message/s:发送消息频率
  • failures/s:失败频率

8.7 在业务系统中接入 OpenIM

场景:业务系统需要站内私聊功能(不需要完整 IM)。要点:

  1. 前端:完成 OpenIM SDK 接入,在网站中嵌入私聊窗口。
  2. 后端:业务系统与 OpenIM 数据同步。

用户数据同步流程(推荐异步):

用户注册 → 业务 DB → Canal 监听 binlog → 消费者 → 调 OpenIM 注册接口

Go 代码示例(用 ekit 的 httpx):

package openim

import (
    "context"
    "net/http"

    "github.com/gotomicro/ekit/httpx"
)

// UserSyncer 把业务系统用户同步到 OpenIM
type UserSyncer struct {
    secret    string // 部署 OpenIM 时配置在 config/config.yaml 的 secret
    apiBase   string // 例如 http://OPENIM_IP:10002
    transport http.RoundTripper
}

// SyncUser 注册用户到 OpenIM
func (s *UserSyncer) SyncUser(ctx context.Context, uid int64, nickname, phone string) error {
    // 步骤 1:构造请求体
    reqBody := map[string]any{
        "secret": s.secret,
        "users": []map[string]any{
            {
                "userID":   uid,
                "nickname": nickname,
                "phone":    phone,
            },
        },
    }

    // 步骤 2:构造请求
    req, err := httpx.NewRequest(ctx, http.MethodPost, s.apiBase+"/user/user_register", reqBody)
    if err != nil {
        return err
    }
    // operationID 用来标识本次请求,推荐用 OTel 的 trace ID
    req.Header.Set("operationID", traceIDFromContext(ctx))

    // 步骤 3:发送请求
    var resp struct {
        ErrCode int    `json:"errCode"`
        ErrMsg  string `json:"errMsg"`
    }
    if err := httpx.DoRequest(s.transport, req, &resp); err != nil {
        return err
    }
    // 步骤 4:校验返回:ErrCode 不为 0 说明失败
    if resp.ErrCode != 0 {
        return fmt.Errorf("openim register failed: %s", resp.ErrMsg)
    }
    return nil
}

关键点

  • 请求要携带 secret(部署时配置在 config/config.yaml)。
  • 要传 operationID 头部,正常用 OTel 的 trace ID。
  • 返回的 ErrCode != 0 即为失败。

8.8 前端 SDK

前端 SDK 是接入 OpenIM 最难的部分(“五花八门”),OpenIM 提供多语言 SDK:Web、iOS、Android、Flutter、Unity 等。具体 API 参考 OpenIM 官方文档。


九、工程实践要点

  1. 网关与后端分离:网关稳定承载连接,后端处理业务,扩缩容互不影响。
  2. 每个网关节点独立消费组:广播机制下,节点之间不能共享消费组,否则消息只会被一个节点消费。
  3. 慢消费者隔离:广播写入要用异步+超时,避免一个慢客户端拖垮整个网关。
  4. 连接数监控:每个节点接入的 WebSocket 数量是核心指标,数量过多会显著影响性能。
  5. WebSocket 调优:调整 read/write buffer、超时设置、Linux TCP 参数(如 tcp_tw_reusesomaxconn)。
  6. 异步优先:业务系统接入 OpenIM 优先走 Canal 监听 binlog,业务服务零侵入。
  7. 不要用 localhost 访问 OpenIM:小部分功能会有问题,统一用 OPENIM_IP

十、自测题与动手练习

自测题(合上书能答出来,才算懂)

  1. 普通 Web 服务用"请求/响应"模型,IM 为什么必须用长连接?轮询模型在 IM 场景下哪两个缺点会被放大?
  2. WebSocket 握手时客户端发了哪个关键请求头?服务端回什么状态码表示"协议切换成功"?
  3. 分布式 IM 下,A 连节点 1、B 连节点 2,A 给 B 发消息有哪两种转发思路?各自的瓶颈在哪?
  4. 基于 Kafka 的群聊里,为什么每个网关节点必须用"不同的消费组"?如果共用一个消费组会怎样?
  5. IM 的"可靠投递"靠什么机制保证不丢不重?离线消息和已读未读分别是怎么实现的?

动手练习(建议真做一遍)

  1. 起一个 WebSocket 服务:用 gorilla/websocket 起服务,让 HandleWS 每秒推时间戳,用 wscat 或浏览器连上去验证能收到推送。
  2. 写单节点 Hub 广播:实现 Hub 的 Register/Unregister/Broadcast,写个测试:A、B 连上后 A 发消息,确认 B 收到而 A 自己没收到。
  3. 跑通 Kafka 群聊广播:起 2 个网关节点消费同一个 topic,故意让它们共用消费组,观察"有用户收不到消息"的现象;改成独立消费组后现象消失。

十一、本章小结

  • IM 的复杂度远超普通 Web 服务,核心难点在于长连接维护分布式下的消息投递
  • 三种协议各有所长:XMPP 适合互通、WebSocket 是主流、WebRTC 适合音视频;本课程主角是 WebSocket(HTTP 升级 + 全双工)。
  • 分布式 IM 跨节点转发有注册机制(精准但压注册中心)和广播机制(广撒网但压中间件)两种思路,生产环境用网关 + 后端分离 + Kafka 广播
  • 群聊广播的关键细节:每个网关节点必须是独立消费组,否则消息只被一个节点消费,其他节点上的用户收不到。
  • OpenIM 是成熟的开源 IM 方案,Docker Compose 一键部署,通过 SDK + 业务后端(推荐 Canal 监听 binlog 异步同步)可接入任何业务系统。
  • IM 的工程难点(可靠投递、顺序性、离线消息、已读未读)是面试深度题,要结合 ACK、单调 ID、先持久化后推送、已读位置等机制讲清楚。

下一章(第21章,也是本课程最后一章)我们做课程总结与进阶路线——把 21章的知识体系串成一张能力地图,复习设计原则、设计模式、缓存与微服务治理,并给出成长路线与面试速查。

About Me

没什么想介绍的,一个很大众的码农…

喜欢代码,车,马,真的是 🐎

讨厌别人让我给自己的代码写注释 最厌烦别人的程序没有写注释

目标

学AI,加油!加油!