如何在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
相关产品推荐
相关产品推荐

