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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 02:10:19