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

Azure Databricks处理大型ZIP文件时遭遇spark.rpc.message.maxSize超限错误的解决求助

Azure Databricks处理大型ZIP文件时遭遇spark.rpc.message.maxSize超限错误的解决求助

大家好,我现在碰到了一个棘手的问题,想请各位帮忙出出主意:

我正在Azure Databricks中处理一个345GB的ZIP压缩包,里面包含单个1.5TB的CSV文件,大概有30亿行数据。我的目标是把这个CSV转换成Delta表,这样后续数据管道里的读取速度能更快,所有文件都存储在Azure Blob存储中。

我的处理流程大致是这样的:

  • 初始化Spark Session
  • 用pd.read_csv(chunksize=CHUNKSIZE)将ZIP内的CSV文件读取为迭代器
  • 遍历迭代器中的每个数据块
  • 对每个数据块执行以下操作:
    • 将Pandas DataFrame转换为Spark DataFrame
    • 追加写入到Azure Blob存储中的Delta表
    • 清理内存

遇到的错误

一开始我设置CHUNKSIZE=5_000_000,处理到第83个块(累计4.15亿行)时,触发了如下错误:

org.apache.spark.SparkException: Job aborted due to stage failure: Serialized task 42776:450 was 282494392 bytes, which exceeds max allowed: spark.rpc.message.maxSize (268435456 bytes). Consider increasing spark.rpc.message.maxSize or using broadcast variables for large values.

之后我把块大小调整为CHUNKSIZE=3_000_000,错误还是会出现,只是延迟到了第139个块(累计4.17亿行),看起来像是处理到一定行数后就会触发这个问题:

org.apache.spark.SparkException: Job aborted due to stage failure: Serialized task 43293:250 was 282494392 bytes, which exceeds max allowed: spark.rpc.message.maxSize (268435456 bytes). Consider increasing spark.rpc.message.maxSize or using broadcast variables for large values.

额外的需求

我明白直接将文件读取为Spark DataFrame是更理想的方案,但我暂时没找到直接读取Azure Blob中ZIP文件为Spark DataFrame的方法,如果大家有这方面的建议,也请不吝赐教。

我的集群配置

  • 2个工作节点:Standard_DS3_v2,14GB内存,4核
  • 驱动节点:Standard_DS13_v2,56GB内存,8核

已尝试的解决方法

我已经试过以下几种方式,但都没能解决问题:

  • 按照错误提示,在SparkSession初始化时配置spark.rpc.message.maxSize为"512"(默认单位应该是MiB),打印配置时能看到设置已生效,但错误依然存在
  • 也尝试过设置为"536870912",担心单位是字节,但结果还是一样
  • 尝试在Spark实例初始化完成后用spark.conf.set("spark.rpc.message.maxSize", "512")修改,结果报错提示该参数无法在Spark启动后更改

完整代码

def convert_zip_to_delta(snapshot_date: str, start_chunk: int = 0):

    # File paths
    zip_file = f"{snapshot_date}.zip"
    delta_file = f"{snapshot_date}_delta"
    delta_table_path = f"wasbs://{CONTAINER_NAME}@{STORAGE_ACCOUNT_NAME}.blob.core.windows.net/{delta_file}/"

    spark = (
        SparkSession.builder.config("spark.sql.shuffle.partitions", "100")
        .config("spark.hadoop.fs.azure.retries", "10")
        .config("spark.rpc.message.maxSize", "536870912")  # 512 MiB in bytes
        .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension")
        .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog")
        .getOrCreate()
    )

    spark.conf.set(f"fs.azure.account.key.{STORAGE_ACCOUNT_NAME}.blob.core.windows.net", BLOB_CREDENTIAL)

    print("**** spark.rpc.message.maxSize = ", spark.conf.get("spark.rpc.message.maxSize"))

    if start_chunk == 0:
        # Delete file if exists
        print("**** Deleting existing delta table")
        if fs.exists(f"{CONTAINER_NAME}/{delta_file}"):
            fs.rm(f"{CONTAINER_NAME}/{delta_file}", recursive=True)

    chunksize = 3_000_000

    with fs.open(f"{CONTAINER_NAME}/{zip_file}", "rb") as file:
        with zipfile.ZipFile(file, "r") as zip_ref:
            file_name = zip_ref.namelist()[0]

            with zip_ref.open(file_name) as csv_file:
                csv_io = TextIOWrapper(csv_file, "utf-8")
                headers = pd.read_csv(csv_io, sep="\t", nrows=0).columns.tolist()
                chunk_iter = pd.read_csv(
                    csv_io,
                    sep="\t",
                    header=None,
                    names=headers,
                    usecols=["col1", "col2", "col3"],
                    dtype=str,  # Read all as strings to avoid errors
                    chunksize=chunksize,
                    skiprows=start_chunk*chunksize
                )

                for chunk in tqdm(chunk_iter, desc="Processing chunks"):
                    # Convert pd DataFrame to Spark DataFrame
                    spark_df = spark.createDataFrame(chunk)

                    (spark_df.repartition(8).write
                     .format("delta")
                     .mode("append")
                     .option("mergeSchema", "true")
                     .save(delta_table_path)
                    )

                    # Clear memory after each iteration
                    spark_df.unpersist(blocking=True)
                    del chunk
                    del spark_df
                    gc.collect()

备注:内容来源于stack exchange,提问作者yellowsubmarine

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.13 18:50:27