PySpark读取Avro为何触发两次Job?如何优化?
问题分析:PySpark读取Avro触发重复Job的原因与优化方案
问题描述
我通过PySpark基于Blob名称列表读取1.3TB的Avro数据时,发现系统触发了两个任务(Job Id 0含11252个任务,耗时约31分钟;Job Id 1含11452个任务,耗时约15分钟),且两个Job的描述完全一致,疑似重复读取数据。怀疑是读取时未指定Schema导致的问题,想知道具体原因以及如何优化实现仅读取一次数据。相关代码如下:
blob_names: List[str] = ... blobs_filtered: List[str] = ... ( spark .read .format('avro') .load(blobs_filtered, infer_schema=True, header=True) .select(*["qname", "user_ip", "config_id"]) .where("qname != '' and qname is not NULL") .withColumn("date", F.lit(config.process_date)) .repartition(200) .write .format("delta") .partitionBy("date") .mode("overwrite") .save(config.output_table_path) )
原因分析
- Schema推断引发重复扫描:你设置了
infer_schema=True,PySpark为了推断Avro的Schema,会先启动一个Job扫描所有指定的Avro文件,解析文件元数据获取Schema信息——这就是第一个Job的来源。之后执行实际的数据读取和处理逻辑时,又会启动第二个Job再次扫描所有文件读取数据,因此出现了两次重复读取操作。 - 无效参数干扰:
header=True参数对Avro格式完全无效,Avro本身不依赖表头定义结构,这个参数可以直接移除。
优化方案
1. 提前指定Schema(最优方案)
预先定义好Avro数据的Schema,让Spark直接使用指定Schema解析数据,避免自动推断带来的额外扫描。示例代码如下:
from pyspark.sql.types import StructType, StructField, StringType # 根据实际数据结构定义对应Schema avro_schema = StructType([ StructField("qname", StringType(), nullable=True), StructField("user_ip", StringType(), nullable=True), StructField("config_id", StringType(), nullable=True) # 若有其他字段,按需补充定义 ]) ( spark .read .format('avro') .schema(avro_schema) # 使用预定义Schema .load(blobs_filtered) .select(*["qname", "user_ip", "config_id"]) .where("qname != '' and qname is not NULL") .withColumn("date", F.lit(config.process_date)) .repartition(200) .write .format("delta") .partitionBy("date") .mode("overwrite") .save(config.output_table_path) )
这样Spark只会执行一次数据扫描,不会额外触发Schema推断的Job,彻底解决重复读取问题。
2. 缓存中间结果(备选方案)
如果无法提前确定Schema,可以在读取数据后缓存DataFrame,避免后续处理重复扫描源文件:
df = ( spark .read .format('avro') .load(blobs_filtered, infer_schema=True) .cache() # 缓存读取后的DataFrame ) ( df .select(*["qname", "user_ip", "config_id"]) .where("qname != '' and qname is not NULL") .withColumn("date", F.lit(config.process_date)) .repartition(200) .write .format("delta") .partitionBy("date") .mode("overwrite") .save(config.output_table_path) ) df.unpersist() # 处理完成后释放缓存,避免占用资源
注意:这种方法仍会执行一次Schema推断的Job,只是后续处理复用缓存数据,效率略低于提前指定Schema。
3. 额外优化建议
- 调整分区数:1.3TB数据设置200个分区,每个分区约6.5GB,可能因单分区数据量过大导致任务执行缓慢。建议根据集群资源调整,将单分区数据量控制在1-2GB左右,提升并行处理效率。
- 合并小文件:如果Blob存储中的Avro文件过小(远小于1GB),可以先合并小文件,减少Task数量,降低调度开销。
内容的提问来源于stack exchange,提问作者Dariusz Krynicki
相关产品推荐
相关产品推荐

