如何通过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
相关产品推荐
相关产品推荐

