多彩编程 多彩编程MZPH · CODE BLOG
ARTICLE DETAIL

文章详情

深耕前端与后端开发技术的一线实战笔记与踩坑复盘。

etcd二次封装实战:解决IM分布式服务发现与配置管理难题

etcd二次封装实战:解决IM分布式服务发现与配置管理难题 做即时通讯项目做到第三期我的节奏反而慢下来了。前面两期把长连接网关、消息投递链路跑通以后服务数量一上来所有依赖集群协作的问题几乎同时冒出来新网关节点怎么被其他服务发现用户 Session 路由表存哪里灰度配置怎么样不下发到每台机器早期我是直接用 etcd 的 clientv3 裸调的代码倒是能跑可等到节点超过几十台问题就开始连续爆炸。所以这期我把重点放在 etcd 二次封装上不是讲 etcd 基础用法而是讲我在即时通讯场景里把 clientv3 包上一层业务接口之后踩过的坑、想清楚的设计、以及最后落地的样子。如果你正在做 IM 或者任何分布式服务这篇文章应该能帮你少走不少弯路。1. 为什么 IM 项目里要上 etcd节点注册、路由发现和在线状态协调1.1 IM 分布式架构里哪些数据必须交给 etcd先把场景铺开。一个稍微正经的 IM 系统不会只有一个网关、一台消息服务。至少会有多台长连接网关客户端随机连到其中一台多台消息推送服务需要找到用户当前挂在哪个网关上多台业务 API 服务可能需要按集群维度做限流或放量若干后台任务比如离线消息清理、会话归档这些任务往往只能有一个执行者。这里的数据可以分成三类。第一类是业务数据比如消息内容、会话列表、用户资料这些必须交给 MySQL、Redis 或者对象存储etcd 不适合存高频写入的海量数据。第二类是集群元数据包括服务节点列表、用户到网关的路由映射、在线状态、配置开关。这类数据的特点是单条很小、读多写少、变化需要被集群里其他节点感知。第三类是临时协调状态比如分布式锁、选主结果、幂等标记它们需要原子性和租约保活。etcd 主要吃后两类。举个例子用户长连接挂在网关 A网络切换后重连到了网关 B。网关 B 在接收消息之前必须知道用户之前的连接已经被摘除否则会出现双端同时收发消息的脏状态。这个处理可以靠业务层单独查数据库但更优雅的做法是维护一个路由 key/im/route/{userID} - gateway-a节点注册时写入连接断开时通过租约自动过期。这样任何服务要看用户在哪个网关直接查 etcd 就行。etcd 在这个架构里其实承担了“总机话务员”的角色。1.2 和 ZooKeeper、Redis 相比我为什么在 IM 场景选 etcd做技术选型时肯定绕不开 ZooKeeper、Redis甚至 Consul。我简单列过一张对比表能力etcdZooKeeperRedis一致性模型Raft 线性一致ZAB 线性一致异步复制为主主从切换有窗口监听机制Watch 带 revision可断点续传Watch 一次性触发需重新注册Pub/Sub 不保证送达Stream 偏消息队列临时节点/租约Lease 自动续约临时节点 Session 超时可借助 TTL key 模拟但语义弱KV 历史版本MVCC可回溯有版本号但能力弱基本不具备事务能力多 key 条件事务有但使用繁琐单 key 原子 / Lua 脚本运维成本单二进制etcdctl 好使配置多客户端协议老高可用拓扑要额外维护客户端模型gRPC多语言官方/社区 SDK独有协议客户端维护成本高生态广但分布式语义弱Redis 做缓存和在线状态很好但它的 pub/sub 本质不保证消息不丢主从切换时分布式锁也存在争议拿它当注册中心有点勉强。ZooKeeper 的临时节点和 Watch 模型很适合服务发现只是运维成本和历史包袱偏重新项目里我一般不太愿意再引入一套 ZK 集群。etcd 胜在“一个系统把多个事干完”线性一致的 KV、带 revision 的 Watch、租约、事务、选主。对 IM 这种既要路由发现、又要配置下发、偶尔还要抢锁的场景它比组件拼盘省心。当然它也不是万能后面我会专门讲哪些事不该往 etcd 里塞。2. 裸用 clientv3 的五个痛点二次封装不是炫技是工程止损2.1 连接对象散落各处没人管生命周期最早期我见过一个服务里十多个文件各自构造 etcd client。看起来没什么因为本地测试都能跑但生产环境一旦要对 etcd 做证书轮换、endpoints 调整、或者统一加超时就需要全局搜索所有clientv3.New的调用点去改。更要命的是如果某个模块构造了 client 却忘记 Closegrpc 连接泄漏会比较隐蔽。后来我把连接收敛成统一的单例整个进程只有一个 etcd client连接由封装层负责建立和关闭业务代码永远拿不到裸的 *clientv3.Client。2.2 Watch 恢复逻辑永远在复制粘贴这是裸用 clientv3 最大的坑。你要实现监听一个配置变更正确的流程是先 Get 一次拿到当前值以及当前 revision再用 Watch 从 revision1 开始监听后续变化。如果 Watch 中途断掉还要从上次收到的 revision 继续。如果 etcd 已经执行了 compaction历史 revision 被清除Watch 会直接报ErrCompacted这时候必须重新降级成全量拉取。这套逻辑几乎没有新手第一次就能写对。问题在于每个业务方都自己实现一份结果就是有人从 revision 0 开始 Watch把所有历史事件重新放一遍有人断线后从头全量拉取导致短暂重复还有人没处理 compaction一报错就重启服务。封装层必须把“全量快照 增量Watch 断线恢复 compaction 降级”收敛成一套机制对外只提供一个回调循环。2.3 租约续约、服务下线、异常清理全凭自觉原生Lease.KeepAlive是异步的返回一个 channel你需要在 goroutine 里持续消费。服务正常退出时还需要手动Revoke租约否则节点会在 etcd 里变成僵尸直到租约 TTL 超时。如果你注册的 TTL 设得比较长比如 60 秒故障节点被集群感知的时间就会很慢如果设得短网络抖动又容易导致误摘除。封装层要做的一件事是把“注册”变成一个带生命周期的对象创建租约、自动续约、服务退出自动 Revoke、进程异常时依赖 TTL 兜底。业务代码只负责调用一次Register不需要关心租约的细节。2.4 错误处理粒度错位业务代码被 etcd 细节污染clientv3 返回的错误类型很底层比如context deadline exceeded、rpc error: code Unavailable、ErrCompacted、ErrLeaseNotFound。业务代码其实只关心三个结果成功、可重试失败、不可重试失败。裸用的话你会看到业务逻辑里到处打日志、判断错误码、自己写重试。这会让代码被分布式系统的复杂性填满真正的业务语义反而没人看得清。封装层应该负责翻译这些底层错误连接失败自动重试超时按策略处理compaction 自动降级最终只给上层返回明确的成功或失败。2.5 测试与灰度切换成本高如果业务代码直接依赖 clientv3写单元测试就必须起一个真实的 etcd 容器或者把整个 client 接口 mock 掉。起容器倒不复杂但测试速度慢mock 整个 client 又特别笨重因为clientv3.Client是大接口字段还很多。封装成面向业务的 interface 之后测试里用一个内存版实现就行这就是二次封装的隐形收益。后面如果要把服务发现换掉只要实现同一套接口业务代码一行都不用动。3. 封装设计的第一刀先定义面向 IM 业务的接口再谈实现3.1 Registry服务注册与发现的最小接口我当时先做的是注册发现接口。命名直接从业务语义出发不要出现Put、Get、Watch这种存储词汇。接口原型大概是type Registry interface { // Register 注册一个服务实例重复注册同一个实例 ID 需幂等 Register(ctx context.Context, instance Instance) error // Deregister 注销一个服务实例 Deregister(ctx context.Context, instance Instance) error // Watch 监听某个服务下的实例列表变化 Watch(ctx context.Context, service string) (-chan []Instance, error) } type Instance struct { ID string // 全局唯一实例 ID Service string // 服务名比如 im-gateway Addr string // 对外可访问地址 Meta map[string]string // 附加信息如区域、权重 }这个接口的好处是调用方完全不需要知道 key 怎么拼、租约怎么续、Watch 断线怎么办。Instance就是 IM 集群里一个节点的自然描述。实现层面内部把Service映射成/im/service/{service}/nodes/{id}把Addr和Meta序列化成 JSON 存入 value。3.2 配置中心接口KV 存储和 Watch 要分开看配置下发是 IM 刚需。黑名单、灰度开关、推送频率、消息超时时间这些都要能随时调整并推送到所有节点。我定义了一个非常收敛的接口type ConfigStore interface { Get(ctx context.Context, key string, v interface{}) error Put(ctx context.Context, key string, v interface{}) error Watch(ctx context.Context, key string, onChange func(version int64, data []byte)) }关键点是version而不是revision。业务不需要知道 etcd 的 revision 语义只需要知道配置是第几版能用来判断是否过期。数据序列化默认用 JSONkey 固定加/im/config/前缀。Watch 回调必须保证顺序调用上层不要自己在回调里开 goroutine否则配置生效顺序可能错乱。3.3 锁与选主接口怎么收口IM 里真正需要分布式锁的场景不算多但一旦需要就必须要稳定。封装层提供两个小接口type Locker interface { // Lock 尝试加锁返回解锁函数如果 ctx 超时则返回错误 Lock(ctx context.Context, key string) (UnlockFunc, error) } type LeaderElector interface { // Campaign 参与竞选成为 leader 后返回会话对象 Campaign(ctx context.Context, name string) (LeaderSession, error) } type LeaderSession interface { Resign(ctx context.Context) error Done() -chan struct{} }Lock内部封装了 etcd 的事务和租约调用方拿到UnlockFunc后必须defer调用。Campaign内部可以基于 etcd 官方 election 包做封装外面只暴露一个会话。锁的 key 也要归入命名空间避免和业务数据冲突。3.4 接口设计的原则不让业务感知 etcd 的存在二次封装最容易犯的错误是把封装层做成“薄薄一层”还是让调用方传各种 etcd 选项。我的原则是业务代码只关心自己的领域对象不关心存储在哪里、如何存储。接口定义阶段就要问如果明天我把 etcd 换成 Consul这个接口是否还成立如果答案是否定的说明接口里混入了底层概念。比如Registry.Watch返回-chan []Instance这就是领域语义如果返回的是 etcd 的WatchChan那就是没封干净。再比如错误处理封装层内部统一把错误分类对外只返回ErrNotFound、ErrTimeout这种业务错误上层不会被迫处理rpc error。4. 连接层与 Watch 机制详解从建模到防翻车4.1 client 初始化endpoints、DialTimeout、AutoSync 和单例连接层是整个封装的地基。初始化配置我一般这样定义type Config struct { Endpoints []string // 至少两个 etcd 节点 DialTimeout time.Duration // 建议 5s太短容易误判故障 Username string Password string TLS *tls.Config AutoSync bool // 跟随 etcd 集群 member 变化 }有几个细节非常重要。第一endpoints 不要只填一个节点etcd 集群至少要三个节点客户端会自己做负载均衡和故障转移。第二DialTimeout别设成 1 秒etcd 集群在做 leader 选举时短时间内的连接建立会变慢太短的超时会让客户端频繁报错。第三如果集群开了鉴权用户名密码和证书配置必须在封装层统一设置不要让每个业务方自己传。生产环境建议做进程级单例。多个业务模块共享一个 etcd client既省连接又能统一做监控上报。我在封装里加了两个指标请求耗时、错误率。有了这两个指标后面排查问题轻松很多。4.2 Watch 的版本号语义从哪开始读才能真正不丢事件很多同学没搞懂 etcd 的 revision 机制导致 Watch 漏事件或者重复消费。etcd v3 里每个修改都会分配一个单调递增的 revision。注意不是每个 key 一个版本号而是全局一个。所以要做“先全量再增量”的监听标准流程是先Get(key, WithPrefix())拿到当前所有 key-value以及这次 Get 返回的 revision从 revision1 开始Watch(key, WithPrefix())这样不会漏掉 Get 之后、Watch 建立之前发生的变更Watch 断线之后用本地记录的最后事件 revision1 重新发起 Watch。用一段伪代码表示func WatchWithSnapshot(ctx context.Context, key string) (-chan Snapshot, error) { resp, err : cli.Get(ctx, key, clientv3.WithPrefix()) if err ! nil { return nil, err } // 先返回全量快照 snap : Snapshot{ Data: resp.Kvs, Revision: resp.Header.Revision, } // 从快照 revision 的下一个开始监听增量 wch : cli.Watch(ctx, key, clientv3.WithPrefix(), clientv3.WithRev(snap.Revision1), ) // 把快照和 watch 事件合并后统一返回... }这段逻辑值得放在封装层里而不是让业务方自己写。因为边界情况很多如果业务拿到快照后还没建立 Watch这个间隙里的变更会不会丢只要从revision1开始 Watch就不会丢。关键是必须用 Get 返回的 revision不是用本地时间或者其他 key 的版本号。4.3 应对 compact、修改冲突和事件乱序etcd 默认会自动 compaction保留最近一段时间的历史版本。如果 Watch 落后太多从某个很老的 revision 开始监听就会收到ErrCompacted。封装层的正确降级策略是重新走一次全量 Get 增量 Watch。也就是说业务侧可能看到一次“重新全量”这是正常现象实现时要保证这次全量不会导致业务状态错乱。另外要小心事件处理顺序。同一个前缀下多个 key 的变更客户端会收到持续的事件流。如果业务回调里异步处理就可能出现 A 事件先到、B 事件后到但因为异步并发导致处理结果反过来。我的经验是监听同一个数据视图的事件必须单 goroutine 按顺序处理如果想提高吞吐可以按 key 的 hash 拆分多个 worker但要保证相同 key 的事件进同一个 worker。4.4 key 目录设计与命名空间直接影响运维效率etcd 里没有真正的目录概念key 是一个扁平字符串但路径式命名是最好的实践。我用这套风格/im/service/{serviceName}/nodes/{instanceID} /im/route/{userID} /im/config/{configKey} /im/election/{leaderName} /im/lock/{resourceKey}好处有三个。第一etcdctl 可以按前缀检索排障时一条命令就能看到某个服务的全量节点。第二权限控制可以按前缀做比如只允许网关服务写/im/service/gateway/前缀。第三Watch 时不会串数据WithPrefix()的边界清晰。还有一个特别容易踩的坑不要在单个 key 里塞大 value。etcd 适合存小数据单个 value 超过几百上千字节就要警惕。我见过有人把某个用户的全量连接信息塞进路由 key一个 key 几十 KBWatch 全量拉取时直接把内存打满。正确做法是 key 里只放最小必要字段比如gatewayID connID详细信息放 Redis 或其他存储。5. 租约、事务与选主IM 里的分布式锁和单点控制实战5.1 Lease KeepAlive 的正确姿势etcd 的租约Lease是服务发现的基础。节点注册后绑定一个租约只要租约不过期节点就认为存活进程挂掉后租约到期自动删除这就实现了“临时节点”。租约的使用有几个关键点。第一续约要持续进行。第二续约 channel 关闭不代表节点立刻失效可能是网络抖动此时要根据业务场景决定是重新续约还是主动下线。第三服务正常退出时要显式调用Revoke不要等 TTL 超时否则下线不够及时。我封装时 TTL 一般设 10~30 秒。太短网络抖动容易造成节点大规模误摘除太长故障节点被发现得慢。IM 场景里网关节点挂了我们希望 30 秒内集群做出反应所以 TTL 15 秒左右比较合理。续约的伪代码leaseResp, err : cli.Grant(ctx, int64(ttl.Seconds())) if err ! nil { return err } keepAliveCh, err : cli.KeepAlive(ctx, leaseResp.ID) if err ! nil { return err } go func() { for { select { case _, ok : -keepAliveCh: if !ok { // 续约异常需要根据业务决定是否重新注册 return } case -ctx.Done(): // 注意退出时不要复用已取消的 ctx cli.Revoke(context.Background(), leaseResp.ID) return } } }()这里有个真实教训Revoke必须在进程退出时执行但很多服务接到退出信号后context.Background()才是安全的因为原来传入的 ctx 可能已经被取消用它发请求会立刻失败。5.2 用 etcd 事务实现服务注册的原子性服务注册要解决一个问题同一个节点 ID 可能在节点重启后再次注册这时候不应该创建重复数据。用 etcd 的 Txn 可以保证“如果不存在才创建”这个操作的原子性nodeKey : /im/service/gateway/nodes/ instanceID txn : cli.Txn(ctx) resp, err : txn. If(clientv3.Compare(clientv3.CreateRevision(nodeKey), , 0)). Then(clientv3.OpPut(nodeKey, value, clientv3.WithLease(leaseID))). Else(clientv3.OpGet(nodeKey)). Commit()CreateRevision为 0 表示 key 从未创建过。如果节点重启时旧 key 因为租约还没过期这个事务会走 Else 分支拿到旧数据后你可以选择复用旧租约或者等待其过期后重新注册。这个操作比“先 Get 再 Put”安全得多因为两步操作之间可能有并发竞争。5.3 选主IM 里什么场景真的需要锁先泼一盆冷水IM 里很多并发问题不需要分布式锁。比如用户双端在线不需要用锁用版本号或者CAS就能解决消息 ID 全局唯一不需要锁用雪花算法或者分片 ID 生成器就够了。真正需要锁或选主的场景是“跨节点的互斥资源访问”。我项目里用到了三个场景离线消息归档任务每天凌晨只允许一个节点跑避免重复归档长连接网关组内选主主节点负责维护全局连接统计和路由表快照用户 Session 迁移同一个用户从一个网关迁移到另一个网关迁移器只能有一个。选主用 etcd 自带 election 包封装即可。election : etcd.NewElection(client, /im/election/gateway-master) if err : election.Campaign(ctx, nodeID); err ! nil { log.Error(campaign failed, err, err) return } defer election.Resign(ctx)Campaign内部会创建租约并写入 leader key其它节点持续监听这个 key。成为 leader 的节点要定期检查Done()channel如果租约断了会收到通知此时必须释放资源或重新竞选。注意选主是“自动故障转移”不是“负载均衡”别想着拿选主做流量分发。6. 落地效果与排障记录我在真实项目里踩过的坑6.1 事件丢失假象原来是 Watch 被 compaction 打断上线后第一个诡异问题某个服务的路由表偶尔会少几条数据但日志里没有任何错误。排查链路是这样的先看 etcd 里的数据发现 key 还在说明不是写入问题再查消费者进程发现它本地缓存的路由表不完整翻日志发现偶尔有ErrCompacted出现在 watch goroutine 里但业务代码只打了watch closed没有继续处理定位到根因Watch 从很老的 revision 开始监听但 etcd 已经执行了 compaction历史数据被清理Watch 被强制中断。修复方案就是封装层加的降级策略捕获ErrCompacted后重新执行一次全量 Get 增量 Watch。这个坑也说明监听逻辑不能只处理“正常断线重连”还要处理“历史数据被清理”这种更隐蔽的情况。6.2 大 key 拖慢 Watch 全量拉取key 瘦身第二次事故是服务启动时内存飙高。当时某个服务注册时把每个用户的连接信息都写到了同前缀的 key 里一个 key 甚至几十 KB。启动时WithPrefix()全量拉取几十万 key 一次性返回直接把内存打满Watch 通道也积压了大量事件。这个问题的本质是数据建模错了。etcd 适合存“路由表”不适合存“详细会话快照”。我当时的整改是把 key 里的扩展信息拆出去etcd 只保存gatewayID|connID用户详情放 Redis等真正需要时再查。改完之后全量拉的体积下降了 90% 以上。不要觉得“etcd 能存就存”等 Watch 积压到几万条事件处理延迟会指数上升。6.3 grpc 连接数飙升与客户端资源泄漏第三次问题是排查 etcd 连接数。某次压测后我在 etcd 节点上看到来自同一个服务的 grpc 连接数异常高远超预期。原因有两层业务代码里有人通过封装层的GetClient()方法拿到了裸 client然后不小心又调用了New构造了新连接某些长连接任务没有正确关闭导致 goroutine 和连接持续堆积。最终我把GetClient()严格收口业务代码只能拿封装接口裸 client 只在封装层内部保留另外所有KeepAlivegoroutine 必须在退出时关闭。加了这些控制后连接数稳定在个位数。grpc 客户端还应该设置 keepalive 参数避免空闲连接被防火墙切断后没有及时重建。etcd 官方文档也推荐设置这些参数但在裸用的时候很容易被忽略封装层正好统一处理。6.4 二次封装的边界别把 etcd 用成消息队列最后想聊聊封装的边界。二次封装不是把 etcd 所有能力都包一遍才算好。很多团队会把 etcd 当 Redis pub/sub 用搞大量 watch 来做事件广播这其实是错误方向。etcd 的 Watch 适合低频、小数据量的状态同步不适合高频消息投递。IM 里的消息推送应该走自己的长连接通道etcd 只负责集群状态同步。另外要克制住“给封装层加功能”的冲动。我看到有人封了一层EventBus又封了一层TaskQueue最后封装类比 etcd 本身还复杂出了问题根本没人敢改。二次封装的目的是让业务简单不是把简单问题复杂化。如果再让我从头做一次这个封装我会先把注册、配置这两条最核心的链路做好再按需补锁和选主同时每个方法都要能打印耗时、错误、重试次数。这看起来是“多写一层”其实是在给线上问题留退路。如果你也在 IM 项目里集成 etcd希望这些思路能帮你绕开我踩过的坑。
返回列表