【问题标题】:RxJS mergeMap() with original order具有原始顺序的 RxJS mergeMap()
【发布时间】:2019-07-15 19:13:04
【问题描述】:

抽象问题

有什么方法可以按照外部 observable 的原始顺序使用 mergeMap 的结果,同时仍然允许内部 observable 并行运行?


更详细的解释

让我们看看两个合并映射运算符:

  • mergeMap

    ...它接受一个映射回调,以及可以同时运行多少个内部可观察对象:

      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

    ...输入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 的好处,即可以并行运行给定数量的请求,同时仍然在原来的顺序?


我的具体问题

上面抽象地描述了这个问题。当您知道手头的实际问题时,有时会更容易推理问题,所以这里是:

  1. 我有一份必须发货的订单清单:

     const orderNumbers = [1, 2, 3, 4, 5, 6, 7, 8, 9, 10];
    
  2. 我有一个shipOrder 方法可以实际发送订单。它返回一个Promise:

     const shipOrder = orderNumber => api.shipOrder(orderNumber);
    
  3. API 最多只能同时处理 5 个订单发货,所以我使用 mergeMap 来处理:

     from(orderNumbers).pipe(
         mergeMap(orderNumber => shipOrder(orderNumber), 5)
     );
    
  4. 订单发货后,我们需要打印其发货标签。我有一个printShippingLabel 函数,给定发货订单的订单号,它将打印其发货标签。所以我订阅了我们的 observable,并在输入值时打印运输标签:

     from(orderNumbers)
         .pipe(mergeMap(orderNumber => shipOrder(orderNumber), 5))
         .pipe(orderNumber => printShippingLabel(orderNumber));
    
  5. 这可行,但现在运输标签打印不正常,因为mergeMap 会根据shipOrder 完成其请求的时间发出值。我想要的是标签以与原始列表相同的顺序打印。

这可能吗?


可视化

查看这里问题的可视化:https://codepen.io/JosephSilber/pen/YzwVYZb?editors=1010

您可以看到,较早的订单在后续订单发货之前就已打印出来。

【问题讨论】:

标签: javascript rxjs rxjs6 rxjs-pipeable-operators


【解决方案1】:

我确实设法部分解决了它,所以我将它发布在这里作为我自己问题的答案。

我还是很想知道处理这种情况的规范方法。


一个复杂的解决方案

  1. 创建一个自定义运算符,该运算符接受具有索引键的值(Typescript 用语中的{ index: number }),并保留值的缓冲区,仅根据它们的index 发出它们顺序。

  2. 将原始列表映射为嵌入了 index 的对象列表。

  3. 将其传递给我们的自定义 sortByIndex 运算符。

  4. 将这些值映射回其原始值。

这是sortByIndex 的样子:

function sortByIndex() {
    return observable => {
        return Observable.create(subscriber => {
            const buffer = new Map();
            let current = 0;
            return observable.subscribe({
                next: value => {
                    if (current != value.index) {
                        buffer.set(value.index, value);
                    } else {
                        subscriber.next(value);
                    
                        while (buffer.has(++current)) {
                            subscriber.next(buffer.get(current));
                            buffer.delete(current);
                        }
                    }
                },
                complete: value => subscriber.complete(),
            });
        });
    };
}

使用sortByIndex 运算符,我们现在可以完成整个管道:

of(1, 2, 3, 4, 5, 6).pipe(
    map((number, index) => ({ number, index })),
    mergeMap(async ({ number, index }) => {
        const doubled = await api.get('/double', { number });
        return { index, number: doubled };
    }, 3),
    sortByIndex(),
    map(({ number }) => number)
);

在此处查看实际操作:https://codepen.io/JosephSilber/pen/zYrwpNj?editors=1010

创建concurrentConcat 运算符

事实上,有了这个sortByIndex 运算符,我们现在可以创建一个通用的concurrentConcat 运算符,它将在内部执行与{ index: number, value: T } 类型之间的转换:

function concurrentConcat(mapper, parallel) {
    return observable => {
        return observable.pipe(
            mergeMap(
                mapper,
                (_, value, index) => ({ value, index }),
                parallel
            ),
            sortByIndex(),
            map(({ value }) => value)
        );
    };
}

然后我们可以使用这个concurrentConcat 运算符而不是mergeMap,它现在会按照原来的顺序发出值:

of(1, 2, 3, 4, 5, 6).pipe(
    concurrentConcat(number => api.get('/double', { number }), 3),
);

在此处查看实际操作:https://codepen.io/JosephSilber/pen/pogPpRP?editors=1010

所以要解决我原来的订单发货问题:

from(orderNumbers)
    .pipe(concurrentConcat(orderNumber => shipOrder(orderNumber), maxConcurrent))
    .subscribe(orderNumber => printShippingLabel(orderNumber));

在此处查看实际操作:https://codepen.io/JosephSilber/pen/rNxmpWp?editors=1010

您可以看到,即使后来的订单最终可能会在较早的订单之前发货,但标签始终按原始顺序打印。


结论

