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

基于PySpark实现Azure存储容器增量文件复制(保留目录结构)

增量复制Azure存储文件并保留目录结构的解决方案

核心思路

通过跟踪已复制文件的元数据(如文件路径、修改时间),对比源文件状态,仅复制新增或更新的文件;同时通过保留相对路径的方式维持原目录层级。以下是基于Databricks PySpark的落地实现:

步骤1:配置基础参数

修正原代码路径错误,定义存储与路径相关配置:

from datetime import datetime
from pyspark.sql.functions import col
from pyspark.sql.types import StructType, StructField, StringType, TimestampType

# 存储账户与容器配置
storage_account = "dexflex"
source_container = "source"
destination_container = "destination"

# 源/目标根路径
source_root = f"abfss://{source_container}@{storage_account}.dfs.core.windows.net/"
dest_root = f"abfss://{destination_container}@{storage_account}.dfs.core.windows.net/"

# 元数据存储路径(记录已复制文件信息)
metadata_path = f"{dest_root}.copied_files_metadata"

步骤2:加载已复制文件的元数据

用Parquet文件持久化已复制文件的记录,首次运行时自动初始化空数据集:

# 定义元数据Schema
metadata_schema = StructType([
    StructField("file_path", StringType(), nullable=False),
    StructField("last_modified", TimestampType(), nullable=False)
])

# 加载已有元数据,若不存在则创建空DataFrame
try:
    copied_files_df = spark.read.schema(metadata_schema).parquet(metadata_path)
except Exception:
    copied_files_df = spark.createDataFrame([], metadata_schema)

步骤3:获取源目录所有文件信息

递归遍历源目录,提取文件相对路径与修改时间:

def list_all_files(path):
    files = []
    for item in dbutils.fs.ls(path):
        if item.isFile():
            # 提取相对路径,用于后续维持目录结构
            relative_path = item.path.replace(source_root, "")
            files.append((relative_path, item.modificationTime))
        else:
            files.extend(list_all_files(item.path))
    return files

# 转换为DataFrame便于后续对比
source_files_df = spark.createDataFrame(list_all_files(source_root), ["relative_path", "last_modified"])

步骤4:筛选需复制的新增/更新文件

对比源文件与已复制记录,找出未复制或修改时间更新的文件:

files_to_copy_df = source_files_df.join(
    copied_files_df,
    source_files_df.relative_path == copied_files_df.file_path,
    "left_outer"
).where(
    copied_files_df.file_path.isNull() | 
    (source_files_df.last_modified > copied_files_df.last_modified)
).select(source_files_df.relative_path, source_files_df.last_modified)

步骤5:执行增量复制并更新元数据

按相对路径复制文件到目标目录,同时更新已复制记录:

# 遍历需复制的文件列表
for row in files_to_copy_df.collect():
    relative_path = row["relative_path"]
    source_file = f"{source_root}{relative_path}"
    dest_file = f"{dest_root}{relative_path}"
    
    # 创建目标目录(避免因目录不存在导致复制失败)
    dest_dir = dest_file.rsplit("/", 1)[0]
    dbutils.fs.mkdirs(dest_dir)
    
    # 执行复制
    dbutils.fs.cp(source_file, dest_file)
    print(f"Copied: {source_file} -> {dest_file}")

# 更新元数据,合并新复制文件与原有记录
new_copied_df = files_to_copy_df.withColumnRenamed("relative_path", "file_path")
updated_metadata_df = copied_files_df.unionByName(new_copied_df).dropDuplicates(["file_path"])
updated_metadata_df.write.mode("overwrite").parquet(metadata_path)

关键说明

  • 目录结构保留:通过提取相对路径,在目标容器中重建与源一致的层级结构,确保文件存放路径不变。
  • 增量逻辑:基于文件修改时间与已复制记录对比,仅处理新增或更新的文件,避免全量复制。
  • 元数据持久化:用Parquet文件存储已复制记录,避免每次运行重复扫描全量文件。
  • 环境适配:该方案依赖Databricks的dbutils,若为纯PySpark环境,可替换为Hadoop FS API或Azure Storage SDK实现文件操作。

替代简化方案

若无需代码实现,可直接使用Azure Data Factory(ADF):

  • 创建复制活动,源选择Azure Blob存储并开启递归遍历。
  • 启用增量复制模式(按最后修改时间或文件列表),目标配置为保留源目录结构,ADF会自动完成增量同步。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 10:13:16