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

如何通过GCP StorageTransfer获取S3最新文件?Airflow/Composer场景

问题

我们在GCP上运行Airflow/Composer,尝试通过Python Operator获取S3中的最新文件。目前每日执行一次S3文件获取操作,但单日可能生成多个文件,我们仅需获取其中最新的一个。请问有哪些实现方案?

当前使用的脚本如下:

create_transfer_job = CloudDataTransferServiceCreateJobOperator(
                task_id="create_transfer_job",
                body=self.create_transfer_body(self),
                aws_conn_id=aws_conn_id(),
                on_success_callback=partial(self.write_logging_url, 
                  self.get_tmp_project_id()),
            )

   
def create_transfer_body(self):
       
        dag_parameters = get_dag_parameters()
        data_transfer = dag_parameters["data_transfer"]
        s3_data_source = data_transfer["awsS3DataSource"]

        s3_data_source_path = ""
        if "path" in s3_data_source:
            s3_data_source_path = s3_data_source["path"]

        # check
        print("s3_data_source_path!:"+s3_data_source_path)

        file_types = file_types()
        includePrefixes = []
        for file_type in file_types:
            includePrefixes.append(f"{file_type}/")

        objectConditions = {
            "includePrefixes": includePrefixes
            ,"lastModifiedSince": "{{ ti.xcom_pull(task_ids='initialize', key='last_modified_since') }}"
            ,"lastModifiedBefore": "{{ ti.xcom_pull(task_ids='initialize', key='last_modified_before') }}"
        }
        if "objectConditions" in data_transfer:
            objectConditions = data_transfer["objectConditions"]

        transfer_body = {
            "description": "DataPipe Data Transfer",
            "status": GcpTransferJobsStatus.ENABLED,
            "projectId": project_id(),
            "name": "transferJobs/dataPipelineJob_" + str(uuid.uuid4()),
            "schedule": {
                "scheduleStartDate": datetime.today(),
                "scheduleEndDate": datetime.today(),
            },
            "transferSpec": {
                "awsS3DataSource": {"bucketName": s3_data_source["bucketName"], "path": s3_data_source_path},
                "gcsDataSink": {"bucketName": self.get_gcs_bucket_name()},
                "objectConditions": objectConditions,
                "transfer_options": {
                    "deleteObjectsFromSourceAfterTransfer": False,
                    "overwriteWhen": "{{ ti.xcom_pull(task_ids='initialize', key='overwrite_when') }}",
                },
            },
            "loggingConfig": {
                "logActions": ["FIND", "DELETE", "COPY"],
                "logActionStates": ["FAILED"],
            },
        }
        return transfer_body

实现方案

方案1:Python Operator枚举S3文件筛选最新项后触发传输

在现有DAG中新增PythonOperator任务,通过boto3调用AWS SDK列出指定S3路径下的当日文件,按最后修改时间排序取最新文件,将其路径传入后续的CloudDataTransferServiceCreateJobOperator,修改传输目标为该具体文件而非前缀路径。

核心代码示例:

from airflow.operators.python import PythonOperator
import boto3
from airflow.models import Variable

def get_latest_s3_file(**context):
    # 初始化S3客户端
    s3_client = boto3.client('s3', 
                             aws_access_key_id=Variable.get("aws_access_key"),
                             aws_secret_access_key=Variable.get("aws_secret_key"))
    # 从DAG参数中获取桶名和前缀
    s3_conf = context["dag_run"].conf["data_transfer"]["awsS3DataSource"]
    bucket_name = s3_conf["bucketName"]
    prefix = s3_conf["path"]
    
    # 列出前缀下所有对象
    response = s3_client.list_objects_v2(Bucket=bucket_name, Prefix=prefix)
    if "Contents" not in response:
        raise ValueError("指定S3路径下无文件")
    
    # 按最后修改时间降序排序,取第一个
    sorted_objects = sorted(response["Contents"], key=lambda x: x["LastModified"], reverse=True)
    latest_object_key = sorted_objects[0]["Key"]
    
    # 将最新文件路径存入XCom
    context["ti"].xcom_push(key="latest_s3_file_key", value=latest_object_key)
    return latest_object_key

# 新增获取最新文件的任务
get_latest_file = PythonOperator(
    task_id="get_latest_s3_file",
    python_callable=get_latest_s3_file,
    provide_context=True,
    dag=dag
)

# 修改create_transfer_body方法,从XCom获取最新文件路径
def create_transfer_body(self, **context):
    # 保留原有逻辑
    dag_parameters = get_dag_parameters()
    data_transfer = dag_parameters["data_transfer"]
    s3_data_source = data_transfer["awsS3DataSource"]
    # ... 其他原有代码 ...
    
    # 替换为最新文件路径
    latest_file_key = context["ti"].xcom_pull(task_ids="get_latest_s3_file", key="latest_s3_file_key")
    transfer_body["transferSpec"]["awsS3DataSource"]["path"] = latest_file_key
    
    return transfer_body

# 调整任务依赖:get_latest_file >> create_transfer_job

方案2:传输后在GCS端筛选保留最新文件

保持现有CloudDataTransferServiceCreateJobOperator逻辑,先将当日所有符合条件的文件传输到GCS,再新增PythonOperator任务,在GCS中筛选出最新文件,删除或归档其余文件。

核心代码示例:

from airflow.operators.python import PythonOperator
from google.cloud import storage

def keep_only_latest_gcs_file(**context):
    gcs_client = storage.Client()
    bucket_name = self.get_gcs_bucket_name()
    bucket = gcs_client.get_bucket(bucket_name)
    
    # 列出GCS目标路径下所有文件
    s3_prefix = context["dag_run"].conf["data_transfer"]["awsS3DataSource"]["path"]
    blobs = bucket.list_blobs(prefix=s3_prefix)
    blob_list = list(blobs)
    
    if len(blob_list) <= 1:
        return
    
    # 按更新时间降序排序,保留最新文件
    sorted_blobs = sorted(blob_list, key=lambda x: x.updated, reverse=True)
    for blob in sorted_blobs[1:]:
        blob.delete()
    
    # 将最新文件路径存入XCom供后续任务使用
    context["ti"].xcom_push(key="latest_gcs_file_path", value=f"gs://{bucket_name}/{sorted_blobs[0].name}")

# 新增清理GCS文件的任务
cleanup_gcs_files = PythonOperator(
    task_id="cleanup_gcs_files",
    python_callable=keep_only_latest_gcs_file,
    provide_context=True,
    dag=dag
)

# 调整任务依赖:create_transfer_job >> cleanup_gcs_files

方案3:S3事件通知触发Airflow记录最新文件

若允许修改S3配置,可为目标桶添加S3事件通知(如文件上传时触发Lambda或SQS),让Airflow监听事件源,每次有新文件上传时记录最新文件路径,每日执行任务时直接调用该路径进行传输。此方案适合文件数量较多的场景,避免每日枚举所有文件的开销。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 15:35:32