使用Python/Pandas并行化大CSV文件的数据摄入与转换
超大CSV文件并行处理与PySpark适用性问题
问题描述
我当前用串行代码处理S3上的50GB超大CSV文件(含50万条记录、3万列),按1000条记录为块读取,转置成键值对格式,最终需导入Redshift。想请教:
- 能否并行处理这些数据块?
- PySpark是否适用于这个场景?
当前串行代码
import pandas as pd import numpy as np import time import boto3 s3 = boto3.client('s3') obj = s3.get_object(Bucket='bucket', Key = 'path') csv_header = pd.read_csv(obj['Body'], nrows=1).columns obj = s3.get_object(Bucket='bucket', Key = 'path') csv_iterator = pd.read_csv(obj['Body'], iterator=True, chunksize=1000,low_memory=False) i=0 c=0 print("Iteration Start!",time.strftime("%H:%M:%S", time.localtime())) for csv_chunk in csv_iterator: column_size = csv_header.size-csv_chunk.columns.size csv_chunk[np.arange(column_size)]=None chunk = pd.DataFrame(csv_chunk.values,columns=csv_header) if i==0: chunk.to_csv('temp_source.csv',mode = 'a') else: chunk.to_csv('temp_source.csv',mode = 'a',header=False) out = pd.melt(chunk, id_vars=['eid'], value_vars=chunk.columns[1:]) if i==0: out.to_csv('temp_key_value.csv',mode='a') else: out.to_csv('temp_key_value.csv',mode='a',header=False) i=i+1 c=c + csv_chunk.shape[0] print("Iteration ",i," Chunk - ",c," - ",time.strftime("%H:%M:%S", time.localtime())) print("Load Complete!",time.strftime("%H:%M:%S", time.localtime()))
问题解答
一、可以并行处理数据块
完全可以将现有串行逻辑改成并行执行,核心注意点如下:
- 避免写入冲突:多个进程同时写入同一个CSV文件会导致数据错乱,建议给每个进程分配独立临时文件,最后统一合并;或者用"任务处理+单进程写入队列"的模式。
- 控制S3并发度:S3对单个对象的并发读取有一定限制,建议进程数控制在20-50之间,避免触发限流。
- 内存阈值把控:每个进程处理1000条3万列的数据,需确保单进程内存足够,避免OOM(可适当调小chunksize)。
Python生态中推荐用concurrent.futures.ProcessPoolExecutor实现多进程并行,因为数据转置、melt属于CPU密集型操作,多进程比多线程效率更高。
二、PySpark非常适用,且是更优方案
PySpark天生为大数据分布式处理设计,完美匹配你的场景,核心优势:
- 自动并行化:无需手动拆分数据块,Spark会自动将S3上的大文件分成多个分区,分布式并行处理。
- S3与Redshift原生集成:直接读取S3 CSV,无需手动调用S3 API;支持通过JDBC或Redshift专用连接器直接写入目标库,跳过本地临时文件环节。
- 高效宽表转长表:通过
stack函数实现类似pandasmelt的功能,处理3万列的宽表性能远高于单机pandas。
简化PySpark示例代码
from pyspark.sql import SparkSession from pyspark.sql.functions import expr # 初始化Spark会话 spark = SparkSession.builder.appName("CSVToKeyValue").getOrCreate() # 读取S3上的CSV,自动处理列数不一致的行 df = spark.read.csv( "s3://bucket/path", header=True, inferSchema=False, # 3万列建议关闭自动推断,统一用string类型 mode="PERMISSIVE", # 列数不匹配的行自动补null nullValue="" ) # 生成宽表转长表的stack表达式 non_id_cols = [col for col in df.columns if col != "eid"] stack_expr = f"stack({len(non_id_cols)}, {', '.join([f'{repr(col)}, `{col}`' for col in non_id_cols])}) as (variable, value)" key_value_df = df.select("eid", expr(stack_expr)) # 写入Redshift key_value_df.write.format("jdbc").options( url="jdbc:redshift://your-redshift-endpoint:5439/your-db", driver="com.amazon.redshift.jdbc.Driver", dbtable="your_target_table", user="your-username", password="your-password" ).mode("append").save() spark.stop()
关键注意事项
- Schema优化:3万列手动指定schema不现实,建议先读取表头,统一将列类型设为string,后续按需转换。
- 并行度调整:根据集群资源修改
spark.sql.shuffle.partitions(默认200),本地测试可设置master="local[*]"利用全部CPU核心。 - 写入性能优化:优先用Redshift的
COPY命令结合S3临时存储,比JDBC写入快数倍,PySpark可通过redshift格式实现。
三、两种方案对比
| 方案 | 优点 | 缺点 |
|---|---|---|
| Python多进程 | 代码改动小,无需额外集群 | 单机器性能上限低,内存压力大 |
| PySpark分布式 | 大数据处理效率高,可扩展 | 需要Spark集群环境,有一定学习成本 |
内容的提问来源于stack exchange,提问作者Will Graham
相关产品推荐
相关产品推荐

