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

并行处理海量KeyGridCells对象时内存占用过高的解决方案咨询

内存占用过高导致多进程处理失败的解决方案

问题背景

迭代器self.__listOfKeyCellsCollector包含数十万KeyGridCellsInTreatmentFilterator类型对象,每个对象通过run()方法在独立进程中处理。随着处理推进,内存持续攀升,因内存耗尽,操作系统无法创建更多进程,导致任务无法完成。

相关代码

主进程代码

with Pool(processes=int(config['MULTIPROCESSING']['proceses_count'])) as KeyGridCellsInTreatmentFilterator.pool:
        for res in KeyGridCellsInTreatmentFilterator.pool.map(func=self.run,iterable=self.__listOfKeyCellsCollector,chunksize=self.__chunckSize):
                            
            fourCornersOfKeyWindows.append(res[0])
            areasOfCoverage.append(res[1])
            areasOfCoverage2.append(res[1])
            interception.append(res[2])
            
            centerPointsInWindowInImageCoordinate.append(res[3])
            centerPointsOfWindowInEPSG3857.append(res[4]) #equal to centerPointsOfKeyWindowInCRSInEPSG25832. they are the same
            centerPointsOfWindowInEPSG4326.append(res[5])
            pixelValuesOfCenterPoints.append(res[6])                    
            centerPointsOfWindowsAsGeoJSONInEPSG4326.append(res[7])
            """pixel values satisfy threshold value in key windows representative to treatment"""
            pixelsValuesSatisfyThreshold.append(res[8])#should be equal to pixelsValuesSatisfyThresholdInTIFFImageDataset
            """average heights of key windows representative to treatment"""
            averageHeights.append(res[9])
            """save distances"""
            distanceFromCenterPointOfKeyWindows.append(res[10])
            PECTerrestrialRisk.append(res[11])
            ETRForPECTerrestrial.append(res[12])
            
            if res[20] is not None:
                veryLowRiskForPECT+=res[20].veryLowRiskForPECTerrestrialAsString()
                lowRiskForPEC+=res[20].lowRiskForPECTerrestrialAsString()
                mediumRiskForPEC+=res[20].mediumRiskForPECTerrestrialAsString()
                highRiskForPEC+=res[20].highRiskForPECTerrestrialAsString()
            
            self.__inputStringToCopyFromStatement+="{0}\t{1}\t{2}\t{3}\t{4}\t{5}\t{6}\t{7}\t{8}\t{9}\t{10}\t{11}\t{12}\t{13}\t{14}\t{15}\t{16}\t{17}\t{18}\t{19}\t{20}\t{21}\t{22}\t{23}\t{24}\t{25}\t{26}\t{27}\t{28}\n".format(
            str(DateTimeUtils.getTimeToMacroSecondsPercision()),
            str(True),
            str(True),
            str(False),
            str(json.dumps(res[0]['features'][0]['geometry'])),
            str(""),
            str(res[10]),
            str(config['DEFAULT']['distance']),
            str(res[1]),
            str(config['DEFAULT']['area_of_coverage']),
            str(res[9]),
            str(config['DEFAULT']['average_height']),
            str(res[2]),
            str(config['DEFAULT']['interception']),
            str(res[13]),
            str(res[14]),
            str(res[15]),
            str(res[16]),
            str(res[17]),
            str(res[18]),
            str(res[19]),
            str(config['DEFAULT']['empty_polygon']),
            str(config['DEFAULT']['empty_polygon']),
            str(NumpyUtils.convertToNumpyArray(res[8][0])),
            str(config['DEFAULT']['pixelValue']),
            str(res[7]),
            str(config['DEFAULT']['empty_string']),
            str(res[6]),
            str(config['DEFAULT']['pixelValue']),
            )
           
    KeyGridCellsInTreatmentFilterator.pool.join()

run方法代码

