使用ADF实现SFTP至ADLS Synapse文件重名时自动递增命名
SFTP上传ADLS Synapse时同名文件递增命名实现方案
核心思路
先检查ADLS目标路径下的现有文件,匹配待上传文件的名称前缀与扩展名,提取已存在的递增序号,生成新的文件名后再执行复制操作,避免直接覆盖原有文件。
方案一:用Synapse Pipeline内置活动实现
步骤1:获取目标目录文件元数据
添加Get Metadata活动,配置ADLS Gen2的连接信息,目标路径设为文件要上传的目录,在"字段列表"中勾选Child items,用于获取该目录下的所有文件列表。
步骤2:生成递增文件名
- 先创建两个变量:
fileNamePrefix:存储文件名前缀,表达式为@split(pipeline().parameters.sourceFileName, '.')[0]fileExtension:存储文件扩展名,表达式为@concat('.', split(pipeline().parameters.sourceFileName, '.')[1])
- 添加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')))) - 添加Set Variable活动,生成最终文件名,表达式示例:
逻辑:如果原文件名已存在,就提取现有递增序号的最大值加1;如果不存在,直接用原文件名。@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 )
步骤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
相关产品推荐
相关产品推荐

