【问题标题】:Emit actions before `toArray` - Redux Observable在 `toArray` 之前发出动作 - Redux Observable
【发布时间】:2019-03-01 11:30:47
【问题描述】:

我正在使用 Redux Observable,需要解决从史诗触发动作时的时间问题。

我有一组项目,我想循环这些项目以便对每个项目进行 AJAX 调用。收到 AJAX 响应后,我想立即执行一些操作。在原始数组中的每个项目的所有 AJAX 响应都返回后,我想触发更多操作。

如何在 timer 过期后立即触发这些操作,即使原始数组尚未完成循环?

const someEpic = action$ => (
    action$
    .pipe(
        ofType(SOME_ACTION),
        switchMap(({ payload }) => (
            from(payload) // This is an array of items
            .pipe(
                mergeMap(() => (
                    timer(5000) // This is my AJAX call
                    .pipe(
                        map(loadedAction),
                    )
                )),
                toArray(),
                map(anotherAction),
            )
        ))
    )
)

【问题讨论】:

    标签: rxjs rxjs5 redux-observable rxjs6 rxjs-pipeable-operators


    【解决方案1】:

    可能最简单的方法实际上是使用tap 在任何你想要的地方发出动作。这是假设您可以访问商店。例如:

    tap(result => this.store.dispatch(...))
    

    然而,更“Rx”的方式是使用multicast 拆分链,然后立即重新发送一部分(即加载进度),另一半与toArray() 链接以收集所有结果,然后将其转换为另一个指示加载完成的动作。

    import { range, Subject, of } from 'rxjs';
    import { multicast, delay, merge, concatMap, map, toArray } from 'rxjs/operators';
    
    const process = v => of(v).pipe(
      delay(1000),
      map(p => `processed: ${p}`),
    );
    
    range(1, 5)
      .pipe(
        concatMap(v => process(v)),
        multicast(
          () => new Subject(), 
          s => s.pipe(
            merge(s.pipe(toArray()))
          )
        ),
      )
      .subscribe(console.log);
    

    现场演示:https://stackblitz.com/edit/rxjs6-demo-k9hwtu?file=index.ts

    【讨论】:

    • 我认为多播在这里是一个很好的概念,但是你如何使用 redux-observable 来做到这一点?最新版本发送state$而不是store所以你不能随意发送。
    【解决方案2】:

    这需要 2 个史诗。

    const ajaxCallerEpic = action$ => (
        action$
        .pipe(
            ofType(AJAX_ACTION),
            switchMap(({
                payload,
                payloadId,
            }) => (
                merge(
                    from(payload) // This is an array of items
                    .pipe(
                        mergeMap(() => (
                            timer(5000) // This is my AJAX call
                            .pipe(
                                switchMap(() => (
                                    merge(
                                        of(loadedAction),
                                        of(
                                            sentData(
                                                payloadId
                                            )
                                        ),
                                    )
                                )),
                            )
                        )),
                    ),
                )
            ))
        )
    )
    
    const ajaxResponsesEpic = action$ => (
        action$
        .pipe(
            ofType(AJAX_ACTION),
            switchMap(({
                payload,
                payloadId,
            }) => (
                action$
                .pipe(
                    ofType(SENT_DATA_ACTION),
                    filter(({ id }) => (
                        id === payloadId
                    )),
                    bufferCount(
                        payload
                        .length
                    ),
                    map(anotherAction),
                )
            ))
        )
    )
    

    重要的部分是第二个SENT_DATA_ACTION。我在调用它时传递了一个唯一的 ID,以确保您正在收听正确的 ID。如果您不发送所有这些,只要浏览器打开,它就会一直在监听。您始终可以在内部侦听器上添加超时以确保它完成。另一个问题是如果ajaxResponsesEpic 设置action$ 侦听器的时间晚于ajaxCallerEpic 运行AJAX 调用的时间。

    它很可能最终成为race 条件。为了解决这些问题,您需要先执行ajaxResponsesEpic,以便它设置动作侦听器,同时在其侦听后启动 AJAX 调用。

    像这样:

    const ajaxCallerEpic = action$ => (
        action$
        .pipe(
            ofType(AJAX_READY_ACTION),
            switchMap(({
                payload,
                payloadId,
            }) => (
                merge(
                    from(payload) // This is an array of items
                    .pipe(
                        mergeMap(() => (
                            timer(5000) // This is my AJAX call
                            .pipe(
                                switchMap(() => (
                                    merge(
                                        of(loadedAction),
                                        of(
                                            sentAjaxData(
                                                payloadId
                                            )
                                        ),
                                    )
                                )),
                            )
                        )),
                    ),
                )
            ))
        )
    )
    
    const ajaxResponsesEpic = action$ => (
        action$
        .pipe(
            ofType(AJAX_ACTION),
            switchMap(({
                payload,
                payloadId,
            }) => (
                merge(
                    (
                        action$
                        .pipe(
                            ofType(SENT_AJAX_DATA_ACTION),
                            filter(({ id }) => (
                                id === payloadId
                            )),
                            bufferCount(
                                payload
                                .length
                            ),
                            map(anotherAction),
                        )
                    ),
                    (
                        of(ajaxReadyAction)
                    ),
                )
            ))
        )
    )
    

    【讨论】:

      猜你喜欢
      • 2022-01-27
      • 2021-10-29
      • 1970-01-01
      • 2018-02-06
      • 2020-12-29
      • 1970-01-01
      • 2017-10-05
      • 1970-01-01
      • 2020-09-14
      相关资源
      最近更新 更多