def run(self,params:CellsInTreatmentInfoCollector):
    PECTerrestrialRiskForKeyGridCellsInTreatment = []
    ETRForPECTerrestrialForKeyGridCellsInTreatment = []
    if params is not None:
        logger.info(f"None-Zero covering cells filteration/separation phase:cell filtered as belongs to treatment:{params.getAreasOfCoveragePerWindow()}")
        fourCornersOfKeyWindowsAsGeoJSONInEPSG4326 = params.getFourCornersOfWindowsAsGeoJSONInEPSG4326()
        areasOfCoveragePerKeyWindow= params.getAreasOfCoveragePerWindow()
        interceptionPerKeyWindow= params.getInterceptionPerWindow()
        selectedSiteID = params.getSelectedSiteID()
        fieldCoordinatesAsTextInWKTInEPSG4326 = params.getTreatmentGeometry()
        threshold = params.getThreshold()
        visOp0 = params.getIsVisualizeAreaOfCoverage()
        visOp1 = params.getIsVisualizeAverageHeights()
        visOp2 = params.getIsVisualizeInterception()
        visOp3 = params.getIsVisualizeEndangeredAreas()
        centerPointsInKeyWindowInImageCoordinateSystem = params.getCenterPointsOfWindowInImageCoordinateSystem()
        centerPointsOfKeyWindowInEPSG3857 = params.getCenterPointsOfWindowInEPSG3857()
        centerPointsOfKeyWindowInEPSG4326 = params.getCenterPointsOfWindowInEPSG4326()
        pixelValuesOfCenterPointsOfKeyWindow = params.getPixelValuesOfCenterPointsOfWindow()
        centerPointsOfKeyWindowsAsGeoJSONInEPSG4326 = params.getCenterPointsOfWindowsAsGeoJSONInEPSG4326()
        pixelsValuesSatisfyThreshold = params.getPixelsValuesSatisfactionToThreshold()
        averageHeightsPerKeyWindow = params.getAverageHeightsPerWindow()
        distancesFromCenterPointsOfKeyWindowsToNearestEdge = params.getDistancesFromCenterPointsOfWindowsToNearestEdge()
        # _ETRRisk:ETRRisk = cellsInTreatmentInfoCollector.getETRRiskObject()
        
        iEnvironmentalRiskParams = params.getIEnvironmentaRiskParamskObject()
        data = params.getDataObject()
        # the following if-statement is only for debugging purposes
        if((config['KEYS_OF_ETR_RISK_CALC']['enableCalcAndPopulateETRRiskTablesInAWANTIVer2WS'] in data) and (data[config['KEYS_OF_ETR_RISK_CALC']['enableCalcAndPopulateETRRiskTablesInAWANTIVer2WS']] == True)):
            AR = params.getApplicationRate()
        
        _ETRRisk = None
        if((config['KEYS_OF_ETR_RISK_CALC']['enableCalcAndPopulateETRRiskTablesInAWANTIVer2WS'] in data) and (data[config['KEYS_OF_ETR_RISK_CALC']['enableCalcAndPopulateETRRiskTablesInAWANTIVer2WS']] == True)):
            """PECTerrestrialRisk for key-grid-cells in treatment"""
            numeratorPECTerrestrialRiskForKeyGridCellsInTreatment = (float(AR) * float(1 - (interceptionPerKeyWindow/100))) * float(config['ENVIRONMENTAL_RISK']['correctionFactor'])                               
            denumeratorPECTerrestrialRiskForKeyGridCellsInTreatment = float(config['ENVIRONMENTAL_RISK']['assumedDepth']) * (float(config['ENVIRONMENTAL_RISK']['soilDensityInKiloGrams']) * 1000000)
            PECTerrestrialRiskForKeyGridCellsInTreatment.append(numeratorPECTerrestrialRiskForKeyGridCellsInTreatment/denumeratorPECTerrestrialRiskForKeyGridCellsInTreatment)   
            """ETR for PECTerrestrial for key grid-cell in treatment"""
            ETRValueForPECTerrestrialForKeyGridCellsInTreatment = (numeratorPECTerrestrialRiskForKeyGridCellsInTreatment/denumeratorPECTerrestrialRiskForKeyGridCellsInTreatment) / float(config['ENVIRONMENTAL_RISK']['EC50EWCO'])
            ETRForPECTerrestrialForKeyGridCellsInTreatment.append( ETRValueForPECTerrestrialForKeyGridCellsInTreatment )            
            """ categorizing ETR-value to corresponding risk-category for PECTerrestrial """
            insecticideSelected = iEnvironmentalRiskParams[config['KEYS_OF_INSECTICIDES_PARAMS']['insecticideSelected']]
            dose = iEnvironmentalRiskParams[config['KEYS_OF_INSECTICIDES_PARAMS']['dose']]
            doseUnit = iEnvironmentalRiskParams[config['KEYS_OF_INSECTICIDES_PARAMS']['doseUnit']]
            dateOfSpray = iEnvironmentalRiskParams[config['KEYS_OF_INSECTICIDES_PARAMS']['dateOfSpray']]
            _ETRRisk = ETRRisk()
            _ETRRisk.categorizeETRFor(
                selectedSiteID=selectedSiteID,
                ETRValue=ETRValueForPECTerrestrialForKeyGridCellsInTreatment,  
                insecticideType=insecticideSelected,
                dose=dose,
                doseUnit=doseUnit,
                dateOfSpray=dateOfSpray,
                isForPECTerrestrial=True,
                isForPECDrift=False,
                isKey=True,
                fourCornersOfWindowCorrespondsToETRValueInEPSG4326=fourCornersOfKeyWindowsAsGeoJSONInEPSG4326['features'][0]['geometry'],
                geometryOfFourCornersOfWindowCorrespondsToETRValueInEPSG4326=config['DEFAULT']['empty_polygon'])
        
        return fourCornersOfKeyWindowsAsGeoJSONInEPSG4326,areasOfCoveragePerKeyWindow,interceptionPerKeyWindow,centerPointsInKeyWindowInImageCoordinateSystem,centerPointsOfKeyWindowInEPSG3857,centerPointsOfKeyWindowInEPSG4326,pixelValuesOfCenterPointsOfKeyWindow,centerPointsOfKeyWindowsAsGeoJSONInEPSG4326,pixelsValuesSatisfyThreshold,averageHeightsPerKeyWindow,distancesFromCenterPointsOfKeyWindowsToNearestEdge,PECTerrestrialRiskForKeyGridCellsInTreatment,ETRForPECTerrestrialForKeyGridCellsInTreatment,fieldCoordinatesAsTextInWKTInEPSG4326,selectedSiteID,threshold,visOp0,visOp1,visOp2,visOp3,_ETRRisk
    else:
        raise Exception ("WTF.")        

