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

使用ADF实现SFTP至ADLS Synapse文件重名时自动递增命名

SFTP上传ADLS Synapse时同名文件递增命名实现方案

核心思路

先检查ADLS目标路径下的现有文件,匹配待上传文件的名称前缀与扩展名,提取已存在的递增序号,生成新的文件名后再执行复制操作,避免直接覆盖原有文件。

方案一:用Synapse Pipeline内置活动实现

步骤1:获取目标目录文件元数据

添加Get Metadata活动,配置ADLS Gen2的连接信息,目标路径设为文件要上传的目录,在"字段列表"中勾选Child items,用于获取该目录下的所有文件列表。

步骤2:生成递增文件名

  1. 先创建两个变量:
    • fileNamePrefix:存储文件名前缀,表达式为 @split(pipeline().parameters.sourceFileName, '.')[0]
    • fileExtension:存储文件扩展名,表达式为 @concat('.', split(pipeline().parameters.sourceFileName, '.')[1])
  2. 添加Filter活动,输入为@activity('Get Metadata').output.childItems,过滤条件设为:
    @or(equals(item().name, pipeline().parameters.sourceFileName), and(startsWith(item().name, concat(variables('fileNamePrefix'), '_')), endsWith(item().name, variables('fileExtension'))))
    
  3. 添加Set Variable活动,生成最终文件名,表达式示例:
    @if(contains(string(activity('Get Metadata').output.childItems), pipeline().parameters.sourceFileName),
        concat(variables('fileNamePrefix'), '_', string(max(union(createArray(1), split(replace(join(activity('Filter').output.value, ','), concat(variables('fileNamePrefix'), '_'), ''), variables('fileExtension'))))), variables('fileExtension')),
        pipeline().parameters.sourceFileName
    )
    
    逻辑:如果原文件名已存在,就提取现有递增序号的最大值加1;如果不存在,直接用原文件名。

步骤3:执行文件复制

添加Copy Data活动,源配置为SFTP的待上传文件,目标配置为ADLS的对应目录,将目标文件名设置为刚才生成的变量值,完成无覆盖的上传。

方案二:用Python脚本实现(适合自定义场景)

如果需要更灵活的逻辑,可在Synapse Notebook或本地脚本中实现:

1. 依赖库安装

pip install azure-storage-file-datalake paramiko

2. 核心代码实现

import os
from azure.storage.filedatalake import DataLakeServiceClient
import paramiko

def generate_new_adls_filename(adls_client, filesystem, target_dir, original_name):
    # 拆分文件名与扩展名
    file_prefix, file_ext = os.path.splitext(original_name)
    # 获取ADLS目标目录下的所有文件
    fs_client = adls_client.get_file_system_client(filesystem=filesystem)
    dir_client = fs_client.get_directory_client(target_dir)
    existing_files = [path.name for path in dir_client.list_paths() if not path.is_directory]
    
    # 匹配同前缀的文件
    matched_files = [f for f in existing_files if f.startswith(file_prefix) and f.endswith(file_ext)]
    if not matched_files:
        return original_name
    # 处理已有原文件的情况
    if len(matched_files) == 1 and matched_files[0] == original_name:
        return f"{file_prefix}_1{file_ext}"
    # 提取最大序号并生成新文件名
    max_seq = 0
    for f in matched_files:
        if f == original_name:
            current_seq = 0
        else:
            seq_part = f.replace(f"{file_prefix}_", "").replace(file_ext, "")
            if seq_part.isdigit():
                current_seq = int(seq_part)
                if current_seq > max_seq:
                    max_seq = current_seq
    return f"{file_prefix}_{max_seq + 1}{file_ext}"

# 示例:连接ADLS并生成新文件名
adls_account_url = "https://<your-adls-account>.dfs.core.windows.net/"
credential = "<your-credential>"  # 可使用SAS token或服务主体密钥
adls_client = DataLakeServiceClient(account_url=adls_account_url, credential=credential)

new_filename = generate_new_adls_filename(
    adls_client,
    filesystem_name="your-filesystem",
    target_dir="upload/employees",
    original_name="Employee_213.parquet"
)

# 从SFTP下载文件并上传到ADLS
def sftp_to_adls(sftp_host, sftp_port, sftp_user, sftp_pass, sftp_file_path, adls_client, filesystem, adls_target_dir, new_filename):
    # 连接SFTP
    ssh_client = paramiko.SSHClient()
    ssh_client.set_missing_host_key_policy(paramiko.AutoAddPolicy())
    ssh_client.connect(hostname=sftp_host, port=sftp_port, username=sftp_user, password=sftp_pass)
    sftp_client = ssh_client.open_sftp()
    
    # 读取SFTP文件内容
    with sftp_client.open(sftp_file_path, 'rb') as f:
        file_content = f.read()
    
    # 上传到ADLS
    fs_client = adls_client.get_file_system_client(filesystem=filesystem)
    dir_client = fs_client.get_directory_client(adls_target_dir)
    file_client = dir_client.create_file(new_filename)
    file_client.upload_data(file_content, overwrite=False)
    
    # 关闭连接
    sftp_client.close()
    ssh_client.close()

# 调用上传函数
sftp_to_adls(
    sftp_host="<your-sftp-host>",
    sftp_port=22,
    sftp_user="<sftp-username>",
    sftp_pass="<sftp-password>",
    sftp_file_path="/remote/path/Employee_213.parquet",
    adls_client=adls_client,
    filesystem="your-filesystem",
    adls_target_dir="upload/employees",
    new_filename=new_filename
)

注意事项

  • 并发处理:如果存在多线程/多进程同时上传同名文件的场景,建议先创建空文件(原子操作)再写入内容,避免出现重复命名。
  • 性能优化:若目标目录文件数量过多,文件列表查询会变慢,建议按日期、业务维度拆分目录,减少单目录文件数。
  • 权限配置:确保执行主体拥有SFTP的读取权限,以及ADLS的读取、写入权限。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 00:45:16