ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Netty发送字符串消息客户端收不到?从ChannelPipeline到粘包半包排查指南

Netty发送字符串消息客户端收不到?从ChannelPipeline到粘包半包排查指南 前两天帮同事排查一个Netty问题现象很典型客户端连上了服务端服务端日志也打了“send success”ctx.writeAndFlush(hello)执行完没有抛任何异常但客户端一个字节都没收到。更迷惑的是同一个服务端用ChannelGroup广播另一条消息时部分客户端又能正常收到。顺着这条线查下去才发现问题根本不是“发不出去”而是String这类消息在Netty的ChannelPipeline里要经过好几道“门”每一道门都可能把消息静默挡住。这篇文章我把完整的排查链路和常见根因拆开讲按这个顺序查基本半小时能定位。标题里提到的channel、channelGroup、ctx.writeAndFlush()是Netty服务端最常用的三个写消息姿势而“发送字符串消息客户端接收不到”几乎每个用Netty做IM、推送、长连接服务的人都遇到过。下面我直接从最外层现象开始一层层往里剥。1. 先别急着改代码确定“没收到”的四种具体表现1.1 你以为的“没收到”到底是哪一种排查Netty问题最怕一上来就改代码。同样是“客户端收不到”背后可能是完全不同的原因必须先对号入座。我把实际工作中遇到的情况归纳成四种第一种服务端控制台直接打印异常比如unsupported message type: String或者编码器相关的ClassCastException。这种情况最好定位基本可以断定是writeAndFlush的参数类型在Pipeline里没人能处理。第二种服务端没有任何异常writeAndFlush返回的ChannelFuture也是成功的但客户端一个字节都接收不到。这种情况最磨人问题往往不在“写”这个动作本身而在channel的状态、缓冲区水位或者连接是否半开。第三种客户端其实已经收到了字节但你的业务Handler没有触发或者打印出来是乱码。很多新人把这个误判为“没收到”实际上是缺少解码器或者解码器在Pipeline里的位置不对。第四种不是完全收不到而是部分客户端收不到。比如A连接能收到B连接收不到或者刚连上能收到过一段时间后收不到。这类问题通常和ChannelGroup里存的Channel“过期”有关。1.2 用一段客户端打印代码快速判断问题层次我建议在客户端Handler里临时加一段最原始的打印用来确认数据到底有没有到达TCP这一层public class DebugHandler extends SimpleChannelInboundHandlerByteBuf { Override protected void channelRead0(ChannelHandlerContext ctx, ByteBuf msg) { // 先别解码直接看原始字节 int len msg.readableBytes(); byte[] bytes new byte[len]; msg.getBytes(msg.readerIndex(), bytes); System.out.println(raw bytes len len , content new String(bytes, StandardCharsets.UTF_8)); // 处理完记得交给后面或者自己释放 ctx.fireChannelRead(msg.retain()); } Override public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) { cause.printStackTrace(); ctx.close(); } }如果这段代码能打印出raw bytes说明数据已经到了操作系统和Netty的接收缓冲区问题出在解码层或者业务Handler层如果连raw bytes都没有再回头查服务端的出站链路、channel状态和缓冲区。这个判断一步就能把排查范围缩小一半。2. 出站链路里没有“翻译官”writeAndFlush(String)为什么写不出去2.1 ChannelPipeline里发生的事情很多人对ctx.writeAndFlush(hello)的理解是把字符串通过Socket直接发出去。实际上Netty里的写操作不是一步到位的它会从当前Handler开始沿着ChannelPipeline向“头部”方向传播经过所有出站Handler最后才到达底层Socket。这个过程可以类比成地铁换乘字符串是一张写着字的纸你把它交给Pipeline这趟地铁但中途能不能被正确“翻译”成字节流取决于Pipeline里有没有装“翻译官”这个角色。Netty自带的StringEncoder就是干这个的它会把String按照指定字符集编码成ByteBuf后续的节点才能正常把ByteBuf写进Socket。如果Pipeline里没有任何出站Handler能处理String类型消息一路传到最底层的HeadContextNetty会发现“这个类型我处理不了”于是抛出类似unsupported message type: String (expected: ByteBuf, FileRegion)的异常。2.2 三种典型配置导致的静默失败我见过最多的三种错误配置第一种服务端压根没加StringEncoder直接ctx.writeAndFlush(hello)。后果就是上面的异常但如果你没有在exceptionCaught里打印日志或者日志级别被调高了这个异常会被吞掉看起来就像是“没报错但客户端没收到”。第二种加了StringEncoder但位置不对。举个例子Pipeline里先加了你的业务Handler再在它后面加StringEncoder。出站消息从你的业务Handler开始往头部方向传播如果你的StringEncoder排在业务Handler后面更靠近尾部那么ctx.writeAndFlush根本不会经过它。第三种服务端发送没问题但客户端没用StringDecoder业务Handler却写成了SimpleChannelInboundHandlerString。一旦收到的消息是ByteBuf而不是StringSimpleChannelInboundHandler的类型检查不通过你的channelRead0永远不会被触发。客户端表现就是“没收到”实际数据早就在缓冲区里躺着了。2.3 ctx.writeAndFlush和channel.writeAndFlush的起点不同这是整个章节里最容易被忽略的知识点。ctx.writeAndFlush和channel.writeAndFlush都会触发写操作但传播起点完全不同ctx.writeAndFlush(msg)从当前Handler的下一个出站Handler开始传播。channel.writeAndFlush(msg)从Pipeline的尾部开始传播。如果你的业务Handler是ServerHandlerPipeline配置是pipeline.addLast(new StringEncoder()); pipeline.addLast(new ServerHandler());此时在ServerHandler里调用ctx.writeAndFlush(hello)消息会从ServerHandler开始向前找出站HandlerStringEncoder在它前面所以会正确编码但如果你在ServerHandler里用channel.writeAndFlush(hello)起点变成尾部也会经过StringEncoder结果一样。反过来如果Pipeline配置成pipeline.addLast(new ServerHandler()); pipeline.addLast(new StringEncoder());那么ctx.writeAndFlush(hello)从ServerHandler向前找根本找不到StringEncoder消息直接裸奔到HeadContext最终触发类型异常。而channel.writeAndFlush(hello)从尾部开始却能经过StringEncoder。同一个项目里混用这两种写法很容易出现“这段代码能发、那段代码不能发”的灵异现象。3. 连接对象不对手动保存channel和ChannelGroup的几个大坑3.1 用Map/List保存Channel的问题很多人的第一版代码喜欢自己维护一个MapString, Channelprivate static final MapString, Channel CHANNEL_MAP new ConcurrentHashMap(); Override public void channelActive(ChannelHandlerContext ctx) { String userId ctx.channel().attr(AttributeKey.valueOf(userId)).get(); CHANNEL_MAP.put(userId, ctx.channel()); }这套写法本身没问题但坑在移除逻辑。如果客户端异常掉线、网络闪断channelInactive没有及时触发或者你忘了在channelInactive里remove这个Map里就留下了死连接。之后你用这个Channel去writeAndFlush方法调用不会马上报错数据只是进入了一个已经无法写入的Channel客户端自然收不到。更隐蔽的情况是channelActive里保存的时机太早。如果连接刚建立还没完成认证和业务属性初始化你在这个时间点把channel放进Map但后续真正要发消息时这个channel可能已经被关闭了。3.2 ChannelGroup的正确打开方式ChannelGroup真正好用的地方在于它内部会自动监听Channel的关闭事件连接断开后会自动把Channel从组里移除不需要你手动维护remove逻辑。标准用法是public class ServerHandler extends SimpleChannelInboundHandlerString { public static final ChannelGroup CHANNEL_GROUP new DefaultChannelGroup(GlobalEventExecutor.INSTANCE); Override public void channelActive(ChannelHandlerContext ctx) { // 加入组这里不用手动removeGroup内部会监听关闭事件 CHANNEL_GROUP.add(ctx.channel()); ctx.writeAndFlush(welcome\n); } Override public void channelInactive(ChannelHandlerContext ctx) { // 手动remove一次加深印象其实DefaultChannelGroup会自动清理 CHANNEL_GROUP.remove(ctx.channel()); } Override protected void channelRead0(ChannelHandlerContext ctx, String msg) { // 广播给所有连接 CHANNEL_GROUP.writeAndFlush([ ctx.channel().remoteAddress() ] msg \n); } }这里有一个细节DefaultChannelGroup构造器必须传一个EventExecutor常用的是GlobalEventExecutor.INSTANCE。它的作用是决定操作ChannelGroup的线程跑在哪个Executor上GlobalEventExecutor是一个全局单例的线程池适合大多数场景。3.3 group.writeAndFlush和channel.writeAndFlush该用谁很多人不理解为什么有了ChannelGroup还要单独发ctx.writeAndFlush。其实两者语义完全不同ChannelGroup.writeAndFlush(msg)广播发给组内所有Channel。ctx.writeAndFlush(msg)定向只发给当前连接。channel.writeAndFlush(msg)定向通过指定Channel发给某个连接。如果你的业务是“私聊”却习惯性用ChannelGroup.writeAndFlush那就会出现“A发消息B和C都收到了目标D却没收到”的错乱。反之如果你的业务是“群聊”却只调ctx.writeAndFlush那只有发送者自己能看到消息。我在实际项目里会把两种方式结合私聊走channel.writeAndFlush群聊走channelGroup.writeAndFlush。但不管用哪种都要先确认目标Channel确实存在于当前Group或Map中不要想当然。4. 写方法没报错不代表数据到了对端连接半开、缓冲水位与EventLoop阻塞4.1 ChannelFuture才是发送结果的真凭据ctx.writeAndFlush(msg)返回一个ChannelFuture这才是判断发送是否成功的“真凭据”。很多人忽略这个返回值只看调用没抛异常就认为成功了。正确做法是给Future挂监听ctx.writeAndFlush(hello\n).addListener((ChannelFutureListener) future - { if (future.isSuccess()) { // 只是写到了Netty的发送缓冲区不等于对端收到了 System.out.println(write success); } else { // 这里才能看到真实失败原因 future.cause().printStackTrace(); } });注意isSuccess()为true只代表Netty成功把数据交给了底层Socket发送缓冲区不保证对端应用层一定读到。如果你要严格确认“对端收到了”需要业务层ACK机制这是上层协议设计的范畴。4.2 isWritable和WriteBufferWaterMarkNetty每个Channel都有一个写缓冲区默认低水位是32KB高水位是64KB。当待写数据超过高水位时channel.isWritable()会变成false。这时候如果你继续writeAndFlush数据不会立刻写进Socket而是先在ChannelOutboundBuffer里排队表现就是“服务端疯狂发送客户端迟迟收不到”。特别是客户端处理速度跟不上服务端生产速度时这个现象特别明显。服务端日志打了一堆write success但客户端消息队列早就塞满了新数据只能排队等。这种问题的排查办法是发送前检查channel.isWritable()服务端做背压控制if (ctx.channel().isWritable()) { ctx.writeAndFlush(msg); } else { // 暂时不发把消息缓存起来或者走降级逻辑 PENDING_QUEUE.offer(msg); }也可以通过调整水位来适配高吞吐场景bootstrap.childOption(ChannelOption.WRITE_BUFFER_WATER_MARK, new WriteBufferWaterMark(64 * 1024, 128 * 1024));4.3 EventLoop线程被阻塞Netty的I/O线程EventLoop既要处理读写也要执行你自己注册的Handler逻辑。如果你在channelRead0里做了耗时操作比如同步查数据库、调外部接口、解析大JSON这个Channel对应的EventLoop线程就会被卡住。后续该Channel上所有读写操作都要排队表现就是消息延迟甚至“收不到”。正确的做法是把耗时不长的任务用业务线程池处理或者至少把耗时操作丢给独立的ExecutorOverride protected void channelRead0(ChannelHandlerContext ctx, String msg) { bizExecutor.execute(() - { // 耗时逻辑 String result longTimeTask(msg); ctx.writeAndFlush(result \n); }); }这里要特别注意ctx在跨线程后仍然可以安全使用但ctx.writeAndFlush之后不要手动释放消息Netty的传播机制会处理引用计数。4.4 半开连接与心跳TCP连接在客户端进程崩溃但没有发FIN包、或者网络中间设备断掉的情况下服务端感知不到连接异常。此时ctx.channel().isActive()依然是truewriteAndFlush也不会抛异常但数据包发出去就石沉大海客户端永远收不到。这种问题只能靠应用层心跳解决。常规做法是在Pipeline里加IdleStateHandler服务端定期检查读空闲超过阈值就判定连接失效并关闭pipeline.addLast(new IdleStateHandler(60, 0, 0, TimeUnit.SECONDS)); Override public void userEventTriggered(ChannelHandlerContext ctx, Object evt) throws Exception { if (evt instanceof IdleStateEvent) { // 60秒没收到客户端任何数据判定半开 ctx.close(); } else { super.userEventTriggered(ctx, evt); } }客户端也需要兜底如果连续N个周期没有收到服务端心跳响应主动重连。5. 粘包半包才是“字符串客户端收不到”的高频元凶5.1 粘包半包是怎么发生的TCP是流协议不关心你的业务消息边界。你连续发送“hello”和“world”接收方可能一次性收到“helloworld”也可能先收到“hell”再收到“oworld”。这就是粘包和半包。对于字符串消息来说粘包半包造成的现象尤其迷惑客户端明明打印出了“收到的内容”但内容里夹杂了其他消息、或者消息被截断导致SimpleChannelInboundHandlerString虽然触发了可你的业务逻辑因为拿不到一个完整的字符串结果跟“没收到”没区别。5.2 按行解码LineBasedFrameDecoder StringDecoder最简单的字符串边界是“一行一条消息”。只要发送方保证每条消息末尾带\n接收方就能用LineBasedFrameDecoder按行切分pipeline.addLast(new LineBasedFrameDecoder(1024)); pipeline.addLast(new StringDecoder(StandardCharsets.UTF_8)); pipeline.addLast(new ServerHandler());服务端发送时ctx.writeAndFlush(hello\n);这个方案的坑在于如果你发送的字符串没有以\n结尾LineBasedFrameDecoder会一直傻等下一个换行符结果就是客户端“什么都收不到”。很多第一次用Netty的人在这上面卡了很久——不是没有数据而是数据在FrameDecoder里攒着等换行。5.3 长度字段方案LengthFieldBasedFrameDecoder按行切分简单但如果消息内容本身包含换行符或者消息长度不固定行分隔符方案就不够用了。更可靠的是长度字段方案每个消息头部用4个字节表示正文长度接收方先截取长度再按长度读取完整消息。pipeline.addLast(new LengthFieldBasedFrameDecoder(2048, 0, 4, 0, 4)); pipeline.addLast(new StringDecoder(StandardCharsets.UTF_8)); pipeline.addLast(new ServerHandler());LengthFieldBaseFrameDecoder的参数解释一下2048单条消息最大长度超过会抛异常。0长度字段的起始偏移量这里从第0字节开始。4长度字段本身占4字节。0长度值跟正文之间没有额外偏移长度字段紧挨着正文。4解析完长度字段后跳过4字节把后面的数据作为正文。服务端发送时不能只写字符串要先把长度拼上去byte[] bytes msg.getBytes(StandardCharsets.UTF_8); ByteBuf buf Unpooled.buffer(4 bytes.length); buf.writeInt(bytes.length); buf.writeBytes(bytes); ctx.writeAndFlush(buf);如果项目里同时用StringEncoder和LengthFieldBasedFrameDecoder需要注意自定义编码器否则StringEncoder生成的ByteBuf没有长度前缀接收端的LengthFieldBasedFrameDecoder会把第一个汉字的高位字节当成长度解析结果彻底错乱。这种情况表现为“客户端打印一堆乱码”或者“解码抛异常”。5.4 字符串消息边界设计的通用建议如果做的是新项目我建议直接采用“长度字段 JSON字符串”的协议格式即4字节长度 JSON字符串。这样既解决了粘包半包问题也给后续的鉴权、版本号、消息类型留出了扩展空间。热词里很多人搜“netty粘包处理”大概率就是遇到了这类问题。字符串的“收到但数据错乱”和“彻底没收到”往往是同一个根因别只看现象表面的差异。6. 一套可以直接抄的String消息收发模板与广播代码6.1 服务端完整示例下面是基于spring-boot之外最朴素的Netty原生写法把编码器、解码器、ChannelGroup、心跳全部配置好public class ChatServer { private final int port; public ChatServer(int port) { this.port port; } public void start() throws Exception { EventLoopGroup bossGroup new NioEventLoopGroup(1); EventLoopGroup workerGroup new NioEventLoopGroup(); try { ServerBootstrap bootstrap new ServerBootstrap(); bootstrap.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) { ChannelPipeline pipeline ch.pipeline(); // 按行解码前置解决粘包半包 pipeline.addLast(new LineBasedFrameDecoder(1024)); pipeline.addLast(new StringDecoder(StandardCharsets.UTF_8)); pipeline.addLast(new StringEncoder(StandardCharsets.UTF_8)); // 60秒读空闲判活 pipeline.addLast(new IdleStateHandler(60, 0, 0, TimeUnit.SECONDS)); pipeline.addLast(new ChatServerHandler()); } }); ChannelFuture future bootstrap.bind(port).sync(); System.out.println(server start at port port); future.channel().closeFuture().sync(); } finally { bossGroup.shutdownGracefully(); workerGroup.shutdownGracefully(); } } public static void main(String[] args) throws Exception { new ChatServer(8080).start(); } }ChatServerHandlerpublic class ChatServerHandler extends SimpleChannelInboundHandlerString { public static final ChannelGroup CHANNEL_GROUP new DefaultChannelGroup(GlobalEventExecutor.INSTANCE); Override public void channelActive(ChannelHandlerContext ctx) { CHANNEL_GROUP.add(ctx.channel()); ctx.writeAndFlush(welcome\n); } Override public void channelInactive(ChannelHandlerContext ctx) { CHANNEL_GROUP.remove(ctx.channel()); } Override protected void channelRead0(ChannelHandlerContext ctx, String msg) { if (whoami.equals(msg)) { // 定向回复当前连接 ctx.writeAndFlush(your channel id ctx.channel().id() \n); } else { // 广播给组内所有连接 CHANNEL_GROUP.writeAndFlush([ ctx.channel().remoteAddress() ] say: msg \n); } } Override public void userEventTriggered(ChannelHandlerContext ctx, Object evt) throws Exception { if (evt instanceof IdleStateEvent) { ctx.close(); } else { super.userEventTriggered(ctx, evt); } } Override public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) { // 千万别删这个方法出站类型错误、解码异常都会在这里暴露 cause.printStackTrace(); ctx.close(); } }这段代码有三个关键细节一是exceptionCaught不能省否则异常被吞掉后很难排查二是LineBasedFrameDecoder(1024)的1024表示单行最大字节数超过会抛异常并连接断开三是ctx.channel().remoteAddress()就是SocketAddress用来标识来源。6.2 客户端完整示例客户端同样需要对称的编解码配置public class ChatClient { public void connect(String host, int port) throws Exception { EventLoopGroup group new NioEventLoopGroup(); try { Bootstrap bootstrap new Bootstrap(); bootstrap.group(group) .channel(NioSocketChannel.class) .option(ChannelOption.TCP_NODELAY, true) .handler(new ChannelInitializerSocketChannel() { Override protected void initChannel(SocketChannel ch) { ChannelPipeline pipeline ch.pipeline(); pipeline.addLast(new LineBasedFrameDecoder(1024)); pipeline.addLast(new StringDecoder(StandardCharsets.UTF_8)); pipeline.addLast(new StringEncoder(StandardCharsets.UTF_8)); pipeline.addLast(new ChatClientHandler()); } }); Channel channel bootstrap.connect(host, port).sync().channel(); // 发送字符串必须带换行否则服务端LineBasedFrameDecoder会一直等 channel.writeAndFlush(hello server\n); channel.closeFuture().sync(); } finally { group.shutdownGracefully(); } } }ChatClientHandlerpublic class ChatClientHandler extends SimpleChannelInboundHandlerString { Override protected void channelRead0(ChannelHandlerContext ctx, String msg) { System.out.println(recv: msg); } Override public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) { cause.printStackTrace(); ctx.close(); } }6.3 验证发送结果与常见改造方向把这些代码跑起来以后你可以在服务端加一行ChannelFuture cf ctx.writeAndFlush(hello\n); cf.addListener((ChannelFutureListener) future - { if (!future.isSuccess()) { future.cause().printStackTrace(); } });用自带的nc -vz或者写一个测试脚本连上来发送几行字符串基本能验证全链路是否正常。如果你不是做纯TCP而是做WebSocket记得把编解码器换成WebSocketServerProtocolHandler和TextWebSocketFrame字符串的收发逻辑思路是一样的只是帧类型从ByteBuf变成了TextWebSocketFrame。7. 排查清单从现象到根因的快速对照表7.1 快速对照表下面这张表是我每次排查时都会对照一遍的清单按顺序从“最常见原因”到“最隐蔽原因”排列现象可能根因验证方法修复方式服务端打印unsupported message type: StringPipeline里缺少StringEncoder查看异常堆栈在Pipeline里补上StringEncoder服务端无异常客户端无任何输出客户端没有StringDecoder或Handler泛型不匹配用DebugHandler打印原始字节补上StringDecoder确保SimpleChannelInboundHandlerString刚连上能收到过一会收不到连接半开未做心跳检查服务端连接数、客户端是否断网配置IdleStateHandler心跳保活部分客户端收不到ChannelGroup里的Channel已失效或目标Channel不在Group中打印Group里的Channel数量使用DefaultChannelGroup自动清理或确认channel加入成功数据明明发送了很多客户端偶尔收到ChannelOutboundBuffer超过水位打印channel.isWritable()做背压控制或调高水位客户端收到乱码或消息拼接在一起粘包半包缺少帧解码器观察原始字节内容加LineBasedFrameDecoder或LengthFieldBasedFrameDecoder明明调用了writeAndFlush走到这一步却什么都没发生EventLoop被耗时任务阻塞看线程栈检查channelRead0里是否有耗时代码耗时逻辑丢业务线程池7.2 最后沉淀的两条排查原则我每次排查这类问题都会提醒自己两句话。第一句写在代码注释里writeAndFlush只是把消息交给了ChannelPipeline真正决定它能不能到达对端的是消息类型有没有被正确处理、channel是不是活的、缓冲区还有没有位置。第二句写在便签上日志里没报错不等于数据到了对端。遇到“发了但没收到”不要永远盯着那行发送代码按本章的链路从头查一遍顺着channel、channelGroup、编解码器和缓冲区的方向走十有八九能在十分钟内找到那个被静默挡住的环节。
返回列表