并行处理海量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
相关产品推荐
相关产品推荐

