RxJava2中Completable链式调用调试及Room数据库执行异常问题
解决Room + RxJava2 Completable链式调用的顺序、线程及执行问题
看起来你遇到的核心问题是Completable链式调度混乱、线程不符合预期,且最后一步操作未执行,这主要是因为你在多个子Completable中重复设置subscribeOn(尤其是误用Schedulers.trampoline()),加上缺少错误处理导致流静默终止。下面一步步帮你解决:
一、问题根源分析
subscribeOn的重复设置:RxJava中subscribeOn只对第一个生效,后续的subscribeOn会被忽略,但你在insertarCompra()的子流中多次设置subscribeOn(Schedulers.trampoline())和Schedulers.io(),导致线程调度混乱,甚至部分操作跑到主线程(比如darCompraPorFecha)。- 缺少错误处理:你的
subscribe()只传入了onComplete回调,没有onError。如果前面任意一步Completable抛出异常(比如Room操作报错),整个流会直接终止,后续的actualizarProductos()根本不会执行,而且你完全看不到错误信息。 Schedulers.trampoline()的误用:这个调度器会在当前线程排队执行任务,很容易导致线程混乱,除非你明确需要同步排队,否则数据库操作应该统一用Schedulers.io()。
二、修复后的代码示例
核心修改点:
- 移除所有子Completable中的
subscribeOn,在整个链式最外层统一设置subscribeOn(Schedulers.io()) - 为
subscribe()添加错误回调,捕获并打印异常 - 简化
insertarCompra()的内部流,避免不必要的线程切换
public void agregarCompraProductoYActualizarProductos() { insertarClienteDummy() .andThen(insertarCompra()) .andThen(actualizarProductos()) .subscribeOn(Schedulers.io()) // 统一指定IO线程执行所有数据库操作 .observeOn(AndroidSchedulers.mainThread()) .subscribe( () -> irAFacturaActivity(), throwable -> { // 必须添加错误处理,排查流终止原因 Log.e("RxError", "操作失败", throwable); } ); } private Completable insertarClienteDummy() { Cliente cliente = new Cliente(111, "Juan"); Repositorio repo = new RepositorioCliente(getApplicationContext()); return repo.agregarElemento(cliente); // 移除内部subscribeOn } private Completable insertarCompra() { String fecha = (String) DateFormat.format("yyyy-MM-dd hh:mm:ss a", Calendar.getInstance().getTime()); Repositorio repo = new RepositorioCompra(getApplicationContext()); return repo.agregarElemento(new Compra(0, 111, fecha)) .andThen( repo.darCompraPorFecha(fecha) // 复用同一个repo实例,避免重复创建 .map(compra -> crearCompraProductosPorCodigoCompra(compra.getCodigo())) .flatMapCompletable(compraProductos -> agregarCompraProductos(compraProductos)) ); // 移除所有内部subscribeOn,统一由外层调度 } // 其他方法保持不变,移除内部subscribeOn private ArrayList crearCompraProductosPorCodigoCompra(int codigoCompra){ ArrayList<CompraProducto> compraProductos = new ArrayList<>(); for (Producto p: productos) { compraProductos.add(new CompraProducto(0, codigoCompra, p.getCodigo(), p.getCantidad())); } return compraProductos; } private Completable agregarCompraProductos(ArrayList<CompraProducto> compraProductos){ RepositorioCompraProducto repo = new RepositorioCompraProducto(getApplicationContext()); return repo.agregarCompraProductos(compraProductos.toArray(new CompraProducto[0])); } private Completable actualizarProductos(){ RepositorioProducto repo2 = new RepositorioProducto(getApplicationContext()); return repo2.actualizarProductos(productos.toArray(new Producto[0])); } public void irAFacturaActivity(){ Intent i = new Intent(this, FacturaActivity.class); Bundle bundle = new Bundle(); bundle.putSerializable(PRODUCTOS, productos); i.putExtras(bundle); limpiarCache(); startActivity(i); }
三、调试RxJava2 Completable的实用技巧
因为Completable没有发射数据,只有完成或错误信号,调试起来确实比Observable麻烦,这里分享几个实用方法:
- 强制添加错误回调:这是最关键的一步!任何Completable的
subscribe()都必须传入onError,否则未捕获的异常会导致流静默终止,你根本不知道哪里出问题。 - 使用
doOnXXX操作符跟踪流程:在链式中插入doOnSubscribe、doOnComplete、doOnError来打印日志,跟踪每个步骤的执行时机和线程:insertarClienteDummy() .doOnSubscribe(disposable -> Log.d("RxTrace", "开始执行insertarClienteDummy,线程:" + Thread.currentThread().getName())) .doOnComplete(() -> Log.d("RxTrace", "完成insertarClienteDummy")) .andThen(insertarCompra()) .doOnComplete(() -> Log.d("RxTrace", "完成insertarCompra")) .andThen(actualizarProductos()) .doOnComplete(() -> Log.d("RxTrace", "完成actualizarProductos")) // ...后续调度和订阅 - 在测试环境用
blockingAwait()同步执行:如果是单元测试或调试阶段,可以用blockingAwait()替代subscribe(),强制同步执行,这样可以直接在控制台看到异常堆栈:try { insertarClienteDummy() .andThen(insertarCompra()) .andThen(actualizarProductos()) .subscribeOn(Schedulers.io()) .blockingAwait(); } catch (Exception e) { e.printStackTrace(); } - 全局捕获RxJava未处理异常:在Application类中设置RxJavaPlugins的错误处理器,避免异常被静默吞噬:
RxJavaPlugins.setErrorHandler(throwable -> { Log.e("GlobalRxError", "未处理的Rx异常", throwable); });
四、验证修复后的预期结果
修复后,所有数据库操作都会在RxCachedThreadScheduler线程执行,顺序严格按照insertarClienteDummy() → insertarCompra()(包含内部的darCompraPorFecha和agregarCompraProductos) → actualizarProductos()执行,最后触发irAFacturaActivity()。如果中间有任何错误,你会在日志中看到完整的异常堆栈,便于排查。
内容的提问来源于stack exchange,提问作者JuanDavid
相关产品推荐
相关产品推荐

