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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 15:55:53