解决方案建议

1. 替换pool.map为pool.imap或pool.imap_unordered

pool.map会一次性将所有迭代器元素加载到内存再分配给进程,改用imap/imap_unordered可实现流式处理,元素按需传递,避免一次性加载数十万对象占用内存:

  • imap保持结果顺序,与原逻辑兼容;
  • imap_unordered不保证顺序,但内存占用更低、处理速度更快,若结果顺序不影响后续逻辑可优先使用。

修改示例:

for res in KeyGridCellsInTreatmentFilterator.pool.imap(func=self.run, iterable=self.__listOfKeyCellsCollector, chunksize=self.__chunckSize):
    # 后续处理逻辑不变

2. 优化chunksize参数

过大的chunksize会导致单个进程一次性处理过多任务,内存飙升;过小则进程间通信开销增大。建议根据总任务数和进程数计算合理值:

total_tasks = len(self.__listOfKeyCellsCollector)
process_count = int(config['MULTIPROCESSING']['proceses_count'])
chunksize = max(1, total_tasks // (process_count * 4))  # 可根据实际情况调整倍数

3. 避免主进程缓存大量结果

主进程将所有结果追加到列表会持续累积内存,若这些列表仅用于生成输出字符串或写入文件,建议处理完单个结果后直接写入文件,不保存到内存:

# 主进程开头打开文件句柄
with open('output_result.txt', 'w') as output_file:
    with Pool(...) as pool:
        for res in pool.imap(...):
            # 生成格式化字符串
            line = "{0}\t{1}\t...\n".format(...)
            # 直接写入文件
            output_file.write(line)
            # 仅保留必须的累加变量(如veryLowRiskForPECT等),删除其他结果列表

4. 减少进程间传递的数据量

run方法返回21个字段,其中存在重复或非必要数据(如注释说明重复的centerPointsOfWindowInEPSG3857)。优化方向:

  • 仅返回主进程真正需要的字段;
  • 若部分字段仅用于生成输出字符串,可在子进程中直接格式化好字符串返回,减少进程间传递的数据体积。

5. 手动清理无用对象

在主进程循环中,处理完res后手动删除对象并触发垃圾回收,及时释放内存:

import gc

for res in pool.imap(...):
    # 处理res的逻辑
    ...
    # 删除res并触发垃圾回收
    del res
    gc.collect()

6. 合理限制并发进程数

若当前proceses_count超过CPU核心数过多,会加剧内存竞争。建议将进程数设置为CPU核心数(或核心数+1),避免创建过多进程消耗内存:

import multiprocessing

process_count = multiprocessing.cpu_count()
# 内存不足时可进一步减少进程数,如process_count = max(1, multiprocessing.cpu_count() - 1)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 17:50:24