PySpark加载150GB多行CSV过慢,求Azure环境下优化方案
解决Spark加载超大多行CSV文件的性能瓶颈问题
问题本质
当开启multiLine=True读取CSV时,Spark无法按字节拆分文件——因为要识别引号包裹的多行内容,必须保证每个分区能完整读取一条或多条完整记录,导致大文件只能由单个Executor处理,这就是150GB文件耗时数小时的核心原因。
针对Azure Synapse+ADLS Gen2的解决方案
1. 预处理拆分大文件(批量场景首选)
利用ADLS Gen2的文件访问能力,先将大文件拆分成多个可并行处理的小文件,同时保证不破坏多行记录的完整性。以下是一个PySpark实现的批量拆分函数:
from pyspark.sql import SparkSession import re def split_large_multiline_csv(adls_path, output_path, chunk_size_mb=100): # 转换块大小为字节 chunk_size = chunk_size_mb * 1024 * 1024 # 读取文件为二进制RDD,按块拆分 file_rdd = spark.sparkContext.binaryFiles(adls_path).flatMap(lambda x: [x[1]]) # 处理每个块,确保拆分点不在引号内 def process_chunk(chunk): chunk_str = chunk.decode('utf-8') # 匹配所有引号位置 quote_positions = [m.start() for m in re.finditer('"', chunk_str)] # 找到最后一个偶数位置的引号(确保拆分在完整记录后) split_pos = chunk_size if quote_positions: for pos in reversed(quote_positions): if pos < chunk_size and len(quote_positions[:quote_positions.index(pos)+1]) % 2 == 0: split_pos = pos + 1 break # 拆分块并返回 return [chunk_str[:split_pos], chunk_str[split_pos:]] if split_pos < len(chunk_str) else [chunk_str] processed_rdd = file_rdd.flatMap(process_chunk) # 将拆分后的内容写入临时目录 processed_rdd.map(lambda x: x.encode('utf-8')).saveAsTextFile(output_path) # 从临时目录读取为DataFrame df = spark.read.format('csv')\ .option('header', True)\ .option('multiLine', True)\ .load(output_path) return df
2. 优化Spark执行配置
针对Azure Synapse Spark集群,调整以下参数提升单Executor处理效率:
- 增大Executor内存:
--executor-memory 32G(根据集群节点规格调整) - 调整内存Overhead:
spark.executor.memoryOverhead 8192(避免OOM) - 启用动态资源分配:
spark.dynamicAllocation.enabled true
3. 批量文件处理优化
对于数千个混合大小的文件,建议:
- 先筛选出大文件(比如>10GB)用上述拆分函数处理
- 小文件直接用常规方式加载后合并
- 利用ADLS Gen2的
glob路径匹配(如path/to/files/*.csv)批量加载
注意事项
- 拆分时的
chunk_size_mb可根据集群节点内存调整,建议100-200MB - 确保CSV文件的编码统一(上述函数默认UTF-8)
- 拆分后的临时文件可在任务完成后删除,避免占用存储
内容的提问来源于stack exchange,提问作者rocket porg
相关产品推荐
相关产品推荐

