Python多进程脚本未进入ComputeResult函数问题求助
我是自学Python的新手,若有明显疏漏还请见谅。以下代码中的文件名已隐去,该部分功能正常,我也简化了计算逻辑以便清晰展示需求。
我有一个文本文件,每行包含数值和字符串字段,格式为0,0,0,a,b,c,d,我将其导入数组并追加计算所需的列。
我需要循环数组中的每一行,与其他所有行进行比对,这在计算和内存层面开销较大。比对完成后,我需要更新数组中这两行的数据。
我的计划是使用multiprocessing pool和共享数组,让每个进程操作同一数组副本以提升速度,但这会存在多个进程同时写入同一行的情况。
我编写了如下代码,但问题是:代码运行无报错,却未进入ComputeResult函数——既不打印Computing result for...,断点也无法触发。请问这是什么原因?恳请帮助。
import numpy as np import multiprocessing as mp def AnalysisSettingsFunction(FileNumber): FilesList = [] FilesList.append(r"\\...File to import.txt") if FileNumber != -1 : ImportFile = (FilesList[FileNumber]) else: ImportFile = len(FilesList) OtherVariable = 21 return ImportFile, OtherVariable def ComputeResult(Args): ProcessLock, FileNumber, DataSet, RowNumber, LineCount = Args print(f"Computing result for {RowNumber}", flush=True) for C in range(RowNumber+1, LineCount): RowData = DataSet[RowNumber] OtherRow = DataSet[C] Value1 = RowData[0] Value2 = RowData[1] Value3 = OtherRow[0] Value4 = OtherRow[1] with ProcessLock: if Value1 > Value3: RowData[12] = Value1 * Value2 OtherRow[12] = Value1 * Value2 else: RowData[12] = Value3 * Value4 OtherRow[12] = Value3 * Value4 def main(FileNumber): AnalysisSettings = AnalysisSettingsFunction(FileNumber) FileToImport = AnalysisSettings[0] #Importing data into array RawDataSet = [] with open(FileToImport) as f: for line in f: StrippedLine = line.rstrip() + ",0,0,0,0,0,0,0,0,0,0,0,0,1000000,0,1000000,0,100,-100,0,100,-100,0,0,0,0,0" RawDataSet.append(StrippedLine.split(",")) LineCount = len(RawDataSet) ProcessLock = mp.Lock() DataSet = mp.Array('c_wchar_p', np.array(RawDataSet, dtype=str).flatten(), lock=ProcessLock) del RawDataSet with mp.Pool() as Pool: Args = [(ProcessLock, FileNumber, DataSet, i, LineCount) for i in range(LineCount)] Pool.starmap_async(ComputeResult, Args) if __name__ == "__main__": print("Script Started") AnalysisSettings = AnalysisSettingsFunction(-1) NumberofFiles = AnalysisSettings[0] for FileNumber in range(NumberofFiles): AnalysisSettings = AnalysisSettingsFunction(FileNumber) main(FileNumber) print("Script Finished")
原因及解决办法
1. 异步调用未等待任务完成
Pool.starmap_async是异步方法,调用后立即返回AsyncResult对象,不会阻塞主进程。你的main函数在调用该方法后,with mp.Pool()代码块随即结束,Pool会被关闭,但主进程可能快速进入下一轮循环或直接终止,导致子进程还未启动就被销毁。
解决: 调用AsyncResult.get()等待所有任务执行完毕:
with mp.Pool() as Pool: Args = [(ProcessLock, FileNumber, DataSet, i, LineCount) for i in range(LineCount)] result = Pool.starmap_async(ComputeResult, Args) result.get() # 阻塞至所有任务完成
2. starmap_async参数格式不匹配
starmap_async会将每个元组参数拆包后传递给目标函数,但你的ComputeResult只接受一个元组参数,这会导致参数不匹配错误。但因为是异步且未等待,你看不到报错信息。
解决二选一:
- 改用
map_async,它会把每个元组作为单个参数传入函数:with mp.Pool() as Pool: Args = [(ProcessLock, FileNumber, DataSet, i, LineCount) for i in range(LineCount)] result = Pool.map_async(ComputeResult, Args) result.get() - 修改
ComputeResult为接受5个单独参数:def ComputeResult(ProcessLock, FileNumber, DataSet, RowNumber, LineCount): print(f"Computing result for {RowNumber}", flush=True) # 后续逻辑不变
3. 共享数组的内存访问问题
mp.Array('c_wchar_p')存储的是字符串指针,多进程环境下子进程无法访问主进程的字符串内存空间,会导致无效指针错误,同样因为异步未等待而无法察觉。
解决: 使用multiprocessing.Manager创建跨进程可共享的列表:
from multiprocessing import Manager def main(FileNumber): # ... 其他代码 ... with Manager() as manager: DataSet = manager.list(RawDataSet) ProcessLock = mp.Lock() with mp.Pool() as Pool: Args = [(ProcessLock, FileNumber, DataSet, i, LineCount) for i in range(LineCount)] result = Pool.map_async(ComputeResult, Args) result.get()
4. 字符串与数值类型混淆
从文件读取的是字符串,直接用Value1 > Value3是字符串比较(按ASCII码顺序),不是数值比较,会导致逻辑错误。赋值时也需要转回字符串类型。
解决: 转换为数值类型后再计算:
Value1 = float(RowData[0]) Value2 = float(RowData[1]) Value3 = float(OtherRow[0]) Value4 = float(OtherRow[1]) # 赋值时转回字符串 with ProcessLock: if Value1 > Value3: calc_result = Value1 * Value2 RowData[12] = str(calc_result) OtherRow[12] = str(calc_result) else: calc_result = Value3 * Value4 RowData[12] = str(calc_result) OtherRow[12] = str(calc_result)
内容的提问来源于stack exchange,提问作者Paul

