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

如何通过Airflow DAG从S3桶文件夹选取最新文件加载至Snowflake表

从S3指定文件夹选取最新文件加载到Snowflake的Airflow实现

需求背景

作为Airflow DAG新手,需要实现从指定S3桶文件夹中选取最新生成的文件(文件名后缀带当前日期时间,格式如original_df_with_clusters_Fiserv_202309070657.csv),并将其数据加载到Snowflake表中。当前代码会每次重建表并加载文件夹下所有文件,不符合需求。

S3信息:

  • Bucket名称:aura-skills-data
  • 目标文件夹:pdl-skills-data/pdl_result_data/

当前代码问题

现有COPY INTO语句会加载指定S3文件夹下的所有文件,且每次执行CREATE OR REPLACE TABLE会清空原有表数据,无法实现仅加载最新文件的需求。

解决方案实现

核心思路

  • 通过Airflow的S3Hook获取目标文件夹下的所有文件列表
  • 解析每个文件名中的时间戳(格式为YYYYMMDDHHMM),筛选出时间戳最大的最新文件
  • 修改COPY INTO语句,仅指向该最新文件
  • 调整表创建逻辑:若表不存在则创建,存在则直接追加数据(避免每次清空原有数据)

完整代码示例

from airflow.providers.amazon.aws.hooks.s3 import S3Hook
from airflow.providers.snowflake.hooks.snowflake import SnowflakeHook
import logging

logger = logging.getLogger(__name__)

def upload_latest_s3_file_to_snowflake(table_name, storage_integration):
    # 1. 获取S3目标文件夹下的所有文件
    s3_hook = S3Hook(aws_conn_id="your_aws_connection_id")  # 替换为你的Airflow AWS连接ID
    bucket_name = "aura-skills-data"
    prefix = "pdl-skills-data/pdl_result_data/"
    
    # 获取文件夹下所有对象,过滤掉前缀本身(避免把文件夹路径当成文件)
    s3_objects = s3_hook.list_keys(bucket_name=bucket_name, prefix=prefix)
    file_list = [obj for obj in s3_objects if obj != prefix]
    
    if not file_list:
        logger.warning("S3目标文件夹下无文件,终止加载流程")
        return
    
    # 2. 解析文件名,找出最新文件
    def extract_timestamp(file_path):
        # 从文件名中提取时间戳部分:例如从"original_df_with_clusters_Fiserv_202309070657.csv"取"202309070657"
        file_name = file_path.split("/")[-1]
        timestamp_part = file_name.split("_")[-1].replace(".csv", "")
        return timestamp_part
    
    # 按时间戳降序排序,取第一个即为最新文件
    file_list_sorted = sorted(file_list, key=lambda x: extract_timestamp(x), reverse=True)
    latest_file_path = file_list_sorted[0]
    logger.info(f"找到最新文件:{latest_file_path}")
    
    # 3. 创建表(若不存在)
    create_table_stmt = f"""
        CREATE TABLE IF NOT EXISTS {table_name} (
            INDEX varchar,
            EMPLOYEE_ID varchar,
            skills varchar,
            hc varchar,
            cluster_descr_hc varchar
        );
    """
    snowflake_hook = SnowflakeHook(snowflake_conn_id="your_snowflake_connection_id")  # 替换为你的Airflow Snowflake连接ID
    snowflake_hook.run(create_table_stmt)
    logger.info(f"确保表{table_name}已存在")
    
    # 4. 仅复制最新文件到Snowflake表
    copy_stmt = f"""
        COPY INTO {table_name}
        FROM 's3://{bucket_name}/{latest_file_path}' 
        STORAGE_INTEGRATION={storage_integration}
        FILE_FORMAT = (
            TYPE = 'CSV' 
            NULL_IF = ('','NULL') 
            BINARY_FORMAT='UTF-8' 
            FIELD_OPTIONALLY_ENCLOSED_BY='"' 
            COMPRESSION='NONE' 
            SKIP_HEADER = 1
        );
    """
    copy_result = snowflake_hook.run(copy_stmt)
    logger.info(f"最新文件加载完成,结果:{copy_result}")

关键注意事项

  • 替换代码中的your_aws_connection_id和your_snowflake_connection_id为你在Airflow中配置的对应连接ID
  • 若需要避免重复加载同一文件,可以在加载完成后将文件移动到S3的归档文件夹(例如pdl-skills-data/archived_data/),后续只处理未归档的文件
  • 时间戳解析逻辑需与你的文件名格式匹配,若后续文件名格式变更,需同步调整extract_timestamp函数

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 16:20:22