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

如何应用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.28 14:45:28