You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Angular实现多依赖异步HTTP请求并实时展示测试结果

用RxJS操作符重构Angular Azure DevOps仪表盘(消除嵌套订阅)

需求背景

我是Angular和JavaScript新手,正在开发首个Angular仪表盘应用,用来调用Azure DevOps API获取测试结果,核心流程:

  • 调用API获取筛选后的发布定义列表(约100条);
  • 为每个发布定义获取其最新发布版本;
  • 为每个发布版本获取对应的全部测试运行集合;
  • 收到第一批测试结果后立即展示结果表格(每个发布定义为一行,测试结果可展开)。

另外需要新增流水线数据流,流程和发布定义一致:获取流水线列表→每个流水线的最新运行记录→对应的全部测试运行。两个数据流异步独立执行,任一数据流拿到第一批结果就展示到仪表盘。

我已经用嵌套订阅实现了功能,但知道这种方式应该尽量避免,尝试用mergeMap/switchMap或forkJoin等RxJS操作符重构,但一直没成功。


原嵌套订阅实现代码

ngOnInit(): void {   
    this.router.paramMap.pipe(takeUntil(this.ngUnsubscribe))
    .subscribe(params => {
      this.teamToFilterOn = params.get('team');
      this.apiService.getReleaseDefinitions(this.teamToFilterOn as string)
      .pipe(takeUntil(this.ngUnsubscribe))
      .subscribe((releaseDefinitions: any)=> {
        if (releaseDefinitions.length === 0) {
          this.isLoading = false
        }
        releaseDefinitions.forEach((releaseDefinition: any) => {
          if (releaseDefinition.lastRelease) {
            this.apiService.getRelease(releaseDefinition.lastRelease.id)
            .pipe(takeUntil(this.ngUnsubscribe))
            .subscribe((info: PipelineOrReleaseInfo) => {
              if (info) {
                this.apiService.getTestRunsByRelease(info.releaseId)
                .pipe(takeUntil(this.ngUnsubscribe))
                .subscribe((testruns: any) => {    
                  this.isLoading = false;              
                  this.results = [...this.results, { info: info, testruns: testruns, totals: this.calculateEnvironmentTotals(testruns.testRunResults)}];             
                  this.dataSource.data = this.results;
                });
              }              
            });
          }
        });            
      });
    });
  }   

尝试的forkJoin代码片段(未完成)

ngOnInit(): void {   
    this.router.paramMap.pipe(takeUntil(this.ngUnsubscribe))
    .subscribe(params => {
      this.teamToFilterOn = params.get('team');      
      
      let releaseDefQuery =  this.apiService.getReleaseDefinitions(this.teamToFilterOn as string)
      let pipelineDefQuery = this.apiService.getPipelineDefinitions(this.teamToFilterOn as string)

      forkJoin([releaseDefQuery, pipelineDefQuery]).subscribe(definitions => {
        let releaseDefinitions = definitions[0];
        let pipelineDefinitions = definitions[1];

        releaseDefinitions.forEach((releaseDefinition: any) => {
          if (releaseDefinition.lastRelease) {
            this.apiService.getRelease(releaseDefinition.lastRelease.id)
            .pipe(takeUntil(this.ngUnsubscribe))
            .subscribe((info: PipelineOrReleaseInfo) => {
...

流水线流程的嵌套订阅代码

pipelineDefinitions.forEach((pipelineDefinition: any) => {
    this.apiService.getLatestPipelineRun(pipelineDefinition.id)
    .pipe(takeUntil(this.ngUnsubscribe))
    .subscribe((info: PipelineOrReleaseInfo) => {
      if (info) {
        this.apiService.getTestRunsByPipeline(info.pipelineRunId)
        .pipe(takeUntil(this.ngUnsubscribe))
        .subscribe((testruns: any) => {    
          this.isLoading = false;              
          this.results = [...this.results, { info: info, testruns: testruns, totals: this.calculateEnvironmentTotals(testruns.testRunResults)}];             
          this.dataSource.data = this.results;
        });
      }              
    });
}

优化后的可行代码

下面是用RxJS操作符重构后的代码,彻底消除嵌套订阅,同时满足两个数据流异步独立、实时展示结果的需求:

ngOnInit(): void {
    this.teamToFilterOn = this.router.snapshot.paramMap.get('team');    

    const releaseResults$: Observable<any> = this.apiService.getReleaseDefinitions(this.teamToFilterOn as string).pipe(
      mergeMap(releaseDefs => releaseDefs), // 将数组拆分为单个元素的Observable,遍历每个发布定义
      filter((releaseDef: any) => releaseDef.lastRelease), // 过滤掉无最新发布记录的条目
      mergeMap((releaseDef: any) => this.apiService.getRelease(releaseDef.lastRelease.id)), // 获取对应最新发布详情
      filter(releaseInfo => !!releaseInfo), // 过滤无效的发布信息
      mergeMap((releaseInfo: PipelineOrReleaseInfo) => this.apiService.getTestRunsByRelease(releaseInfo.releaseId)
      .pipe(map(testruns => ({ testruns, info: releaseInfo })))) // 获取测试运行并和发布信息关联
    );

    const pipelineResults$: Observable<any> = this.apiService.getPipelineDefinitions(this.teamToFilterOn as string).pipe(
      mergeMap(pipelineDefs => pipelineDefs), // 将数组拆分为单个元素的Observable,遍历每个流水线定义
      mergeMap((pipelineDef: any) => this.apiService.getLastPipelineRun(pipelineDef.id)), // 获取对应最新流水线运行记录
      filter(pipelineInfo => !!pipelineInfo), // 过滤无效的流水线信息
      mergeMap((pipelineInfo: PipelineOrReleaseInfo) => this.apiService.getTestRunsByPipeline(pipelineInfo.pipelineRunId)
      .pipe(map(testruns => ({ testruns, info: pipelineInfo })))) // 获取测试运行并和流水线信息关联
    );
   
    merge(releaseResults$, pipelineResults$) // 合并两个数据流,任一有结果立即触发回调
    .pipe(takeUntil(this.ngUnsubscribe)) // 组件销毁时取消所有订阅
    .subscribe(({ testruns, info }) => {
      this.isLoading = false;              
      this.results = [...this.results, { info: info, testruns: testruns, totals: this.calculateEnvironmentTotals(testruns.testRunResults)}];             
      this.dataSource.data = this.results;
    });
  }

关键操作符说明

  • mergeMap:既可以将数组拆分为单个元素的Observable(遍历列表),也能串联后续的异步API请求,完美替代嵌套订阅;
  • filter:提前过滤无效数据,减少不必要的API调用;
  • merge:合并发布和流水线两个独立数据流,只要其中一个产生结果,就会立即执行订阅回调,实现"收到第一批结果就展示"的要求;
  • takeUntil:保留原有的内存泄漏防护逻辑,确保组件销毁时所有订阅都被取消。

内容的提问来源于stack exchange,提问作者jvanderwoude

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.28 13:30:01