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

