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

如何使用ADF实现数据湖文件名与CSV的匹配校验及后续处理

技术方案:数据湖文件名比对与批量复制

核心逻辑

  1. 读取指定CSV文件中的目标文件名列表
  2. 遍历数据湖目标文件夹,收集所有实际存在的文件名
  3. 双向比对后执行两个操作:
    • 匹配的文件:复制到指定目标路径
    • 不匹配的文件(数据湖中存在但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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 01:36:21