如何使用ADF实现数据湖文件名与CSV的匹配校验及后续处理
技术方案:数据湖文件名比对与批量复制
核心逻辑
- 读取指定CSV文件中的目标文件名列表
- 遍历数据湖目标文件夹,收集所有实际存在的文件名
- 双向比对后执行两个操作:
- 匹配的文件:复制到指定目标路径
- 不匹配的文件(数据湖中存在但CSV未收录的):收集后写入新CSV文件存储到数据湖
分步实现(Python示例)
以AWS S3作为数据湖存储为例,使用boto3处理存储操作,pandas处理CSV数据:
1. 依赖安装
pip install boto3 pandas
2. 代码实现
import boto3 import pandas as pd from io import StringIO # 配置参数(根据实际环境修改) S3_BUCKET = "your-data-lake-bucket" SOURCE_FOLDER = "raw-data/unprocessed/" # 数据湖待比对文件夹 TARGET_FOLDER = "processed/matched-files/" # 匹配文件的目标存储路径 TARGET_CSV_KEY = "config/target-filenames.csv" # 存储目标文件名的CSV路径 UNMATCHED_CSV_KEY = "logs/unmatched-filenames.csv" # 不匹配文件名的输出路径 # 初始化S3客户端 s3 = boto3.client('s3') def load_target_filenames(): """从数据湖CSV读取目标文件名集合""" csv_obj = s3.get_object(Bucket=S3_BUCKET, Key=TARGET_CSV_KEY) csv_content = csv_obj['Body'].read().decode('utf-8') df = pd.read_csv(StringIO(csv_content)) # 假设CSV有一列名为"filename"存储待匹配文件名 return set(df['filename'].dropna().tolist()) def scan_datalake_filenames(): """遍历数据湖文件夹,获取所有实际存在的文件名""" paginator = s3.get_paginator('list_objects_v2') response_iterator = paginator.paginate(Bucket=S3_BUCKET, Prefix=SOURCE_FOLDER) datalake_files = set() for response in response_iterator: if 'Contents' in response: for obj in response['Contents']: # 提取纯文件名(去除路径前缀) file_key = obj['Key'] filename = file_key.split('/')[-1] if filename: # 排除空文件夹条目 datalake_files.add(filename) return datalake_files def copy_matched_items(matched_files): """将匹配的文件复制到目标路径""" for filename in matched_files: source_key = f"{SOURCE_FOLDER}{filename}" target_key = f"{TARGET_FOLDER}{filename}" # 执行S3内部复制(无需下载再上传) s3.copy_object( Bucket=S3_BUCKET, Key=target_key, CopySource={'Bucket': S3_BUCKET, 'Key': source_key} ) print(f"完成复制: {filename}") def save_unmatched_list(unmatched_files): """将不匹配的文件名写入数据湖CSV""" df = pd.DataFrame({'unmatched_filename': list(unmatched_files)}) csv_buffer = StringIO() df.to_csv(csv_buffer, index=False) s3.put_object( Bucket=S3_BUCKET, Key=UNMATCHED_CSV_KEY, Body=csv_buffer.getvalue(), ContentType='text/csv' ) print(f"不匹配文件列表已保存至: {UNMATCHED_CSV_KEY}") if __name__ == "__main__": target_files = load_target_filenames() datalake_files = scan_datalake_filenames() # 计算匹配和不匹配的文件集合 matched = target_files.intersection(datalake_files) unmatched = datalake_files.difference(target_files) copy_matched_items(matched) save_unmatched_list(unmatched)
适配其他数据湖存储
Azure ADLS Gen2
替换存储客户端为azure-storage-file-datalake,核心逻辑不变:
from azure.storage.filedatalake import DataLakeServiceClient # 初始化ADLS客户端 service_client = DataLakeServiceClient(account_url="https://<account-name>.dfs.core.windows.net/", credential="<your-token>") file_system_client = service_client.get_file_system_client(file_system="<file-system-name>")
文件遍历、读取、复制操作对应替换为ADLS的API即可。
本地模拟数据湖
如果用本地文件夹模拟数据湖,直接用os和shutil处理:
import os import shutil # 读取本地CSV df = pd.read_csv("/local-datalake/config/target-filenames.csv") target_files = set(df['filename'].dropna().tolist()) # 遍历本地文件夹 datalake_files = set(os.listdir("/local-datalake/raw-data/unprocessed/")) # 复制匹配文件 for filename in matched: shutil.copy(f"/local-datalake/raw-data/unprocessed/{filename}", f"/local-datalake/processed/matched-files/{filename}")
关键注意事项
- 大小写敏感问题:不同数据湖存储对文件名大小写的处理不同(如S3默认敏感,ADLS可配置),需统一比对规则
- 大数量文件处理:文件数过万时,务必用分页遍历+批量操作,避免内存溢出
- 权限校验:确保执行代码的账号拥有数据湖的读、写、复制权限
- CSV格式校验:提前清理CSV中的多余空格、换行符,避免匹配失败
- 幂等性设计:可给已复制文件添加前缀(如
copied-),避免重复执行时重复复制
内容的提问来源于stack exchange,提问作者Anonymous
相关产品推荐
相关产品推荐

