RxJava2中FlowableEmitter多次onNext后订阅者onNext无法执行求助
问题分析与解决建议
针对你遇到的FlowableEmitter.onNext多次执行后订阅者onNext无日志输出的问题,核心原因大概率是异常未捕获导致上游静默终止、线程调度阻塞或背压配置冗余/下游消费停滞,以下是具体排查和解决步骤:
1. 强制捕获读取异常并通知下游
蓝牙Socket的InputStream.read()会频繁抛出IOException(比如连接断开、读取超时),如果未捕获处理,会直接终止上游线程,导致Flowable静默停止,下游收不到任何终止通知。必须在循环中捕获异常,并通过emitter.onError()告知下游:
try { while (!emitter.isDisposed()) { int bytesRead = inputStream.read(buffer); if (bytesRead == -1) { // 输入流正常关闭,通知下游完成 emitter.onComplete(); break; } emitter.onNext(Arrays.copyOf(buffer, bytesRead)); } } catch (IOException e) { if (!emitter.isDisposed()) { emitter.onError(e); // 必须抛出异常,避免静默终止 } }
2. 循环中检查Emitter的销毁状态
每次循环前务必检查emitter.isDisposed(),如果订阅者已取消订阅(比如页面销毁),立即终止循环并释放资源,避免无效操作占用线程:
while (!emitter.isDisposed()) { // 读取逻辑... }
3. 修正线程调度与下游消费逻辑
- 上游蓝牙读取必须放在IO线程执行,用
subscribeOn(Schedulers.io()),避免阻塞主线程; - 下游如果更新UI,用
observeOn(AndroidSchedulers.mainThread()),但绝对不能在下游onNext中做阻塞操作(比如大量计算、同步IO),否则会导致缓冲区越积越多,最终上游虽然在发送,但下游无法及时处理,看起来像是无日志输出。
4. 移除冗余的背压配置
你同时使用了BackpressureStrategy.BUFFER和onBackpressureBuffer(),这属于重复配置背压策略,不仅无意义,还可能导致双重缓冲增加内存压力,直接移除onBackpressureBuffer()即可。
5. 排查上游线程是否存活
如果以上步骤都做了仍无输出,可在读取循环中添加日志,确认上游线程是否还在运行:
bytesRead = inputStream.read(buffer); Log.d(TAG, "Read " + bytesRead + " bytes"); // 打印读取状态,确认线程未挂起
修正后的完整伪代码示例
Flowable.create((FlowableEmitter<byte[]> emitter) -> { InputStream inputStream = bluetoothSocket.getInputStream(); byte[] buffer = new byte[1024]; int bytesRead; try { while (!emitter.isDisposed()) { bytesRead = inputStream.read(buffer); if (bytesRead == -1) { emitter.onComplete(); break; } Log.d(TAG, "Upstream send data: " + bytesRead + " bytes"); // 上游发送日志 emitter.onNext(Arrays.copyOf(buffer, bytesRead)); } } catch (IOException e) { if (!emitter.isDisposed()) { emitter.onError(e); } } finally { // 确保资源释放 try { inputStream.close(); bluetoothSocket.close(); } catch (IOException ignored) {} } }, BackpressureStrategy.BUFFER) .subscribeOn(Schedulers.io()) .observeOn(AndroidSchedulers.mainThread()) .subscribe( data -> Log.d(TAG, "Downstream received: " + data.length + " bytes"), error -> Log.e(TAG, "Read failed", error), () -> Log.d(TAG, "Stream completed") );
内容的提问来源于stack exchange,提问作者Taku shousei
相关产品推荐
相关产品推荐

