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

