Netty新线程调用writeAndFlush无响应问题及主动发消息方案问询
问题场景
在Java Netty开发中,自定义新线程内调用Channel的writeAndFlush方法时,客户端无法收到响应;但在EventLoop线程(如普通业务方法,当前线程属于Netty的EventLoop)中调用则能正常接收消息。相关代码如下:
private READER_ERR initReader() { // 客户端收不到消息 new Thread(new Runnable() { @Override public void run() { Channel channel = NettyChannelMap.get(clientId); if (channel != null) { ChatDto returnDto = new ChatDto(); returnDto.setClientId(clientId).setMsgType("READ").setMsg("返回数据"); channel.writeAndFlush(JSON.toJSONString(returnDto)); } } }).start(); // 客户端能正常收到消息 Channel channel = NettyChannelMap.get(clientId); if (channel != null) { ChatDto returnDto = new ChatDto(); returnDto.setClientId(clientId).setMsgType("READ").setMsg("返回数据"); channel.writeAndFlush(JSON.toJSONString(returnDto)); } }
原因分析
通过调试Netty源码可知,在AbstractChannelHandlerContext类中,Netty会判断当前线程是否属于Channel对应的EventLoop线程:
if (executor.inEventLoop()) { if (flush) { next.invokeWriteAndFlush(m, promise); } else { next.invokeWrite(m, promise); } } else { AbstractChannelHandlerContext.WriteTask task = AbstractChannelHandlerContext.WriteTask.newInstance(next, m, promise, flush); if (!safeExecute(executor, task, promise, m, !flush)) { task.cancel(); } }
其中inEventLoop方法直接对比当前线程与EventLoop绑定的线程:
public boolean inEventLoop(Thread thread) { return thread == this.thread; }
当在自定义新线程调用writeAndFlush时,inEventLoop返回false,Netty会将写任务放入EventLoop的任务队列,但如果后续没有EventLoop线程触发的任务(如无EventLoop线程发起的write/flush操作),队列中的任务可能无法及时执行;若最后一条消息由EventLoop线程发送,之前队列中的任务会被一并执行,导致客户端收到消息顺序颠倒(例如先收到EventLoop发送的"end",再收到之前队列中的1、2、3)。
解决方案
要在非EventLoop线程(如自定义新线程、回调方法)中安全给客户端发消息,必须将写操作提交到Channel对应的EventLoop线程执行,常用两种方式:
方式一:使用Channel的eventLoop().execute()提交任务
将写操作封装为Runnable,提交到Channel绑定的EventLoop线程执行:
new Thread(new Runnable() { @Override public void run() { Channel channel = NettyChannelMap.get(clientId); if (channel != null && channel.isActive()) { // 提交任务到EventLoop线程 channel.eventLoop().execute(() -> { ChatDto returnDto = new ChatDto(); returnDto.setClientId(clientId).setMsgType("READ").setMsg("返回数据"); channel.writeAndFlush(JSON.toJSONString(returnDto)); }); } } }).start();
方式二:使用ChannelHandlerContext的executor().execute()
如果持有ChannelHandlerContext实例,也可通过它的executor提交任务,效果与方式一一致:
ctx.executor().execute(() -> { ChatDto returnDto = new ChatDto(); returnDto.setClientId(clientId).setMsgType("READ").setMsg("返回数据"); ctx.writeAndFlush(JSON.toJSONString(returnDto)); });
注意事项
- 执行写操作前需判断Channel是否处于
active状态,避免对已关闭的Channel执行无效操作 - 禁止在非EventLoop线程中直接调用Channel的write/flush方法,必须通过EventLoop提交任务,既符合Netty线程模型规范,也能避免线程安全问题
内容的提问来源于stack exchange,提问作者ylq5848409

