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

Angular中mergeMap与concatMap链式处理引发内存泄漏问题

问题描述

需要将1TB文件夹上传至Blob Storage,设计了并行处理流水线:用mergeMap实现文件并行处理,每个文件通过concatMap执行一系列返回Observable的步骤。流水线可正常处理文件,但存在内存泄漏问题:每处理一个文件内存都无法释放,仅处理1GB文件就导致浏览器抛出内存不足异常。

HTML代码

<!--INPUT FILE START-->
<label for="folderUploadInput" class="folder-input-label hero-background">
  <span>{{selectedRootFolder}}</span>
  <mat-icon>folder</mat-icon>
  <mat-spinner diameter="50" class="mt-1" *ngIf="inProcessFiles.length>0"></mat-spinner>
  <input id="folderUploadInput" class="hide-element" type="file"
         (change)="onFolderSelected($event)" webkitdirectory directory/>
</label>
<!--INPUT FILE END-->

TypeScript组件代码

onFolderSelected(event: any) {
  if (event.target.files.length > 0) {
    const splitFolder = this.splitFolderStructure(event.target.files[0]);
    this.selectedRootFolder = 'Folder selected - ' + splitFolder[0] + ', Total files - ' + event.target.files.length;
    this.processArrayOfFiles(event.target.files);
  } else {
    this.selectedRootFolder = 'No files in selected folder'
  }
}

splitFolderStructure(file: File) {
  return file.webkitRelativePath.split('/');
}

processArrayOfFiles(fileList: File[]) {
  this.totalFiles = fileList.length;
  if (this.totalFiles > 0) {
    const fileListToProcess: FileUploadTracking[] = this.initFileList(fileList);
    const parallelProcessingOfFiles = 10
    // Entire design is a pipeline per file basis. each file goes through a series of steps.
    of(fileListToProcess).pipe(
      // take one element at a time
      mergeMap(fileUploadTrackingList => fileUploadTrackingList),
      // process
      mergeMap((fileUploadTracking: FileUploadTracking) => {
        this.inProcessFiles.push(fileUploadTracking);
        return this.checkIfMediaIsAlreadyPresent(fileUploadTracking).pipe(
          concatMap((fileUploadTracking: FileUploadTracking) => {
            // 1. Check if file present and creating it if not present
            if (this.checkIfStepComplete(STEPS_LIST.MEDIA_CREATED, fileUploadTracking)) {
              return of(fileUploadTracking);
            } else {
              return this.createMediaUploadDetails(fileUploadTracking);
            }
          }),
          concatMap((fileUploadTracking: FileUploadTracking) => {
            // 2. update state to media details created
            if (this.checkIfStepComplete(STEPS_LIST.MEDIA_CREATED, fileUploadTracking)) {
              return of(fileUploadTracking);
            } else {
              return this.updateStateToMediaCreated(fileUploadTracking);
            }
          }),
          concatMap((fileUploadTracking: FileUploadTracking) => {
            // 3. Check if the patient record is found, else stop the pipeline
            if (this.checkIfStepComplete(STEPS_LIST.PATIENT_FOUND, fileUploadTracking)) {
              return of(fileUploadTracking);
            } else {
              return this.checkPatientExistForFile(fileUploadTracking);
            }
          }),
          concatMap((fileUploadTracking: FileUploadTracking) => {
            // 4. update to patient found state
            if (this.checkIfStepComplete(STEPS_LIST.PATIENT_FOUND, fileUploadTracking)) {
              return of(fileUploadTracking);
            } else {
              return this.updateStateOfCurrentFileToPatientFound(fileUploadTracking);
            }
          }),
          concatMap((fileUploadTracking: FileUploadTracking) => {
            // 5. create thumbnail all the time becuase thumbnail is not saved till file is uploaded,
            // if the file is already uploaded skip the step
            if (this.checkIfStepComplete(STEPS_LIST.FILE_UPLOADED, fileUploadTracking)) {
              return of(fileUploadTracking);
            } else {
              return this.createThumbNailOfFile(fileUploadTracking);
            }
          }),
          concatMap((fileUploadTracking: FileUploadTracking) => {
            // 6. update thumbnail state, if the file is already uploaded skip the step
            if (this.checkIfStepComplete(STEPS_LIST.FILE_UPLOADED, fileUploadTracking)) {
              return of(fileUploadTracking);
            } else {
              return this.updateStateOfCurrentFileToThumbnailGenerated(fileUploadTracking);
            }
          }),
          concatMap((fileUploadTracking: FileUploadTracking) => {
            // 7. upload file to blob
            if (this.checkIfStepComplete(STEPS_LIST.FILE_UPLOADED, fileUploadTracking)) {
              return of(fileUploadTracking);
            } else {
              return this.uploadFileUsingGenericFileHandlerService(fileUploadTracking);
            }
          }),
          concatMap((fileUploadTracking: FileUploadTracking) => {
            // 8. update uploaded state
            if (this.checkIfStepComplete(STEPS_LIST.FILE_UPLOADED, fileUploadTracking)) {
              return of(fileUploadTracking);
            } else {
              return this.updateFileDetailsIdAndUpdateState(fileUploadTracking);
            }
          }),
          concatMap(fileUploadTracking => {
            // 9. Create media activity log entry
            if (this.checkIfStepComplete(STEPS_LIST.ACTIVITY_LOG_CREATED, fileUploadTracking)) {
              return of(fileUploadTracking);
            } else {
              return this.createANewMediaActivityLogEntry(fileUploadTracking);
            }
          }),
          // WTF why the hell is this returning any. doesn't make any sense
          concatMap(fileUploadTracking => {
            // 10. update activity log state
            if (this.checkIfStepComplete(STEPS_LIST.ACTIVITY_LOG_CREATED, fileUploadTracking)) {
              return of(fileUploadTracking);
            } else {
              return this.updateActivityLogIdState(fileUploadTracking);
            }
          }),
          concatMap((fileUploadTracking: FileUploadTracking) => {
            // 11. Link media to patient media assoc
            if (this.checkIfStepComplete(STEPS_LIST.MEDIA_LINKED, fileUploadTracking)) {
              return of(fileUploadTracking);
            } else {
              return this.linkToPatientMediaAssoc(fileUploadTracking);
            }
          }),
          concatMap((fileUploadTracking: FileUploadTracking) => {
            // 12. Link media to patient media assoc
            if (this.checkIfStepComplete(STEPS_LIST.MEDIA_LINKED, fileUploadTracking)) {
              return of(fileUploadTracking);
            } else {
              return this.updateStateToPatientMediaAttached(fileUploadTracking);
            }
          }),
          catchError((fileUploadTracking: FileUploadTracking) => {
            console.error(`large-uploader.component:catchError:fileUploadTracking - `, fileUploadTracking);
            this.erroredFiles.push(fileUploadTracking);
            return this.updateStateOfMediaUploadDetails(fileUploadTracking, {isFailed: true}, fileUploadTracking.error.message, true);
          })
        )
      }, parallelProcessingOfFiles),
    ).subscribe({
      // weird error, had to forcefully put res:any even though its clearly res:FileUploadTracking
      next: (res: any) => {
        this.inProcessFiles = this.inProcessFiles.filter(item => item.mediaUploadDetails.OriginalFilePath !== res.mediaUploadDetails.OriginalFilePath);
        this.processedFiles++;
        this.overAllProcess = +(((this.processedFiles / this.totalFiles) * 100).toFixed(1));
        this.getTableData();
      },
      error: err => {
        console.error("large-uploader.component:processArrayOfFiles:error:  err ", err);
        this.insightService.logException(err, 'large-uploader.component', 'error ', '306', SeverityLevel.Error);
      },
      complete: () => {
        console.log(`large-uploader.component:complete:processArrayOfFiles:complete - completed`);
        this.callAPIMetadataOperation();
        this.snackBar.open('Upload complete', 'OK');
      }
    })
  }
}

