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

如何在Databricks Notebook中并行执行多个DataFrame加载任务?

在Databricks Notebook中并行加载多个DataFrame的方法

在Databricks Notebook里并行加载DataFrame来提升ETL效率是个很实用的需求,我平时处理这类场景时常用下面几种方法,亲测有效:

方法一:使用Python concurrent.futures.ThreadPoolExecutor

加载数据(比如读CSV、Parquet、JDBC等)大多属于IO密集型操作,线程池可以让多个读取任务同时执行,充分利用集群的空闲资源,减少总加载时间。

代码示例:

from concurrent.futures import ThreadPoolExecutor

# 定义两个独立的DataFrame加载函数
def load_user_data():
    return spark.read.format("delta").load("/databricks-datasets/learning-spark-v2/userdata/userdata1.parquet")

def load_sales_data():
    return spark.read.csv("/databricks-datasets/retail-org/sales/sales.csv", header=True, inferSchema=True)

# 初始化线程池,设置并行数为2(对应两个DF加载任务)
with ThreadPoolExecutor(max_workers=2) as executor:
    # 提交加载任务到线程池
    future_user_df = executor.submit(load_user_data)
    future_sales_df = executor.submit(load_sales_data)
    
    # 获取加载完成的DataFrame
    user_df = future_user_df.result()
    sales_df = future_sales_df.result()

# 验证加载结果
print("用户数据行数:", user_df.count())
print("销售数据行数:", sales_df.count())

方法二:利用Spark作业的并行调度

如果你的数据源是分布式存储(比如S3、ADLS、HDFS),Spark本身会并行处理数据分区,但如果是加载多个独立的数据集,你可以通过提交多个Spark作业来实现并行。不过这种方法需要注意集群的资源分配,避免资源竞争。

关键注意事项:

  • SparkSession是线程安全的,所以在多线程环境中可以安全使用,无需重新创建。
  • 如果加载的是JDBC数据源,要注意数据库的连接数限制,避免并行任务过多导致数据库拒绝连接。
  • 根据集群的核心数调整max_workers的值,不要设置过大,否则会导致资源耗尽,反而降低效率。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 08:10:39