如何高效将大型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自动调整执行计划与分区数。
- 增大executor内存与核心数,同时调整
内容的提问来源于stack exchange,提问作者JGW
相关产品推荐
相关产品推荐

