如何在RxJava/Reactor响应式流水线中集成异步API调用
方案合理性
将异步API集成到响应式流水线是完全合理的技术选型,甚至可以说是响应式编程的典型使用场景。异步API本身非阻塞的特性和响应式流的设计理念高度匹配,相比阻塞API+多线程池的实现,不需要占用大量线程等待IO返回,资源利用率更高,在高吞吐IO密集型场景下优势尤其明显。
对应算子说明
你示例中的magicMethod是官方标准算子flatMap,RxJava和Reactor都提供了这个算子,使用逻辑完全一致,只需要先把回调风格的异步API包装成对应响应式框架的类型即可。
RxJava 实现示例
首先将回调式异步API封装为Single(单返回值场景)或Observable(多返回值场景):
private Single<Data> wrapAsyncCall(Key key) { return Single.create(emitter -> { client.callAsync(key) .onResult(data -> { if (!emitter.isDisposed()) { emitter.onSuccess(data); } }) .onError(throwable -> { if (!emitter.isDisposed()) { emitter.onError(throwable); } }); }); }
流水线写法如下:
source // 如果需要限制并发异步请求数,可传入第二个参数:flatMap(this::wrapAsyncCall, 100) .flatMap(this::wrapAsyncCall) .buffer(500, TimeUnit.MILLISECONDS, 100) .subscribe(dataList -> storage.batchSave(dataList));
如果需要严格保持上游key的顺序和返回结果的顺序对应,可以替换flatMap为concatMap(串行执行异步请求)或者concatMapEager(并行执行但按上游顺序输出结果)。
Reactor 实现示例
封装异步API为Mono:
private Mono<Data> wrapAsyncCall(Key key) { return Mono.create(sink -> { client.callAsync(key) .onResult(data -> sink.success(data)) .onError(throwable -> sink.error(throwable)); }); }
流水线写法如下:
source // 同样支持传入并发参数控制最大同时请求数:flatMap(this::wrapAsyncCall, 100) .flatMap(this::wrapAsyncCall) .buffer(Duration.ofMillis(500), 100) .subscribe(dataList -> storage.batchSave(dataList));
内容的提问来源于stack exchange,提问作者Nikita Rybak
相关产品推荐
相关产品推荐

