资讯详情

资讯详情

从零搭建IM后台服务端:Netty长连接、自定义协议与心跳保活实战

做即时通讯系统后台服务端开发这事我最早是有些轻敌的。做过几年业务系统增删改查写得飞起心里默认它不就是维持个长连接嘛客户端发消息服务端转一下顺手存个库。真上手之后才发现光“连接”这一个词就能拆出接得住、管得牢、断得干净三个坎每一个都够你排查到凌晨。这篇是这个系列的第一篇不急着上分布式那套高深架构先把一个能稳定跑起来、能登录、能点对点发消息的后台服务端从零到一搭出来。适合刚准备做IM、或者被业务推着接手IM模块的兄弟看完你能理解IM后台的核心构成、协议设计、Netty骨架搭建以及前期最常踩的坑。1. 先盘清楚一套IM后台服务端到底要管哪几件事1.1 连接、路由、存储三条主线别被细节带偏IM后台从宏观上可以切成三块接入层、逻辑层、存储层。接入层干的事情很纯粹把客户端的TCP连接稳定地接住维护好连接状态处理网络报文在网络层的边界。逻辑层负责报文内容的解读、路由、鉴权比如一条消息从A用户发出要判断B用户在哪台机器、在不在线、需不需要回执。存储层负责把关键数据沉淀下来消息记录、离线消息、会话列表这些都得落库或落到KV里。很多新手写IM后台会一头扎进“消息转发代码”里其实最应该先理解这三个层次的边界因为后续所有问题——连接风暴、消息乱序、离线堆积——本质上都是某一层的职责没划清楚。如果把三层理解为一条流水线接入层是工厂门口的门卫只认“这个连接是不是合法”不关心消息内容逻辑层是分拣中心负责把快件送到对应的人手中存储层是档案室所有快件都有副本备案。三个层之间有明确的接口后面想加分布式、加多机房也只是把某一层替换成更牛的实现接口不动上层不感知。这就是为什么我建议第一天就按分层思想建目录而不是等业务复杂之后再重构。1.2 为什么长连接是IM的必选项HTTP轮询怎么就不行可能有人会想客户端定时用HTTP请求拉一下新消息不也能做IM吗早期网页版确实这么干过。但轮询有两个硬伤一是延迟不可控轮询间隔设短了服务器压力大设长了消息不及时体验卡顿二是大量请求是无用开销哪怕没有新消息每次请求也要完整走一遍HTTP头、握手、业务鉴权。IM的流量模型是“闲时静默、忙时突发”群聊、直播弹幕这种场景下如果靠轮询服务器会在无消息时被轮询请求打穿有消息时又因为轮询间隔而延迟。长连接TCP或WebSocket是在连接建立后一直保持服务端可以主动往客户端推数据客户端也可以随时发数据不需要每次重新建连。这样消息的实时性由服务端主动推送来保证空闲时的开销也只是一两个心跳包。代价就是服务端必须维护海量长连接管理连接生命周期这部分的复杂度就是IM后台区别于普通Web应用的核心所在。所以整篇文章都会围绕长连接的“生命周期”展开建立、鉴权、收发、保活、断开、清理。如果做网页端后面可以再加一个WebSocket网关专门适配浏览器但服务端核心逻辑不用变只是换一层接入协议而已。1.3 项目目录结构从第一天就按这个来一个IM服务端的代码目录我推荐按照模块来分而不是按数据表来分。下面这个结构是单人项目也能直接用的im-server ├── im-common │ ├── protocol # 协议定义、报文解析 │ ├── model # 消息实体、用户实体 │ └── util # 序列化、压缩工具 ├── im-connector # 接入层 │ ├── server # Netty启动、ChannelInitializer │ ├── handler # 编解码、心跳、分发Handler │ └── session # 连接会话管理 ├── im-logic # 逻辑层 │ ├── auth # 登录鉴权 │ ├── message # 消息路由与转发 │ └── push # 写回客户端 └── im-store # 存储层 ├── dao # 数据库访问 ├── cache # Redis缓存 └── queue # 异步消息队列初学阶段不一定要拆成多模块工程同一个工程下的包结构也可以这样分。但每层职责要遵守接入层不允许直接写数据库逻辑层不允许直接操作Channel所有跨层调用走接口。这句话听起来像教条但IM项目最容易腐化的地方就是这里——图省事在Handler里直接怼了一个Redis查询结果后续想换存储、想加限流改得你怀疑人生。2. 技术选型和协议设计Netty 自定义TCP协议2.1 放弃了MINA和裸NIO选择Netty的理由Java做IM服务端可选的老牌方案有Apache MINA和Netty也有人直接撸Java NIO。裸NIO最大的问题不是不能用而是你得自己处理半包、粘包、ByteBuffer翻转flip、Selector的空转、多线程并发注册这些底层细节。写一个能跑通demo的NIO程序不难但让它在几千连接下稳定运行没有专门技术积累会非常痛苦。我见过有人在Selector.select()里不做空转判断结果忙等占满一个CPU核这种问题不压测根本发现不了。Netty对NIO做了非常合理的封装字节缓冲用ByteBuf自带动态扩容和引用计数不需要flip线程模型用EventLoop的概念一个连接的读写始终由同一个线程执行天然避免并发问题pipeline的Handler机制让编解码、心跳、业务分发可以解耦。之所以没选MINA主要是生态和维护节奏Netty在RPC、分布式框架里的使用量更大遇到问题容易搜到方案社区活跃度也高。而且Netty对连接异常、半关闭、优雅停机的处理都考虑得比较周全这些在生产环境里全是硬指标。2.2 自定义报文协议的详细字段设计服务端和客户端之间要约定一套二进制协议。协议设计的基本原则是能解析、能校验、能扩展。我的这套协议格式如下自定义二进制协议14字节头 变长消息体字段长度说明magic4字节魔数0xC0DE01标识IM协议version1字节协议版本号当前为1command1字节消息指令如1登录2心跳3单聊seq4字节客户端生成的序列号用于去重和ACKlength4字节消息体长度不包含14字节头最大1MB消息体具体是什么由command决定。比如登录指令的消息体是token单聊指令的消息体是{fromId, toId, msgType, content}心跳指令可以不带消息体。这里有两个容易被忽略的细节魔数一定要校验。有个真实案例有人把监控脚本也接到了IM端口上服务端把玄学数据当成消息体去解析直接抛异常把连接关了还带着一堆堆栈日志。加了魔数校验非本协议流量在解码阶段就被丢弃。其次是序列号seq必须存在后面做ACK、去重、消息记录补拉全靠这个字段串起来否则客户端重发消息和服务端历史记录之间就对不上账。2.3 序列化方案取舍JSON还是Protobuf消息体的序列化我见过两种主流做法一种是直接JSON序列化简单直观、可读性强调试的时候能打印出来人眼看另一种是Protobuf编译期生成类体积小、解析快适合对流量敏感的高并发场景。第一阶段做功能验证我建议先用JSON因为日志、抓包、排查问题都很方便。等后续消息量上来、确确实实成为瓶颈再切Protobuf消息体积小个几百字节对大部分IM系统初期来说感知并不明显但错误日志你能一眼看懂的价值是实打实的。这里要提醒一句如果你用了Java生态的Jackson或Gson序列化对象时一定要把类结构定义得扁平一点避免一个消息体里嵌套好几层Map。IM消息体后面要进数据库、要进MQ、要做离线存储结构越简单后续做数据迁移和协议兼容越省心。3. 把Netty骨架跑起来从启动类到第一个客户端上线3.1 启动类和EventLoopGroup的线程分配接下来就是动手环节。我用的是Java 11 Netty 4.1.x先看一眼pom依赖dependency groupIdio.netty/groupId artifactIdnetty-all/artifactId version4.1.100.Final/version /dependency dependency groupIdcom.fasterxml.jackson.core/groupId artifactIdjackson-databind/artifactId version2.15.2/version /dependency dependency groupIdorg.slf4j/groupId artifactIdslf4j-api/artifactId version2.0.7/version /dependency启动类代码public class ImServer { private final int port; public ImServer(int port) { this.port port; } public void start() throws InterruptedException { EventLoopGroup bossGroup new NioEventLoopGroup(1); EventLoopGroup workerGroup new NioEventLoopGroup(); try { ServerBootstrap b new ServerBootstrap(); b.group(bossGroup, workerGroup) .channel(NioServerSocketChannel.class) .option(ChannelOption.SO_BACKLOG, 1024) .childOption(ChannelOption.TCP_NODELAY, true) .childOption(ChannelOption.SO_KEEPALIVE, true) .childHandler(new ChannelInitializerSocketChannel() { Override protected void initChannel(SocketChannel ch) { ch.pipeline() .addLast(new LengthFieldBasedFrameDecoder(1024 * 1024, 10, 4, 0, 0)) .addLast(new MessageDecoder()) .addLast(new MessageEncoder()) .addLast(new IdleStateHandler(120, 0, 0)) .addLast(new HeartbeatHandler()) .addLast(new MessageDispatcher()); } }); ChannelFuture f b.bind(port).sync(); System.out.println(IM server started on port: port); f.channel().closeFuture().sync(); } finally { bossGroup.shutdownGracefully(); workerGroup.shutdownGracefully(); } } }几个参数的含义要说透。bossGroup只开一个线程因为accept连接本身是轻操作一个线程完全够用开多了反而有锁竞争。workerGroup默认线程数是CPU核心数乘2一般不用手动调大。TCP_NODELAY这个选项很容易被忽略它负责禁用Nagle算法。IM场景里大量小包交互如果Nagle算法开着小消息会在缓冲区里攒着等确认造成几十毫秒的额外延迟客户端感知极其明显。SO_BACKLOG是内核accept队列长度并发握手上来的时候超过队列长度的连接会被直接拒绝客户端会看到Connection reset设成1024是个相对稳妥的值。3.2 Pipeline里的编解码器顺序错了就串包Pipeline是Netty里最核心的概念。入站事件从前往后经过Handler出站事件从后往前经过Handler。所以编解码器的顺序不能乱放LengthFieldBasedFrameDecoder和MessageDecoder必须放在最前面业务Handler放后面。如果反了业务Handler收到的是还没解码的原始ByteBuf整个逻辑就全乱了。LengthFieldBasedFrameDecoder看起来参数复杂本质作用只有一个从TCP字节流里精确地切出“一帧完整数据”。它先读第10到13字节的长度字段知道这帧总长是多少再攒够14length字节才往下一个Handler传。这样粘包被拆成多个帧半包被缓存到齐再发下游Handler永远只收到完整消息。MessageDecoder里要做的事包括校验魔数、校验版本号、校验长度合法性再根据command把body解析成对应的消息对象。这里必须说一个项目里最容易犯的错误在Decoder里读数据时一定不要把readerIndex搞乱。如果发现长度字段超过配置的最大值不要继续读直接关连接因为这多半是对端在恶意构造报文想用超大长度声明撑爆内存。3.3 会话管理用户ID和Channel的绑定与解绑服务端要能根据用户ID找到他的连接也要能根据连接找到用户ID。前者是消息路由的需要后者是连接断开后清理会话的需要。我用一个双向Map管理public class SessionManager { private static final ConcurrentHashMapLong, Channel USER_CHANNEL new ConcurrentHashMap(); private static final ConcurrentHashMapChannel, Long CHANNEL_USER new ConcurrentHashMap(); public static void bind(Long userId, Channel channel) { USER_CHANNEL.put(userId, channel); CHANNEL_USER.put(channel, userId); } public static Channel getChannel(Long userId) { return USER_CHANNEL.get(userId); } public static Long getUserId(Channel channel) { return CHANNEL_USER.get(channel); } public static void unbind(Channel channel) { Long userId CHANNEL_USER.remove(channel); if (userId ! null) { USER_CHANNEL.remove(userId); } } }为什么要维护两个方向的映射因为通道关闭时我们只拿到Channel对象没有UserId如果不做反向Map就没法知道该清理哪条用户记录时间一长就会出现大量僵尸会话。很多初版IM的死连接问题根子就在这里。绑定的时候还要处理“顶号”场景同一个用户ID在新设备上登录正常策略是“新挤旧”先把旧Channel发一个踢下线通知再关闭旧连接。如果不做这一步同一个用户会有多个连接同时在线消息就会重复推送客户端出现消息错乱。顶号逻辑可以放在登录认证成功的分支里调用SessionManager之前先查一次现有通道。4. 打通一个完整的消息闭环登录认证 单聊转发4.1 登录请求从解码到会话注册的完整链路长连接建立后第一个包必须是登录请求。为什么要单独做登录指令而不是建立连接时就把token塞在参数里因为TCP连接本身不具备身份信息必须靠应用层握手来确认“你是谁”而且token有可能过期连接保持期间需要重新登录续期独立登录指令才是可扩展的方案。登录链路是这样的客户端发来command1的报文服务端解析body里的token到Redis或鉴权服务里校验token是否有效、是否过期。校验通过后调用SessionManager.bind(userId, channel)在Channel的attr里也存一份UserId然后返回登录成功响应并附上服务端当前时间方便客户端做时钟同步。校验失败就直接回错误码然后关闭连接。注意认证通过后最好把登录Handler从Pipeline中移除后续业务包不需要再走一遍鉴权逻辑这也是一种性能优化。4.2 单聊消息的路由逻辑与写回策略登录之后消息流转就进入核心逻辑了。单聊消息的body大概长这样{ fromId: 123, toId: 456, msgType: 1, content: hello, msgId: a1b2c3d4 }服务端收到上行消息后先回一个ACK给发送端表示“服务端已经收到这条消息”然后再尝试推给接收端。很多初版IM把“服务端收到消息”和“对方收到消息”混为一谈导致丢消息时无从排查。正确的认知是发送端只有收到服务端ACK才代表这条消息进入了IM系统而真正推到接收端客户端还需要接收端上屏后的确认那是后面讲的“已读回执”。路由逻辑很简单从消息体里拿到toId调用SessionManager.getChannel(toId)查目标连接。如果目标在线且Channel活跃直接writeAndFlush推过去同时给消息加上serverTime字段。如果目标不在线把消息存进离线消息表同时给发送端回一个ACK内容附带“对方离线”的状态标记。这个分支很重要不处理离线消息的话用户重新上线后会永久丢失那段时间的消息。路由时还有一个容易被忽略的坑目标Channel存在但不一定可写。如果对端的TCP接收窗口满了write方法会一直积压。Netty里可以通过channel.isWritable()判断如果不可写就不能死等要把消息先转离线策略不然内存会被OutboundBuffer撑爆。4.3 消息异步落库先记内存队列再批量写每条消息都同步写数据库这在在线用户量上来之后会很快成为瓶颈。常见做法是把落库操作从消息处理链路里拆出去。消息处理线程只负责路由和推送然后把需要存档的消息丢进一个内存队列由专门的落库线程批量消费。private final BlockingQueueMessageRecord queue new ArrayBlockingQueue(100000); private final ExecutorService storeExecutor Executors.newSingleThreadExecutor(); public void startStoreWorker() { storeExecutor.submit(() - { ListMessageRecord batch new ArrayList(200); while (true) { MessageRecord record queue.poll(200, TimeUnit.MILLISECONDS); if (record ! null) { batch.add(record); } if (batch.size() 200) { flush(batch); } if (queue.isEmpty() !batch.isEmpty()) { flush(batch); } } }); }队列一定要用有界队列防止消息生产速度远大于消费速度时OOM。落库失败不能直接丢要积攒起来重试。这套机制能扛住日常流量但要注意一个权衡如果服务器在消息落库前崩溃这部分消息就丢了。生产环境如果要更高的可靠性需要先写WAL或者同步到副本但第一阶段完全可以接受这个简单方案知道取舍在哪里就行。5. 线上才教你的三件事粘包半包、心跳保活、连接清理5.1 粘包半包解码器的长度字段设计TCP是字节流协议本身没有消息边界。发送方可能把多个小包合并发送粘包也可能把一个大包拆成多个TCP段传输半包。如果不做拆包处理服务端读到的数据就是错乱的。解决办法就是前面提到的LengthFieldBasedFrameDecoder它靠长度字段把字节流重新切成一条条完整帧。再回来看一遍参数配置参数值含义maxFrameLength1MB单帧最大长度防止恶意构造大长度字段导致OOMlengthFieldOffset10长度字段在报文中的偏移量字节索引lengthFieldLength4长度字段本身占4字节lengthAdjustment0长度字段值需要加多少等于整体长度我们的长度就是body长度所以为0initialBytesToStrip0解码后是否去掉头部我们还需要自己解析头部所以保留全部14字节报文字节序可以这样记前4字节是魔数第5字节是版本号第6字节是指令第7到10字节是seq第11到14字节是length第15字节开始是body。这里的length只统计body的长度所以lengthAdjustment为0。如果长度字段统计的是整包长度那lengthAdjustment就得设置成负的头部长度这个细节最容易算错。5.2 心跳检测与空闲踢除的姿势TCP自带的keepalive默认间隔是2小时对IM来说太慢必须用应用层心跳。客户端每隔30秒发一个心跳包服务端用IdleStateHandler做读空闲检测设置120秒内没读到任何数据就触发READER_IDLE事件。这里有个经验客户端心跳间隔必须明显小于服务端空闲阈值。如果客户端30秒一个心跳服务端120秒内至少能收到3到4个包即使网络抖动丢了一两个也不会误杀。反过来如果客户端和服务端都用60秒网络稍有抖动就可能造成大量误断开体验极差。触发READER_IDLE后的处理策略建议不要直接一刀切关闭。稳妥的做法是先发一个心跳探测包再给客户端一个空闲周期的时间回PONG如果连续两个周期都没有任何数据进来再关闭连接并清理会话。直接关闭最简单但碰上客户端只是短暂休眠的情况会反复重连反而消耗更多资源。5.3 连接异常和Channel泄漏怎么排查IM上线后最常见的线上事故就是“在线人数看着很多但用户早已离线”。这通常意味着SessionManager里堆积了大量僵尸连接。排查思路是先对比两个数字系统层面的ESTABLISHED连接数和SessionManager里的会话数。如果后者明显偏大说明unbind逻辑没触发。为什么unbind可能不触发大概率是channelInactive回调里抛了异常导致后续清理代码没执行或者代码里用了ChannelGroup但只加了没移除或者有地方在别处持有Channel引用连接关闭后依然没有被回收。遇到这类问题先把Netty的内存泄漏检测开到最严-Dio.netty.leakDetection.levelparanoid日志里如果出现“LEAK: ByteBuf.release() was not called before its garbage-collected”按堆栈定位到具体Handler多半是有地方拿到ByteBuf后没有release或者异步使用时没有retain。ByteBuf的引用计数是Netty最容易翻车的点但只要你记住一条规则谁创建谁负责释放在入站Handler中用SimpleChannelInboundHandler让它自动释放一般就不会出大问题。还有一个看似不起眼但影响巨大的点进程退出时的优雅停机。不要直接kill -9应该发送停止信号后调用shutdownGracefully它会经历一个2秒的安静期期间不再接受新任务但允许旧任务完成。否则大量连接正在收发消息时被硬杀客户端会看到成片掉线线上事故就是这么来的。压测的时候也别只看CPU和内存还要盯着线程数。正常情况连接数涨到一万EventLoop线程数应该基本不变如果线程数随连接数线性上涨说明有地方每个连接都创建了线程这种实现上生产必炸。最后说点个人的体会。我刚开始调IM服务端时最不习惯的是没有一个清晰的“请求-响应周期”所有消息都在异步回调里流动日志一多就不知道这条消息从哪个连接来、往哪个连接去。后来养成一个习惯工程里统一封装一个TraceId字段放在协议扩展头里贯穿每一个Handler日志尽量打成一条可追踪的链路。这个方法帮我省掉了无数次“这个包到底走没走到这里”的现场排查。这个项目刚开始后面还有群聊、离线消息拉取、多端同步、消息已读回执这些大头等着填。骨架建得稳后面加功能其实都是往这几条主线上挂模块。下一篇我会接着写登录网关的多端互踢和离线消息的持久化方案先把地基夯扎实再往上盖楼。
觉得有用,分享给同行:

为您的企业打造数字门面

稳重轻奢商务风格,端正雅致视觉,长效耐看不易过时。

立即咨询 →