如何用EmbeddedChannel测试RxNetty TcpClient?数据丢弃问题求助
看起来你遇到的问题是:虽然已经订阅了客户端的输入流,但使用EmbeddedChannel.writeInbound()写入数据时,RxNetty提示没有订阅者,数据被丢弃。这大概率是因为数据写入的时机早于输入流的订阅完成,或者你的ChannelProvider配置有细节问题,导致RxNetty无法正确关联channel的入站数据到Observable流。
下面是具体的原因和修复方案:
1. 核心原因:Connection还没完成初始化就写入数据
TcpClient.createConnectionRequest()返回的是Observable<Connection>,它的发射是异步的——哪怕你用Observable.just(channel)返回已有的EmbeddedChannel,RxNetty内部还是要完成一些初始化步骤(比如应用pipeline配置、绑定Connection的输入输出流)。如果在Connection还没被发射、flatMap(c -> c.getInput())还没订阅输入流的时候就调用writeInbound(),数据到达channel时自然没有订阅者处理,就会触发那个日志。
修复方案:确保Connection初始化完成后再写入数据
你可以先同步获取Connection,再订阅它的输入流,这样就能保证顺序:
// 先同步获取Connection,确保初始化完成 Connection<Object, Object> connection = TcpClient.newClient(...) .pipelineConfigurator(...) .channelProvider(providerFactory) .createConnectionRequest() .blockingFirst(); // 订阅输入流 TestSubscriber<Object> subscriber = new TestSubscriber<>(); connection.getInput().subscribe(subscriber); // 现在写入数据,确保输入流已经有订阅者 embeddedChannel.writeInbound(serverMessage); // 断言 subscriber.assertValueCount(1);
或者用TestSubscriber等待Connection的发射:
TestSubscriber<Connection<Object, Object>> connSubscriber = new TestSubscriber<>(); TcpClient.newClient(...) .pipelineConfigurator(...) .channelProvider(providerFactory) .createConnectionRequest() .subscribe(connSubscriber); // 等待Connection被发射(createConnectionRequest会发射一次然后完成) connSubscriber.awaitTerminalEvent(); connSubscriber.assertValueCount(1); // 获取Connection并订阅输入流 Connection<Object, Object> connection = connSubscriber.getOnNextEvents().get(0); TestSubscriber<Object> inputSubscriber = new TestSubscriber<>(); connection.getInput().subscribe(inputSubscriber); // 写入数据并断言 embeddedChannel.writeInbound(serverMessage); inputSubscriber.assertValueCount(1);
2. 次要原因:EmbeddedChannel未使用RxNetty的EventLoop
你的ChannelProvider直接返回了预先创建的EmbeddedChannel,它默认使用自己的ImmediateEventLoop,而RxNetty的订阅流程可能运行在另一个EventLoop上,导致事件处理顺序混乱,订阅还没完成数据就到了。
修复方案:创建EmbeddedChannel时传入RxNetty提供的EventLoop
修改你的ChannelProviderFactory实现,用传入的EventLoop创建EmbeddedChannel:
ChannelProviderFactory providerFactory = (a,b,c,d) -> eventLoop -> { EmbeddedChannel channel = new EmbeddedChannel(eventLoop); return Observable.just(channel); };
3. 排查点:Pipeline配置是否正确
如果上面的方案都不行,检查你的pipelineConfigurator:
- 不要直接调用
pipeline.clear()清除所有Handler,这样会删掉RxNetty默认的RxHandler(它负责将channel的入站数据转换成Observable流)。 - 确保你添加了解码器,能将
serverMessage(ByteBuf)正确解码成业务对象,否则getInput()的Observable不会收到任何数据。
比如,如果你的serverMessage是字符串,需要添加StringDecoder:
.pipelineConfigurator(pipeline -> { // 保留RxNetty的默认Handler,再添加自己的解码器 pipeline.addLast(new StringDecoder()); })
内容的提问来源于stack exchange,提问作者artur

