Netty多客户端消息接收异常:消息被覆盖问题求助
问题:Netty多客户端消息存入队列时出现覆盖
环境与业务逻辑
- 使用Netty v4.1.90.Final开发客户端-服务端应用,本地端口9001
- 客户端:模拟新闻推送服务,先发初始化对象,收到服务端"200"响应后,定期发送包含int和string的POJO(
NewsItem) - 服务端:新闻分析器,收到初始化对象后返回"200",将业务对象存入
LinkedBlockingQueue供后续分析
预期与实际结果
预期
启动2个及以上客户端时,Netty异步接收消息,解码后按接收顺序加入队列,队列长度随消息数量递增(3个客户端各发20条,队列应存60条)
实际
消息到达NewsAnalyserHandler后出现覆盖,最终队列仅存20条,与单客户端发送的消息数量一致
关键代码片段
NewsAnalyserHandler
public class NewsAnalyserHandler extends ChannelInboundHandlerAdapter { private static final Logger logger = LoggerFactory.getLogger(NewsAnalyserHandler.class); private static final String OK_TO_SEND = "200"; private static final String STOP_SENDING = "507"; private final BlockingQueue<NewsItem> messageQueue = new LinkedBlockingQueue<>(); @Override public synchronized void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception { NewsItem request = (GeneratedNewsItem) msg; // 处理初始化对象 if (request.getHeadline().equals("INIT") && request.getPriorty() == -1) { ctx.write(OK_TO_SEND); return; } // 存入队列并返回响应 if (messageQueue.offer(request)) { logger.info("Received news item: {}", request); logger.info("Sending 200 response"); ctx.write(OK_TO_SEND); } logger.info("Number of messages: {}", messageQueue.size()); ReferenceCountUtil.release(msg); } //... }
解码器(两种尝试过的实现)
NewsItemByteDecoder
public class NewsItemByteDecoder extends ByteToMessageDecoder { @Override protected void decode(ChannelHandlerContext ctx, ByteBuf in, List<Object> out) throws Exception { while (in.isReadable()) { int strLen = in.readInt(); byte[] headlineBytes = new byte[strLen]; in.readBytes(headlineBytes); String headline = new String(headlineBytes, StandardCharsets.UTF_8); int priority = in.readInt(); NewsItem ni = GeneratedNewsItem.createNewsItem(priority, headline); ctx.fireChannelRead(ni); } } }
NewsItemDecoder
public class NewsItemDecoder extends ReplayingDecoder<NewsItem> { @Override protected void decode(ChannelHandlerContext ctx, ByteBuf in, List<Object> out) throws Exception { int strLen = in.readInt(); byte[] headlineBytes = new byte[strLen]; in.readBytes(headlineBytes); String headline = new String(headlineBytes, StandardCharsets.UTF_8); int priority = in.readInt(); NewsItem ni = GeneratedNewsItem.createNewsItem(priority, headline); out.add(ni); } }
服务端启动类NewsAnalyser
public final class NewsAnalyser { private static final int DEFAULT_PORT = 9001; private static final Logger logger = LoggerFactory.getLogger(NewsAnalyser.class); private int port; public NewsAnalyser(int port) { this.port = port; } public static void main(String[] args) throws Exception { int port = (args.length > 0) ? Integer.parseInt(args[0]) : DEFAULT_PORT; new NewsAnalyser(port).run(); } public void run() throws InterruptedException { EventLoopGroup bossGroup = new NioEventLoopGroup(); EventLoopGroup workerGroup = new NioEventLoopGroup(); try { ServerBootstrap b = new ServerBootstrap(); b.group(bossGroup, workerGroup) .channel(NioServerSocketChannel.class) .childHandler(new ChannelInitializer<SocketChannel>() { @Override protected void initChannel(SocketChannel ch) throws Exception { ch.pipeline().addLast(new NewsItemByteDecoder()) .addLast(new ServerResponseEncoder(), new NewsAnalyserHandler()); } }) .option(ChannelOption.SO_BACKLOG, 128) .childOption(ChannelOption.SO_KEEPALIVE, true); ChannelFuture future = b.bind(port).sync(); logger.info("Starting News Analyser on localhost:{}", port); future.channel().closeFuture().sync(); } finally { logger.info("Shutting down News Analyser on localhost:{}", port); workerGroup.shutdownGracefully(); bossGroup.shutdownGracefully(); } } }
已尝试的无效方案
- 切换
ReplayingDecoder<T>和ByteToMessageDecoder,效果一致 - 更换队列类型:
List、BlockingQueue、ConcurrentLinkedQueue - 给
channelRead()方法添加synchronized关键字或用ReentrantLock同步代码块
问题分析与解决建议
核心原因
Netty中每个客户端连接都会创建独立的ChannelPipeline,而你在ChannelInitializer中每次初始化Channel时都会新建一个NewsAnalyserHandler实例——也就是说,3个客户端对应3个独立的NewsAnalyserHandler对象,每个对象都有自己的messageQueue。你看到的日志是每个Handler自己的队列长度,自然只会显示20条。
解决方法
让所有NewsAnalyserHandler实例共享同一个队列,有两种实现方式:
方案1:将队列改为静态变量
public class NewsAnalyserHandler extends ChannelInboundHandlerAdapter { // ...其他代码不变 // 改为静态变量,所有Handler实例共享同一个队列 private static final BlockingQueue<NewsItem> messageQueue = new LinkedBlockingQueue<>(); // ...其他代码不变 }
方案2:使用单例队列管理类
// 单独的队列管理单例类 public class NewsQueueManager { private static final BlockingQueue<NewsItem> INSTANCE = new LinkedBlockingQueue<>(); private NewsQueueManager() {} public static BlockingQueue<NewsItem> getInstance() { return INSTANCE; } } // 在Handler中使用单例队列 public class NewsAnalyserHandler extends ChannelInboundHandlerAdapter { // ...其他代码不变 private final BlockingQueue<NewsItem> messageQueue = NewsQueueManager.getInstance(); // ...其他代码不变 }
额外注意事项
LinkedBlockingQueue本身是线程安全的,无需额外同步,可去掉channelRead()方法的synchronized关键字NewsItemByteDecoder存在逻辑问题:ByteToMessageDecoder的decode方法会被Netty自动循环调用,直到没有足够数据,你在方法内加的while(in.isReadable())可能导致重复读取或解析错误,建议改用NewsItemDecoder(ReplayingDecoder)的实现,它会自动处理半包问题
内容的提问来源于stack exchange,提问作者Logan Duff
相关产品推荐
相关产品推荐

