Netty单通道下大消息与循环小消息发送阻塞问题咨询
Netty大消息被循环小消息阻塞的场景分析与测试验证
场景解释
Netty的出站消息处理遵循**FIFO(先进先出)**原则:所有通过writeAndFlush提交的消息都会被加入到ChannelOutboundBuffer队列中,由IO线程按队列顺序逐个发送。
你观察到的“大消息被小消息阻塞”,核心原因是小消息的写入操作先于大消息进入出站队列。虽然你先提交了大消息的线程任务,但线程池的任务调度是异步的:主线程在提交大消息任务后,立刻循环提交1000个小消息任务,第二个单线程池会快速将这些小消息任务排入执行队列,大概率先于第一个线程池的大消息任务执行,导致大量小消息先被加入到Netty的出站队列中,大消息只能排队等待前面的小消息全部发送完成。
另外,大消息在Netty中可能会被拆分为多个符合TCP MTU的片段,但这些片段依然会按顺序排在出站队列中,不会抢占已入队的小消息的发送顺序。
测试是否存在问题
你的测试存在两个关键问题:
- 线程调度的不确定性:使用两个独立的单线程池提交任务,无法保证大消息的
writeAndFlush先于小消息执行。线程池的任务执行顺序依赖JVM的线程调度,导致测试结果存在随机性。 - “同时发送”的模拟不准确:你想模拟的是“同时发起大消息和小消息的写入”,但实际是异步提交任务,任务的执行时间点无法同步,不能真实反映“并发写入”的场景。
这种情况是否合理
这种FIFO的行为是合理且符合Netty设计初衷的:
- 大多数网络通信场景要求消息发送顺序与写入顺序一致,避免业务逻辑因乱序出现错误。
- 如果业务存在“大消息优先”的需求,可以通过自定义
ChannelHandler实现优先级调度:比如使用优先级队列管理出站消息,在write方法中根据消息类型调整排队顺序,或者直接在IO线程中优先处理大消息。
测试用例是否有效
测试用例能复现你观察到的现象,但有效性不足:
- 由于线程调度的随机性,测试结果可能不稳定,无法稳定复现“大消息被阻塞”的场景。
- 若要准确模拟“并发写入”,可以调整测试逻辑:使用
CountDownLatch让两个写入任务同时启动,避免线程调度的干扰。例如:
CountDownLatch latch = new CountDownLatch(1); ExecutorService executor = Executors.newFixedThreadPool(2); // 大消息任务 executor.submit(() -> { try { latch.await(); connect.writeAndFlush(largeBytes).addListener(f -> { if (f.isSuccess()) { log.info("channel write large message success"); } else { log.error("write large message error:", f.cause()); } }); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } }); // 小消息任务 executor.submit(() -> { try { latch.await(); int times = 1000; do { connect.writeAndFlush("rpcMessage"); times--; } while (times > 0); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } }); // 同时启动两个任务 latch.countDown();
原始测试代码
Channel connect = nettyClient.connect(); ExecutorService executorService = Executors.newSingleThreadExecutor(); executorService.submit(() -> connect.writeAndFlush(largeBytes) .addListener(new FutureListener<Void>() { public void operationComplete(Future<Void> f) throws Exception { if (f.isSuccess()) { log.info("channel write message success"); } else { log.error("write message error:", f.cause()); } } })); int times = 1000; ExecutorService executorService1 = Executors.newSingleThreadExecutor(); do { executorService1.submit(() -> connect.writeAndFlush("rpcMessage")); times--; } while (times > 0);
内容的提问来源于stack exchange,提问作者cookies_yon
相关产品推荐
相关产品推荐

