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

并行处理下日志文件格式异常问题求助

问题分析与解决方案

问题场景

使用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)

问题根源

  1. 线程安全问题:多个线程同时读写同一日志文件时,write操作会出现数据交错、截断——比如A线程写了一半,B线程的写入插入进来,导致行格式损坏。
  2. 竞态条件:多个线程同时执行aggiorna_tit时,会出现线程A读取文件后,线程B先完成写入替换,线程A再写入自己处理后的内容,直接覆盖掉线程B的修改,导致日志丢失、行数不增。
  3. 进程池适配问题:改用进程池后,进程间内存不共享,若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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 16:25:24