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

使用Python/Pandas并行化大CSV文件的数据摄入与转换

超大CSV文件并行处理与PySpark适用性问题

问题描述

我当前用串行代码处理S3上的50GB超大CSV文件(含50万条记录、3万列),按1000条记录为块读取,转置成键值对格式,最终需导入Redshift。想请教:

  1. 能否并行处理这些数据块?
  2. 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天生为大数据分布式处理设计,完美匹配你的场景,核心优势:

  1. 自动并行化:无需手动拆分数据块,Spark会自动将S3上的大文件分成多个分区,分布式并行处理。
  2. S3与Redshift原生集成:直接读取S3 CSV,无需手动调用S3 API;支持通过JDBC或Redshift专用连接器直接写入目标库,跳过本地临时文件环节。
  3. 高效宽表转长表:通过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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 18:35:13