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'); } }) } }
内存泄漏原因分析及修复方案
核心问题点
- 订阅未清理:整个上传流水线的Observable订阅未保存,组件销毁时无法取消订阅,导致File对象、跟踪对象等被持续引用,无法被GC回收。
- File引用残留:
FileUploadTracking对象持有原始File引用,处理完成后未主动释放,即使从inProcessFiles移除,erroredFiles或其他集合仍可能持有引用。 - 并行度过高:同时处理10个大文件,每个文件的处理步骤(如缩略图生成、文件上传)可能占用大量内存,且多个文件的资源无法及时释放。
- 中间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
相关产品推荐
相关产品推荐

