如何统计多个tar.gz文件中各CSV列的空值数量?
高效处理10000+压缩CSV的一致性校验方案
针对你手里10000+个tar.gz压缩的气象CSV文件,要高效完成列一致性校验、空值统计,下面是经过优化的实操方案:
1. 先锁定基准列集合
从你提供的master_file里提取所有标准列名,作为后续所有文件的对齐基准——这是解决列顺序不一致问题的核心,确保每个CSV最终都按这个列集合输出统计结果。
2. 直接读取压缩包内的CSV,跳过磁盘解压
没必要把所有tar.gz都解压到磁盘,用tarfile直接读取压缩包内的CSV文件,能大幅节省磁盘IO和存储空间:
import tarfile import pandas as pd def load_csv_from_tar(tar_path, master_columns): with tarfile.open(tar_path, 'r:gz') as tar: # 假设每个tar.gz仅包含一个CSV文件 csv_file = next(member for member in tar.getmembers() if member.name.endswith('.csv')) with tar.extractfile(csv_file) as f: # 只读取基准列内的字段,过滤多余列 df_chunk = pd.read_csv(f, usecols=lambda col: col in master_columns, dtype=str) # 对齐基准列顺序,缺失列自动补NaN return df_chunk.reindex(columns=master_columns)
3. 分块读取超大CSV,控制内存占用
单个CSV200-600MB,直接全量读入内存容易撑爆,用pandas的chunksize分块读取,逐块统计空值后合并结果:
def calculate_null_counts(tar_path, master_columns, chunksize=100000): # 初始化空值统计容器 null_stats = pd.Series(0, index=master_columns) extra_cols = set() with tarfile.open(tar_path, 'r:gz') as tar: csv_file = next(member for member in tar.getmembers() if member.name.endswith('.csv')) with tar.extractfile(csv_file) as f: # 先读取首行获取当前文件的列名,检查是否有额外列 header_line = f.readline().decode('utf-8').strip().split(',') extra_cols = set(header_line) - set(master_columns) # 重置文件指针到开头 f.seek(0) # 分块读取并统计空值 for chunk in pd.read_csv(f, usecols=lambda col: col in master_columns, dtype=str, chunksize=chunksize): chunk = chunk.reindex(columns=master_columns) null_stats += chunk.isnull().sum() # 组装结果为DataFrame,方便后续合并 result = pd.DataFrame({ 'filename': tar_path, 'extra_columns': ','.join(extra_cols) if extra_cols else '无', **null_stats.to_dict() }, index=[0]) return result
这里额外加入了额外列检测,能帮你发现不符合基准的异常列。
4. 多进程并行处理,提速10倍+
10000+文件单线程处理太慢,用concurrent.futures.ProcessPoolExecutor做多进程并行(IO+CPU密集型任务,多进程比多线程效率更高):
from concurrent.futures import ProcessPoolExecutor import os def batch_process_files(root_dir, master_columns, max_workers=8): # 遍历目录下所有tar.gz文件 all_tar_files = [os.path.join(root_dir, fname) for fname in os.listdir(root_dir) if fname.endswith('.tar.gz')] # 并行处理,max_workers根据CPU核心数设置(建议核心数的1-2倍) with ProcessPoolExecutor(max_workers=max_workers) as executor: # 映射任务到进程池 results = list(executor.map(lambda path: calculate_null_counts(path, master_columns), all_tar_files)) # 合并所有结果为一张汇总表 summary_table = pd.concat(results, ignore_index=True) return summary_table
5. 结果保存与后续分析
处理完成后,把汇总表保存为Parquet格式(比CSV节省70%以上空间,且支持快速查询),或者CSV:
# 替换成你的文件目录和基准列列表 master_columns = ['date', 'temperature', 'humidity', 'wind_speed', ...] # 从master_file提取 summary_table = batch_process_files('/your/tar_files_directory', master_columns) # 保存为Parquet(推荐) summary_table.to_parquet('weather_data_summary.parquet', index=False) # 或者保存为CSV summary_table.to_csv('weather_data_summary.csv', index=False)
在Jupyter Notebook里可以直接加载汇总表,做进一步分析:
- 筛选空值率过高的文件
- 按年份/月份统计空值趋势
- 排查存在额外列的异常文件
6. 额外优化建议
- 指定数据类型:提前知道各列的dtype(比如日期列用
parse_dates,数值列用float32),在read_csv里指定,能进一步减少内存占用和读取时间。 - 异常捕获:在
calculate_null_counts里加入try-except块,遇到损坏的tar.gz或CSV时,记录错误信息并跳过,避免整个流程中断。 - 缓存中间结果:用
joblib缓存每个文件的处理结果,万一处理中断,不用重新跑所有文件。
内容的提问来源于stack exchange,提问作者Jason Rogers
相关产品推荐
相关产品推荐