内存泄漏原因分析及修复方案

核心问题点

  1. 订阅未清理:整个上传流水线的Observable订阅未保存,组件销毁时无法取消订阅,导致File对象、跟踪对象等被持续引用,无法被GC回收。
  2. File引用残留:FileUploadTracking对象持有原始File引用,处理完成后未主动释放,即使从inProcessFiles移除,erroredFiles或其他集合仍可能持有引用。
  3. 并行度过高:同时处理10个大文件,每个文件的处理步骤(如缩略图生成、文件上传)可能占用大量内存,且多个文件的资源无法及时释放。
  4. 中间Observable未正确完成:部分步骤的Observable未触发complete,导致订阅长期挂起,占用内存。

具体修复步骤

1. 管理订阅生命周期

  • 在组件中添加销毁信号Subject和订阅对象:
    import { Subscription, Subject } from 'rxjs';
    import { takeUntil } from 'rxjs/operators';
    
    private uploadSubscription?: Subscription;
    private destroy$ = new Subject<void>();
    
    ngOnDestroy() {
      this.destroy$.next();
      this.destroy$.complete();
      this.uploadSubscription?.unsubscribe();
    }
    
  • 在流水线的pipe中添加takeUntil(this.destroy$),确保组件销毁时自动终止所有订阅:
    of(fileListToProcess).pipe(
      mergeMap(fileUploadTrackingList => fileUploadTrackingList),
      mergeMap((fileUploadTracking: FileUploadTracking) => {
        // 原代码逻辑
      }, parallelProcessingOfFiles),
      takeUntil(this.destroy$) // 添加此行
    ).subscribe(...)
    
  • 保存订阅对象:
    this.uploadSubscription = of(fileListToProcess).pipe(...).subscribe(...)
    

2. 主动释放File引用

  • 文件处理完成(成功/失败)后,清除FileUploadTracking中的File引用:
    // next回调中
    next: (res: FileUploadTracking) => {
      const index = this.inProcessFiles.findIndex(item => item.mediaUploadDetails.OriginalFilePath === res.mediaUploadDetails.OriginalFilePath);
      if (index !== -1) {
        this.inProcessFiles.splice(index, 1);
      }
      // 释放File引用
      res.file = null; // 假设FileUploadTracking包含file属性
      this.processedFiles++;
      this.overAllProcess = +(((this.processedFiles / this.totalFiles) * 100).toFixed(1));
      this.getTableData();
    }
    
  • 在catchError中同样释放:
    catchError((fileUploadTracking: FileUploadTracking) => {
      console.error(`large-uploader.component:catchError:fileUploadTracking - `, fileUploadTracking);
      fileUploadTracking.file = null; // 释放File引用
      this.erroredFiles.push(fileUploadTracking);
      return this.updateStateOfMediaUploadDetails(fileUploadTracking, {isFailed: true}, fileUploadTracking.error.message, true);
    })
    

3. 降低并行度

将parallelProcessingOfFiles从10调整为3-5,减少同时在内存中处理的大文件数量:

const parallelProcessingOfFiles = 4;

4. 确保中间Observable正确完成

检查createThumbNailOfFile、uploadFileUsingGenericFileHandlerService等方法,确保内部的Observable都能触发complete。例如,文件上传完成后要主动调用complete(),避免订阅长期存在。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 03:30:56