【发布时间】:2019-07-15 19:13:04
【问题描述】:
抽象问题
有什么方法可以按照外部 observable 的原始顺序使用 mergeMap 的结果,同时仍然允许内部 observable 并行运行?
更详细的解释
让我们看看两个合并映射运算符:
-
...它接受一个映射回调,以及可以同时运行多少个内部可观察对象:
of(1, 2, 3, 4, 5, 6).pipe( mergeMap(number => api.get('/double', { number }), 3) );在此处查看实际操作:https://codepen.io/JosephSilber/pen/YzwVYNb?editors=1010
这将分别触发
1、2和3的3 个并行请求。一旦其中一个请求完成,它将触发另一个对4的请求。以此类推,始终保持 3 个并发请求,直到处理完所有值。但是,由于先前的请求可能在后续请求之前完成,因此生成的值可能会乱序。所以而不是:
[2, 4, 6, 8, 10, 12]...我们实际上可能得到:
[4, 2, 8, 10, 6, 12] // or any other permutation -
...输入
concatMap。该运算符确保所有可观察对象都按原始顺序连接,因此:of(1, 2, 3, 4, 5, 6).pipe( concatMap(number => api.get('/double', { number })) );...总会产生:
[2, 4, 6, 8, 10, 12]在此处查看实际操作:https://codepen.io/JosephSilber/pen/OJMmzpy?editors=1010
这是我们想要的,但现在请求不会并行运行。正如the documentation 所说:
concatMap等价于mergeMap,concurrency参数设置为1。
回到问题:是否有可能获得mergeMap 的好处,即可以并行运行给定数量的请求,同时仍然在原来的顺序?
我的具体问题
上面抽象地描述了这个问题。当您知道手头的实际问题时,有时会更容易推理问题,所以这里是:
-
我有一份必须发货的订单清单:
const orderNumbers = [1, 2, 3, 4, 5, 6, 7, 8, 9, 10]; -
我有一个
shipOrder方法可以实际发送订单。它返回一个Promise:const shipOrder = orderNumber => api.shipOrder(orderNumber); -
API 最多只能同时处理 5 个订单发货,所以我使用
mergeMap来处理:from(orderNumbers).pipe( mergeMap(orderNumber => shipOrder(orderNumber), 5) ); -
订单发货后,我们需要打印其发货标签。我有一个
printShippingLabel函数,给定发货订单的订单号,它将打印其发货标签。所以我订阅了我们的 observable,并在输入值时打印运输标签:from(orderNumbers) .pipe(mergeMap(orderNumber => shipOrder(orderNumber), 5)) .pipe(orderNumber => printShippingLabel(orderNumber)); -
这可行,但现在运输标签打印不正常,因为
mergeMap会根据shipOrder完成其请求的时间发出值。我想要的是标签以与原始列表相同的顺序打印。
这可能吗?
可视化
查看这里问题的可视化:https://codepen.io/JosephSilber/pen/YzwVYZb?editors=1010
您可以看到,较早的订单在后续订单发货之前就已打印出来。
【问题讨论】:
标签: javascript rxjs rxjs6 rxjs-pipeable-operators