Python处理大型CSV导出:如何分块写入及查看进度
解决方案:大数据CSV处理的进度监控与分块写入
问题背景
需处理含140万行数据的CSV文件,计算后导出新CSV。测试100行数据时输出达2400行,但处理大数据时无法知晓进度及运行状态,需实现进度可视化和分块写入CSV的功能。
实现方案
1. 进度监控
使用tqdm库为数据处理的关键步骤添加进度条,实时显示处理进度、剩余时间等信息。先安装依赖:
pip install tqdm
2. 分块写入CSV
针对最终生成的大体积tech数据,避免一次性写入导致内存过载,采用分块迭代写入的方式,通过mode='a'追加写入,并控制仅在第一次写入时添加表头。
修改后的完整代码
import pandas as pd import numpy as np import csv from tqdm import tqdm # 开启tqdm对pandas的支持 tqdm.pandas() # ---------------------- 数据读取与预处理 ---------------------- print("开始读取并预处理数据...") patents = pd.read_csv('E:/Work/final.csv') # 执行数据清洗与筛选 patents = patents[['publn_nr', 'permno_adj', 'publn_year', 'IPC1']].dropna().drop_duplicates().reset_index(drop=True) patents = patents[(patents['publn_year'] >= 1980) & (patents['publn_year'] < 2016)].reset_index(drop=True) patents['permno_adj'] = patents['permno_adj'].astype(str) + patents['publn_year'].astype(str) print("数据预处理完成,前5行数据:") print(patents.head()) # ---------------------- 数据聚合计算 ---------------------- print("开始分组聚合计算...") patents = patents.groupby(['permno_adj', 'IPC1'])['publn_nr'].nunique().reset_index() patents.columns = ['permno_adj', 'IPC1', 'ipc_patents'] patents['total_patents'] = patents.groupby(['permno_adj'])['ipc_patents'].transform('sum') patents['share'] = patents['ipc_patents'] / patents['total_patents'] # ---------------------- 核心计算与进度监控 ---------------------- for v in ['IPC1']: temp = patents.copy() print(f"开始处理{v}维度的 pivot 转换...") T = temp.pivot(index=f'{v}', columns='permno_adj', values='share').fillna(0) X_t = temp.pivot(index='permno_adj', columns=f'{v}', values='share').fillna(0) print("开始标准化T矩阵...") T_t = T.copy() # 为列循环添加进度条 for column in tqdm(list(T_t), desc="标准化列"): norm_val = np.sqrt(np.dot(T_t[[column]].values.T, T_t[[column]].values)[0][0]) T_t[column] = T_t[column] / norm_val print("计算相似度矩阵...") om_f = X_t.T.dot(X_t) om_s = X_t.T.dot(X_t) # 为双重循环添加进度条 sic_list = list(om_s) for sic1 in tqdm(sic_list, desc="处理相似度矩阵行"): for sic2 in sic_list: om_s.loc[sic1][sic2] = om_s.loc[sic1][sic2] / (np.sqrt(om_f[sic1][sic1]) * np.sqrt(om_f[sic2][sic2])) print("计算技术相似度...") tech = T_t.T.dot(om_s).dot(T_t) tech = tech.unstack().reset_index(level=1) if 'IPC1' in v: tech.columns = ['permno_adj_pat', 'tech_mahal_sim'] tech = tech.reset_index() tech = tech[(tech['permno_adj'] != tech['permno_adj_pat'])].sort_values( ['permno_adj', 'permno_adj_pat']).reset_index(drop=True) # ---------------------- 分块写入CSV ---------------------- print("开始分块写入CSV文件...") chunk_size = 100000 # 每块写入10万行,可根据内存调整 total_chunks = len(tech) // chunk_size + 1 output_file = 'gajuf.csv' # 先清空目标文件(若已存在) with open(output_file, 'w', newline='', encoding='utf-8') as f: pass for i in tqdm(range(total_chunks), desc="写入CSV块"): start = i * chunk_size end = min((i+1)*chunk_size, len(tech)) chunk = tech.iloc[start:end] # 仅第一次写入时添加表头 chunk.to_csv(output_file, mode='a', index=False, header=(i==0), encoding='utf-8') # 写入stata文件 tech.to_stata('gajuf.dta', write_index=False) print("处理完成!")
关键改进说明
- 进度监控:通过
tqdm为列循环、相似度矩阵循环、CSV写入循环添加进度条,实时显示处理进度;关键步骤添加打印提示,明确当前执行阶段。 - 分块写入:将
tech数据按指定大小拆分,逐块追加写入CSV,避免一次性写入大文件导致的内存压力,同时保证表头仅写入一次。 - 内存优化(可选):若原数据读取时内存不足,可改用分块读取模式替换原
pd.read_csv:chunks = [] for chunk in tqdm(pd.read_csv('E:/Work/final.csv', chunksize=100000), desc="读取CSV块"): processed_chunk = chunk[['publn_nr', 'permno_adj', 'publn_year', 'IPC1']].dropna().drop_duplicates() processed_chunk = processed_chunk[(processed_chunk['publn_year'] >= 1980) & (processed_chunk['publn_year'] < 2016)] processed_chunk['permno_adj'] = processed_chunk['permno_adj'].astype(str) + processed_chunk['publn_year'].astype(str) chunks.append(processed_chunk) patents = pd.concat(chunks).reset_index(drop=True)
内容的提问来源于stack exchange,提问作者Gaju_masare
相关产品推荐
相关产品推荐

