Reactor中如何避免Flux出错时下游订阅者继续执行?
方案说明
你的需求本质是所有外部接口请求全部成功完成后,才允许下游保存逻辑执行,这种场景下不管用什么操作符,都必须等待上游所有fetch请求结束(无报错)后再向下游发射数据,否则就会出现部分数据提前进入保存步骤的问题。
为什么concatMap不符合预期
你替换为concatMap后出现部分数据已保存的问题,原因是concatMap处理完单个元素就会立刻向下游发射结果,无需等待上游所有元素处理完成。如果第N个元素的fetchFromExternalService报错,前N-1个已经处理完成的元素早就流入下游的saveSequentiallyToDatabase执行完成了,自然不符合你的要求。
可用实现方案
你当前使用的collectList + flatMapMany(Flux::fromIterable)是官方推荐的标准实现,逻辑清晰易懂,没有额外的冗余开销,不需要刻意修改。
如果不想显式使用collectList,也可以用buffer()操作符实现完全等价的效果,buffer()无参调用时会等待上游Flux全部完成后,一次性发射包含所有元素的List:
getFluxIdsFromDatabase() .flatMap(this::fetchFromExternalService) .buffer() .flatMapIterable(list -> list) .concatMap(this::saveSequentiallyToDatabase)
两种方案的底层逻辑完全一致,都是缓冲所有fetch结果直到上游无报错完成,再向下游逐个发射元素执行保存。
内容的提问来源于stack exchange,提问作者user2635874
相关产品推荐
相关产品推荐

