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

如何高效将大型PostgreSQL数据加载转换至Databricks避免超时

问题:Databricks合并PostgreSQL数据后写入Azure Blob Storage超时优化

需要将三个PostgreSQL表的数据加载至Databricks,合并后写入Azure Blob Storage为CSV文件。目前已成功从PostgreSQL直接加载生成df1和df2,单独展示这两个DataFrame正常,但合并得到final_df后,执行display(final_df)或写入Parquet文件时出现超时。尝试按列分区写入仍报错:

org.apache.spark.SparkException: Job aborted due to stage failure: Task 0 in stage 19.0 failed 4 times, most recent failure: Lost task 0.3 in stage 19.0 (TID 36) (10.139.64.11 executor 4): ExecutorLostFailure (executor 4 exited caused by one of the running tasks) Reason: Executor heartbeat timed out after 125857 ms

已增加计算资源但问题依旧,寻求高效优化方案。


相关代码

加载df1的代码

# Connection properties
jdbc_url = "jdbc:postgresql://server1.postgres.database.azure.com:port_num/database1"
connection_properties = {
    "user": "username1@server1",
    "password": "password1",
    "driver": "org.postgresql.Driver",
    "sslmode": "sslmodeoption"
}

#Schema / Table details
schema_name = "schemaname1"
table1_name = '"table1"'
table2_name = '"table2"'

# Construct table paths
table1_path = f"{schema_name}.{table1_name}"
table2_path = f"{schema_name}.{table2_name}"

# SQL query with join operations to join data on load
query = f"(SELECT t1.*, t2.\"TotalValue\", t2.\"TotalTime\" FROM {table1_path} AS t1 INNER JOIN {table2_path} AS t2 ON t1.\"PKey\" = t2.\"Id\") AS joined_data"

# Read data from the specified tables and create spark dataframe
df1 = spark.read.jdbc(url=jdbc_url, table=query, properties=connection_properties)

加载df2的代码

# Connection properties
jdbc_url = "jdbc:postgresql://server1.postgres.database.azure.com:port_num/database2"
connection_properties = {
    "user": "username2@server1",
    "password": "password2",
    "driver": "org.postgresql.Driver",
    "sslmode": "sslmodeoption"
}

# Schema / table properties
schema_name = "schemaname1"
table_name = '"table3"'

# Construct table paths
table_path = f"{schema_name}.{table_name}"

df2 = spark.read.jdbc(url=jdbc_url, table=table_path, properties=connection_properties)

合并数据与写入代码

final_df = spark.sql("""SELECT d1.*,d2.name FROM df1 as d1 INNER JOIN df2 as d2 ON df1.a_ID = df2.ID""")

# filepath
file_path = "/mnt/area1/Staging"

# Define the columns to partition by
partition_columns = ["name", "id2"]

# Write the DataFrame to Azure Blob Storage partitioned by the specified columns
final_df.write \
        .format("parquet") \
        .mode("overwrite") \
        .partitionBy(*partition_columns) \
        .save(file_path)

优化方案
  • 数据库端提前完成合并:将final_df的连接逻辑迁移至PostgreSQL执行,直接读取最终合并后的数据集,减少Databricks侧的数据处理量。例如修改df1的查询,加入table3的连接逻辑,避免在Spark中进行大表关联。
  • 提升JDBC读取并行度:读取PostgreSQL数据时,添加分区参数让Spark并行加载,避免单分区承载过大数据量。示例:
    df2 = spark.read.jdbc(url=jdbc_url, 
                          table=table_path, 
                          properties=connection_properties,
                          partitionColumn="ID",
                          lowerBound=1,
                          upperBound=1000000,
                          numPartitions=10)
    
  • 调整Spark分区策略:合并后对final_df重分区,根据集群资源和数据量设置合理的分区数,避免分区过大导致任务超时:
    final_df = final_df.repartition(20)  # 数值根据实际情况调整
    
  • 优化写入逻辑:
    • 避免使用高基数列作为分区键,name如果是高基数列会生成大量小文件,增加写入开销,可更换为低基数列或减少分区层级。
    • 写入前执行final_df.cache()缓存数据,避免重复计算;写入完成后执行final_df.unpersist()释放内存。
    • 若目标是CSV文件,直接指定.format("csv")写入,无需先写Parquet再转换,减少中间步骤。
  • 调整Spark核心配置:
    • 增大executor内存与核心数,同时调整spark.sql.shuffle.partitions(默认200),根据数据量设置为500-1000区间。
    • 开启自适应执行:spark.sql.adaptive.enabled=true,让Spark自动调整执行计划与分区数。

内容的提问来源于stack exchange,提问作者JGW

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 07:31:37