Databricks中PySpark并行复制文件时dbutils使用报错咨询
Databricks中PySpark并行文件复制报错问题
问题背景
想要在Databricks中借助多Executor实现文件复制的并行执行,编写了PySpark代码后,在DataFrame的foreach算子中调用dbutils.fs.cp时出现报错,但相同逻辑的Scala代码可以正常运行,疑问是否PySpark API未更新支持该场景。
问题代码(PySpark)
def parallel_copy_execution(src_path: str, target_path: str): files_in_path = dbutils.fs.ls(src_path) file_paths_df = spark.sparkContext.parallelize(files_in_path).toDF() file_paths_df.foreach(lambda x: dbutils.fs.cp(x.path.toString(), target_path, recurse=True))
报错信息
You cannot use dbutils within a spark job or otherwise pickle it. If you need to use getArguments within a spark job, you have to get the argument before using it in the job. For example, if you have the following code: myRdd.map(lambda i: dbutils.args.getArgument("X") + str(i)) Then you should use it this way: argX = dbutils.args.getArgument("X") myRdd.map(lambda i: argX + str(i))
可行的Scala代码
def parallel_copy_execution(p: String, t: String): Unit = { dbutils.fs.ls(p).map(_.path).toDF.foreach { file => dbutils.fs.cp(file(0).toString,t , recurse=true) println(s"cp file: $file") } }
原因解析
不是PySpark API未更新,而是PySpark与Scala的序列化机制存在差异:
dbutils是Databricks提供的工具对象,在PySpark中,该对象无法被序列化(pickle)。而Spark的分布式算子(如foreach)会将传入的函数序列化后分发到各个Executor节点执行,因此直接在算子中调用dbutils会触发序列化失败的报错。- Scala版本的
dbutils经过Databricks的特殊处理,支持在Executor端反序列化并使用,因此Scala代码可以正常运行。
解决方案
方法1:使用Hadoop原生FileSystem API(推荐,利用Executor并行)
通过Hadoop的FileSystem API实现文件复制,该API可被序列化,能在Executor节点上执行,真正实现分布式并行复制:
from py4j.java_gateway import java_import from pyspark.sql import SparkSession def copy_files_partition(iterator, target_path): # 获取Hadoop FileSystem实例 spark = SparkSession.getActiveSession() java_import(spark._jvm, "org.apache.hadoop.fs.FileSystem") java_import(spark._jvm, "org.apache.hadoop.fs.Path") fs = spark._jvm.FileSystem.get(spark._jsc.hadoopConfiguration()) for row in iterator: src_path = row.path target_full_path = f"{target_path.rstrip('/')}/{src_path.split('/')[-1]}" # 根据文件系统类型调整复制方法,此处以本地为例 fs.copyToLocalFile(False, spark._jvm.Path(src_path), spark._jvm.Path(target_full_path), True) def parallel_copy_execution(src_path: str, target_path: str): files_in_path = dbutils.fs.ls(src_path) file_paths_df = spark.sparkContext.parallelize(files_in_path).toDF() # 使用foreachPartition减少序列化开销 file_paths_df.foreachPartition(lambda iter: copy_files_partition(iter, target_path))
方法2:Driver端多线程并行(仅利用Driver资源)
如果不需要利用Executor资源,可将文件路径收集到Driver端后,用Python多线程实现并行复制:
import concurrent.futures def copy_single_file(src, target): dbutils.fs.cp(src, target, recurse=True) def parallel_copy_execution(src_path: str, target_path: str): files_in_path = dbutils.fs.ls(src_path) file_paths = [f.path for f in files_in_path] with concurrent.futures.ThreadPoolExecutor(max_workers=8) as executor: futures = [executor.submit(copy_single_file, path, f"{target_path.rstrip('/')}/{path.split('/')[-1]}") for path in file_paths] concurrent.futures.wait(futures)
内容的提问来源于stack exchange,提问作者Nandini Raja
相关产品推荐
相关产品推荐