这个解决方案甚至不完整(因为它不处理发出多个值的内部可观察对象),但它需要一堆自定义代码。这是一个常见的问题,我觉得必须有一种更简单(内置)的方法来解决这个问题:|

【讨论】:

  • 我正在尝试解决同样的问题。我想运行并行 HTTP 请求以最小化加载时间,但我只希望它们按顺序发出:``` lang-js const first = of('res1').pipe(delay(1250)); const second = of('res2').pipe(delay(1000)); const third = of('res3').pipe(delay(1500)); whatGoesHere( first, second, third ).subscribe(res => { console.log(res) });``` 输出将是 @1250ms 'res1'-->'res2' @1500ms --> 'res3'
【解决方案2】:

您可以使用此运算符:sortedMergeMap、example。

const DONE = Symbol("DONE");
const DONE$ = of(DONE);
const sortedMergeMap = <I, O>(
  mapper: (i: I) => ObservableInput<O>,
  concurrent = 1
) => (source$: Observable<I>) =>
  source$.pipe(
    mergeMap(
      (value, idx) =>
        concat(mapper(value), DONE$).pipe(map(x => [x, idx] as const)),
      concurrent
    ),
    scan(
      (acc, [value, idx]) => {
        if (idx === acc.currentIdx) {
          if (value === DONE) {
            let currentIdx = idx;
            const valuesToEmit = [];
            do {
              currentIdx++;
              const nextValues = acc.buffer.get(currentIdx);
              if (!nextValues) {
                break;
              }
              valuesToEmit.push(...nextValues);
              acc.buffer.delete(currentIdx);
            } while (valuesToEmit[valuesToEmit.length - 1] === DONE);
            return {
              ...acc,
              currentIdx,
              valuesToEmit: valuesToEmit.filter(x => x !== DONE) as O[]
            };
          } else {
            return {
              ...acc,
              valuesToEmit: [value]
            };
          }
        } else {
          if (!acc.buffer.has(idx)) {
            acc.buffer.set(idx, []);
          }
          acc.buffer.get(idx)!.push(value);
          if (acc.valuesToEmit.length > 0) {
            acc.valuesToEmit = [];
          }
          return acc;
        }
      },
      {
        currentIdx: 0,
        valuesToEmit: [] as O[],
        buffer: new Map<number, (O | typeof DONE)[]>([[0, []]])
      }
    ),
    mergeMap(scannedValues => scannedValues.valuesToEmit)
  );

【讨论】:

    【解决方案3】:

    结果

    https://youtu.be/NEr6qfPlahY

    request 1
    request 2
    request 3
    response 3
    request 4
    response 1
    request 5
    1
    response 4
    request 6
    response 2
    request 7
    2
    3
    4
    response 6
    request 8
    response 5
    request 9
    5
    6
    response 7
    request 10
    7
    response 9
    response 10
    response 8
    8
    9
    10
    

    代码

    https://stackblitz.com/edit/js-5kvwl6?file=index.js

    import { range, Subject, from, of } from 'rxjs';
    import { concatMap, share, map, concatAll, delayWhen } from 'rxjs/operators';
    
    const pipeNotifier = new Subject().pipe(share());
    
    range(1, 10)
      .pipe(
        // 1. Make Observable controlled by pipeNotifier
        concatMap((v) => of(v).pipe(delayWhen(() => pipeNotifier))),
        // 2. Submit the request
        map((v) =>
          from(
            (async () => {
              console.log('request', v);
              await wait();
              console.log('response', v);
    
              pipeNotifier.next();
    
              return v;
            })()
          )
        ),
        // 3. Keep order
        concatAll()
      )
      .subscribe((x) => console.log(x));
    
    // pipeNotifier controler
    range(0, 3).subscribe(() => {
      pipeNotifier.next();
    });
    
    function wait() {
      return new Promise((resolve) => {
        const random = 5000 * Math.random();
        setTimeout(() => resolve(random), random);
      });
    }
    

    【讨论】:

      【解决方案4】:

      你想要的是这个:

      from(orderNumbers)
        .pipe(map(shipOrder), concatAll())
        .subscribe(printShippingLabel)
      

      说明:

      管道中的第一个运算符是ma​​p。它立即为每个值调用 shipOrder(因此后续值可能会启动并行请求)。

      第二个运算符 concatAll 将解析后的值按正确的顺序排列。

      (我简化了代码;concatAll() 等价于 concatMap(identity)。)

      【讨论】:

      • IMO 这将是完美的解决方案 如果 shipOrder 负责“批处理”可以同时“发货”的最大订单数量。换句话说:这个答案没有考虑到不能超过n 订单可以同时“发货”的要求。
      • 是的,你是对的。并发控制仅与 mergeMap(或 mergeAll)运算符一起提供,但是,当我们需要顺序处理响应时,它们是不切实际的。如果需要限制并行请求的数量,恐怕没有简单的解决方案(见上面的其他答案)。
      猜你喜欢
      • 1970-01-01
      • 2020-03-28
      • 2021-07-18
      • 2019-01-31
      • 2018-12-07
      • 2018-07-31
      • 2022-01-24
      • 2022-08-06
      • 2018-09-16
      相关资源
      最近更新 更多