并行处理下日志文件格式异常问题求助
问题分析与解决方案
问题场景
使用ThreadPoolExecutor并行执行脚本时,日志文件出现截断行、行数不增等异常;改用ProcessPoolExecutor后流程无法正常执行,问题也未解决。
核心代码片段
线程执行逻辑:
while(conf['loop']): for (_,_,whole_tit) in os.walk(dir_base+conf['dirr']): break # with cf.ProcessPoolExecutor(max_workers=cores) as executor: with cf.ThreadPoolExecutor(max_workers=cores) as executor: executor.map(parallel_master_routine, whole_tit)
日志处理相关代码:
# parallel_master_routine内的日志写入逻辑 an = dir_base+conf['dir_conf']+'Analisi_models'+'.csv' try: if not os.path.isfile(an): open(an,'w').close() aggiorna_tit(an,dir_base+conf['dir_conf']+'_Analisi_model_'+'.csv',tit,net) except: pass with open(an,'a') as f: dt = datetime.datetime.now() f.write(str(dt)+'§'+tit+'§'+net+'§'+'pred/orig'+'§'+",".join([str(i) for i in data_predict])+'§'+",".join([str(i) for i in y_test])+'§'+str(pct_nul)+"\n") # 去重日志的函数 def aggiorna_tit(t_in,t_out,tit,net): if not os.path.exists(t_in): open(t_in,'a').close() with open(t_in, "r") as input: input = input.readlines() ex = [e for e in input if not '§'+str(tit)+'§'+str(net)+'§' in e ] x = "".join(ex) with open(t_out, "w") as output: output.write(x) # replace file with original name os.replace(t_out, t_in)
问题根源
- 线程安全问题:多个线程同时读写同一日志文件时,
write操作会出现数据交错、截断——比如A线程写了一半,B线程的写入插入进来,导致行格式损坏。 - 竞态条件:多个线程同时执行
aggiorna_tit时,会出现线程A读取文件后,线程B先完成写入替换,线程A再写入自己处理后的内容,直接覆盖掉线程B的修改,导致日志丢失、行数不增。 - 进程池适配问题:改用进程池后,进程间内存不共享,若
parallel_master_routine或依赖函数存在未序列化的对象、全局变量依赖,会导致流程执行异常。
解决方案
方案1:给文件读写加线程锁
为日志文件的读写操作添加全局锁,确保同一时间只有一个线程操作文件:
# 模块级别初始化全局锁 import threading log_lock = threading.Lock() # 修改parallel_master_routine内的日志逻辑 an = dir_base+conf['dir_conf']+'Analisi_models'+'.csv' try: if not os.path.isfile(an): open(an,'w').close() # 给去重操作加锁 with log_lock: aggiorna_tit(an,dir_base+conf['dir_conf']+'_Analisi_model_'+'.csv',tit,net) except: pass # 给追加写入加锁 with log_lock: with open(an,'a') as f: dt = datetime.datetime.now() f.write(str(dt)+'§'+tit+'§'+net+'§'+'pred/orig'+'§'+",".join([str(i) for i in data_predict])+'§'+",".join([str(i) for i in y_test])+'§'+str(pct_nul)+"\n")
该方案适配线程池场景,可彻底解决线程间的文件读写竞态问题。
方案2:改用进程安全的文件操作(针对ProcessPoolExecutor)
若坚持使用进程池,需做以下调整:
- 确保所有传递给
executor.map的参数可序列化(pickle兼容); - 使用进程锁(
multiprocessing.Lock),且需通过参数传递给子进程:
import multiprocessing as mp # 初始化进程锁 log_lock = mp.Lock() # 修改parallel_master_routine,新增lock参数 def parallel_master_routine(tit, lock): # ... 原有计算逻辑 ... an = dir_base+conf['dir_conf']+'Analisi_models'+'.csv' try: if not os.path.isfile(an): open(an,'w').close() with lock: aggiorna_tit(an,dir_base+conf['dir_conf']+'_Analisi_model_'+'.csv',tit,net) except: pass with lock: with open(an,'a') as f: dt = datetime.datetime.now() f.write(str(dt)+'§'+tit+'§'+net+'§'+'pred/orig'+'§'+",".join([str(i) for i in data_predict])+'§'+",".join([str(i) for i in y_test])+'§'+str(pct_nul)+"\n") # 调用时传递锁 with cf.ProcessPoolExecutor(max_workers=cores) as executor: executor.map(parallel_master_routine, whole_tit, [log_lock]*len(whole_tit))
同时需检查parallel_master_routine内的所有依赖对象,确保均可被pickle序列化,避免进程池执行异常。
方案3:异步日志队列(更优雅的方式)
创建单独的日志线程,所有线程/进程将日志消息发送到队列,由该线程统一处理文件写入:
import queue import threading # 初始化日志队列 log_queue = queue.Queue() # 日志处理线程函数 def log_worker(log_file_path): while True: msg = log_queue.get() if msg is None: # 收到终止信号 break action, data = msg with open(log_file_path, 'r+') as f: if action == 'update': tit, net = data lines = f.readlines() lines = [e for e in lines if not f'§{tit}§{net}§' in e] f.seek(0) f.truncate() f.writelines(lines) elif action == 'append': dt, tit, net, data_predict, y_test, pct_nul = data line = f"{dt}§{tit}§{net}§pred/orig§{','.join(map(str, data_predict))}§{','.join(map(str, y_test))}§{pct_nul}\n" f.write(line) log_queue.task_done() # 启动日志线程 an = dir_base+conf['dir_conf']+'Analisi_models'+'.csv' if not os.path.isfile(an): open(an,'w').close() log_thread = threading.Thread(target=log_worker, args=(an,)) log_thread.daemon = True log_thread.start() # 修改parallel_master_routine内的逻辑 def parallel_master_routine(tit): # ... 原有计算逻辑 ... try: # 发送去重请求到队列 log_queue.put(('update', (tit, net))) except: pass # 发送追加请求到队列 dt = datetime.datetime.now() log_queue.put(('append', (dt, tit, net, data_predict, y_test, pct_nul))) # 线程池调用不变 with cf.ThreadPoolExecutor(max_workers=cores) as executor: executor.map(parallel_master_routine, whole_tit) # 程序结束时发送终止信号 log_queue.put(None) log_thread.join()
该方案将所有文件操作集中到单个线程,从根源上避免竞态问题,同时保持并行计算的效率。
内容的提问来源于stack exchange,提问作者Simon
相关产品推荐
相关产品推荐

