PySpark读取仅4行csv小文件时spark.read.csv耗时过长如何解决?
PySpark读取极小CSV文件耗时过高解决方案
Spark作为分布式计算引擎,本身存在任务调度、执行计划生成的固定开销,该开销针对TB级大文件可忽略,但对于仅4行的极小文件会成为耗时主体,你观测到的3~5秒耗时绝大多数属于该固定开销,而非文件读取本身的耗时。
优化建议
- 关闭Schema自动推断,手动指定Schema
默认配置下Spark读取CSV会默认开启inferSchema参数,需要扫描全量文件推断字段类型,额外增加耗时。手动提前定义Schema后可跳过该流程,代码示例如下:
from pyspark.sql.types import StructType, StringType, IntegerType # 按照实际CSV的字段顺序、类型自定义Schema custom_schema = StructType() \ .add("col1", StringType(), nullable=True) \ .add("col2", IntegerType(), nullable=True) # 按实际业务补充剩余字段定义 step2 = datetime.now() df = spark.read.options(header=True, delimiter="|", inferSchema=False)\ .schema(custom_schema)\ .csv(csv_data[0]) print(F"\nStep-2 | {(datetime.now() - step2).total_seconds()}\n")
- 极小文件优先用Pandas读取后转Spark DataFrame
如果业务场景固定读取行数极少的小文件,可跳过Spark的文件读取流程,直接用单机IO开销更低的Pandas读取后转换为Spark DataFrame,代码示例如下:
import pandas as pd step2 = datetime.now() pd_df = pd.read_csv(csv_data[0], sep="|", header=0) df = spark.createDataFrame(pd_df) print(F"\nStep-2 | {(datetime.now() - step2).total_seconds()}\n")
优化本地Spark配置降低调度开销
如果是本地测试场景,可调整Spark配置减少不必要的资源调度损耗:- 启动Spark时设置
spark.master = local[1],仅用单核心运行避免多线程调度开销 - 调低Spark日志级别到WARN或ERROR,减少日志打印的IO开销
- 启动Spark时设置
重复使用的小文件读取后持久化缓存
如果同一份小文件需要多次计算使用,读取完成后调用df.cache()将其持久化到内存,后续使用无需重复走文件读取、Schema校验流程,直接从内存读取。
内容的提问来源于stack exchange,提问作者Ankit Tiwari
相关产品推荐
相关产品推荐

