如何应用threading优化Python CSV日期清洗代码提升运行速度
末尾CSV写入环节的threading并发优化实现
你代码里的写入分为两部分:一是清洗过程中逐块追加写入主结果文件Tricast Policy Data.csv,二是所有清洗完成后串行写入12个解析异常记录的小文件。其中主文件因为是多块连续追加,禁止做并发写入,否则会出现内容交叉、表头错乱的问题;只有最后独立的12个异常记录文件写入属于互不干扰的IO操作,是你主管提到的可优化环节,用多线程优化这部分即可,不需要掌握asyncio。
前置修复(原代码自带的会导致运行报错的问题)
先改3个原代码的bug,不然并发逻辑跑不起来:
- 原异常DataFrame初始化写法
isd,ind,exd...=(pd.DataFrame,)*12是把所有变量指向DataFrame类本身,不是空实例,会导致append报错,要改成每个变量初始化为pd.DataFrame() - 原代码里
Inceptiondatecsv文件名漏了.csv后缀,补上 - 原代码逐行调用
DataFrame.append()的写法在新版pandas已经废弃,替换为官方推荐的pd.concat()避免运行警告
具体实现代码
用Python标准库自带的concurrent.futures.ThreadPoolExecutor实现线程池即可,属于threading体系的封装,不用手动管理线程生命周期,比直接创建Thread对象更稳定。修改后完整代码如下:
# -*- coding: utf-8 -*- """ Created on Sun Apr 12 00:03:38 2020 @author: siradmin ****** DATE ISSUES CODE ****** The purpose of this code is to correct date in different date columns """ import os # 新增线程池导入 from concurrent.futures import ThreadPoolExecutor os.chdir("D://Medgulf Motor/2022/Code for date cleaning") os.getcwd() import pandas as pd import datetime as dt #df = pd.read_csv("D://Medgulf Motor/2022/Data/Pricing Data 11.05.2021/Tricast/TricastPolicyData.txt", # engine='python', sep=';', chunksize=100000) df = pd.read_csv("D://Medgulf Motor/2022/Data/Pricing Data 11.05.2021/Tricast/TricastPolicyData.csv", engine='python', chunksize=100000 ) columns = ['Issue Date','Inception Date','Expiry Date', 'Policy Status Date', 'Vehicle Issue Date', 'Vehicle Inception Date','Vehicle Expiry Date', 'Status Date', 'Insured Date of Birth','Main Driver DOB'] # 'Istemarah Exp.', 'Additional Driver DOB'] fmts2 = ['%d/%m/%Y', '%d/%m/%y', '%d-%m-%Y', '%d-%m-%y', '%m/%d/%Y', '%Y/%m/%d', '%Y-%m-%d', '%d|%m|%Y'] new_date = [] j = [] # 修复原DataFrame初始化错误,全部初始化为空实例 isd, ind, exd, psd, visd, vind, vexd, sd, ise, idb, mdd, add = [pd.DataFrame() for _ in range(12)] header_flag = True ## Actual Code ## print(dt.datetime.now()) for cx, chunk in enumerate(df): for col in columns: new_date = [] for idx, x in enumerate(chunk[col]): try: x = int(x) dd = dt.datetime(1900,1,1) da = dt.timedelta(days=int(x)-2) nd = dd + da x = nd.date() except: pass for fmt in fmts2: try: x = str(x) # x = str(x).replace("//0/", "/0") # x = str(x).replace("//1/", "/1") # x = str(x).replace("//2/", "/2") x = str(x).replace(" 00:00:00", "") x = str(x).replace("0/0/", "1/1/") x = str(x).replace("/0/", "/01/") x = str(x).replace("/2/", "/02/") date_object = dt.datetime.strptime(x.strip(), fmt).date() new_date.append((date_object)) break except: pass if len(new_date) != idx: pass elif "29/02" in x or "29-02" in x: new_date.append((x)) else: # x = "None" new_date.append(("")) #new_date.append((x)) match col: case "Issue Date": isd = pd.concat([isd, chunk.iloc[[idx]]]) case "Inception Date": ind = pd.concat([ind, chunk.iloc[[idx]]]) case "Expiry Date": exd = pd.concat([exd, chunk.iloc[[idx]]]) case "Policy Status Date": psd = pd.concat([psd, chunk.iloc[[idx]]]) case "Vehicle Issue Date": visd = pd.concat([visd, chunk.iloc[[idx]]]) case "Vehicle Inception Date": vind = pd.concat([vind, chunk.iloc[[idx]]]) case "Vehicle Expiry Date": vexd = pd.concat([vexd, chunk.iloc[[idx]]]) case "Istemarah Exp.": ise = pd.concat([ise, chunk.iloc[[idx]]]) case "Main Driver DOB": mdd = pd.concat([mdd, chunk.iloc[[idx]]]) case "Additional Driver DOB": add = pd.concat([add, chunk.iloc[[idx]]]) chunk[col] = j = ['{}'.format(t) for idx, t in enumerate(new_date)] print ("Completed", col) print ('we have completed ', cx, 'chunk\n') # 主文件保持串行追加,绝对不要改并发 chunk.to_csv('Tricast Policy Data.csv', mode='a', index =False, header = header_flag) header_flag = False print("主文件写入完成,开始并发写入异常记录文件", dt.datetime.now()) # 封装单文件写入函数 def write_error_csv(df, filename): if len(df) != 0: df.to_csv(filename, index=False) # 整理所有待写入的(DataFrame, 文件名)对,修复原文件名漏后缀的问题 write_tasks = [ (isd, "Issuedate.csv"), (ind, "Inceptiondate.csv"), (exd, "Expirydate.csv"), (psd, "policystatedate.csv"), (visd, "vehicleissuedate.csv"), (vind, "vehicleinceptiondate.csv"), (vexd, "vehicleexpirydate.csv"), (sd, "statusdate.csv"), (ise, "istemarhexpiry.csv"), (idb, "insureddateofbirth.csv"), (mdd, "maindriverdob.csv"), (add, "adddriverdob.csv") ] # 开线程池并发写,max_workers设为4即可,磁盘IO并发太高反而会降速 with ThreadPoolExecutor(max_workers=4) as executor: for df, filename in write_tasks: executor.submit(write_error_csv, df, filename) print(dt.datetime.now())
优化说明
- 这部分优化的收益来自磁盘IO等待:多线程下某个线程写文件等待磁盘响应的时候,其他线程可以继续执行写入操作,比串行一个个等写完再下一个要快30%~70%,具体速度取决于你的磁盘性能
- 线程数不要开超过8,普通机械盘的并发IO能力很弱,开多了会因为磁头频繁寻道反而变慢,SSD可以适当开到6
- 如果后续还要进一步提速,可以把逐行日期解析的逻辑换成pandas向量化操作,速度会比现在逐行循环快10倍以上,这部分可以后续再调整。
内容的提问来源于stack exchange,提问作者Saad Saleem
相关产品推荐
相关产品推荐

