嵌套for循环转换为Project Reactor响应式Flux的问题咨询
Project Reactor 改造方案及问题解答
以下是匹配原业务逻辑的响应式实现代码:
Flux.fromIterable(getC()) // 对应问题1的校验逻辑实现 .filter(c -> Valid(c)) // 嵌套flatMap利用作用域保留cId,解决跨步骤参数传递问题 .flatMap(c -> { String cId = c.getId(); // 将getM返回的集合转成Flux,逐个下发M对象 return Flux.fromIterable(getM(cId)) // 对应问题3:在流内生成q参数,和cId一起传入getD .map(m -> { String q = MessageFormat.format("query:{}", m); return getD(cId, q); }) // 执行发送Kafka的副作用操作 .doOnNext(dList -> sendDToKafka(dList)); }) .subscribe();
如果你的getM、getD本身是返回Flux/Mono的响应式方法,直接替换代码中Flux.fromIterable的部分即可,无需额外做集合转响应式流的处理。
问题逐一解答
- 如何集成IsValid()校验方法:直接使用
filter操作符即可,将Valid(c)作为判断条件传入,不满足校验的C对象会直接被过滤,不会进入后续处理,完全等价于原循环中的if判断逻辑。 - 如何保留传递cId到后续步骤:有两种常用方案,一是用
Tuples.of(cId, 其他参数)把多个需要传递的参数打包成Tuple对象往下游传递;二是用嵌套flatMap的作用域天然保留上层的cId变量,后者写法更简洁,不需要额外处理参数的打包和拆解,更适配你当前的业务场景。 - 能否在流程内生成q变量传入getD:完全可以,你可以在流的任意环节生成临时变量,只要在对应作用域内就可以作为参数传入后续方法,和你在循环里定义q变量的逻辑完全一致。
- 当前代码的问题和优化建议:你当前写的代码无法实现原业务逻辑,存在多处错误:
doOnNext是副作用操作符,只会执行传入的逻辑,不会把返回值传递到下游,你调用getM(cId)之后流里的元素还是cId,后续map(m -> m.trim())实际处理的是字符串类型的cId,根本拿不到M对象- 没有对集合类型的返回值做展平处理,
getM返回的是List,你没有转成Flux下发单个M元素,后续无法遍历M集合 - 没有做参数传递的设计,cId和q都无法传入
getD方法
按上方给出的方案实现即可匹配原业务逻辑,且是Reactor的标准写法,无额外性能损耗。
内容的提问来源于stack exchange,提问作者perplexedDev
相关产品推荐
相关产品推荐

