代码之家  ›  专栏  ›  技术社区  ›  LppEdd

RxJs-将观察值数组转换为发射值数组

  •  1
  • LppEdd  · 技术社区  · 7 年前

    对不起,我想不出更好的标题了。
    我有一段代码,基本上:

    1. 用于有效(非空)cron-epressions数组的筛选器
    2. 将每个cron表达式映射到对服务的调用

    this.formGroup.valueChanges.pipe(
        op.filter(v => !!v.cronExpressions),
        op.map((v): string[] => v.cronExpressions),
        op.map((v: string[]) => v.map(cron =>
                this.cronService.getReadableForm(cron).pipe(
                    op.map(this.toDescription),
                    op.map((description): CronExpressionModel => ({ cron, description }))
                )
            )
        ),
        // What now?
    ).subscribe((cronExpressions: CronExpressionModel[]) => ...) // Expected result
    

    我想上车 subscribe() ,数组 CronExpressionModel 从所有服务调用返回。


    根据Martin的回答,当前解决方案:

    filter(v => !!v.cronExpressions),
    map(v => v.cronExpressions),
    map(cronExprs => cronExprs.map(c => this.invokeCronService(c))),
    mergeMap(serviceCalls => forkJoin(serviceCalls).pipe(defaultIfEmpty([])))
    
    2 回复  |  直到 7 年前
        1
  •  2
  •   madjaoue    7 年前

    要将流转换为数组,可以使用 toArray 操作人员

    this.formGroup.valueChanges.pipe(
        filter(v => !!v.cronExpressions),
        // transform [item1, item2...] into a stream ----item1----item2----> 
        concatMap((v): Observable<string> => from(v.cronExpressions).pipe(
            // for each of the items, make a request and wait for it to respond
            concatMap((cron: string) => this.cronService.getReadableForm(cron)),
            map(this.toDescription),
            map((description): CronExpressionModel => ({ cron, description })),
            // wait for observables to complete. When all the requests are made, 
            // return an array containing all responses
            toArray()
          )
        ).subscribe((cronExpressions: CronExpressions[]) => ...) // Expected result
    

    注:

    你可以用 mergeMap 而不是 concatMap 使请求并行化。但是你需要知道你在做什么;)

        2
  •  1
  •   martin    7 年前

    forkJoin 如果您不介意并行运行所有请求:

    switchMap(observables => forkJoin(...observables))
    

    或者,如果要按顺序运行所有这些程序:

    switchMap(observables => concat(...observables).pipe(toArray()))
    

    而不是 switchMap concatMap mergeMap 取决于你想要什么样的行为。