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

Airflow实现S3桶文件列表与跨桶复制失败问题求助

解决Airflow中动态创建S3复制任务不执行的问题

问题场景

我是Airflow新手,尝试编写DAG实现列出S3源桶所有文件并逐个复制到目标桶的功能,代码能正常遍历文件,但复制操作始终不执行,日志显示已获取文件列表但无复制动作。

原代码如下:

from airflow.models import DAG
from airflow.decorators import task
from datetime import datetime
from airflow.models import Variable
import logging
from airflow.providers.amazon.aws.operators.s3 import S3ListOperator
from airflow.providers.amazon.aws.operators.s3 import S3CopyObjectOperator
from airflow.operators.dummy import DummyOperator


default_args = {
    'owner': 'airflow',
    'start_date': datetime(2023, 2, 16),
    'email_on_failure': False,
    'email_on_success': False,
    'email_on_retry': False,
    'schedule': "@daily"
}

dag = DAG(
    dag_id='myFirstDag',
    start_date=datetime(2023, 5, 15),
    default_args= default_args,
    catchup=False
)    

@dag.task()
def print_objects(objects):
    print("All Keys", objects)
    last_task = None
    for key in objects:
        print("Current key", key)
        s3Copy = S3CopyObjectOperator(
        task_id= key,
        source_bucket_key=key,
        dest_bucket_key=key,
        source_bucket_name="s3-bukcet-for-airflow-in",
        dest_bucket_name="s3-bukcet-for-airflow-out",
        aws_conn_id="vivek_aws",
        dag=dag
        )
        if last_task:
            last_task >> s3Copy
        last_task = s3Copy               

list_bucket = S3ListOperator(
    task_id='list_files_in_bucket',
    bucket='s3-bukcet-for-airflow-in',
    aws_conn_id='vivek_aws'
)
print_objects(list_bucket.output)

问题原因

你在@dag.task()装饰的print_objects函数里实例化S3CopyObjectOperator是错误的:Airflow的任务(Operator)是在DAG解析阶段被定义并加入任务流的,而@task装饰的函数是在DAG运行阶段才执行的代码,运行阶段创建的Operator不会被Airflow调度执行,只会作为普通Python对象处理,自然不会触发复制操作。

修正方案

使用Airflow 2.2+支持的**动态任务映射(Dynamic Task Mapping)**来实现动态生成复制任务,这是官方推荐的动态任务实现方式。

修正后的完整代码:

from airflow.models import DAG
from datetime import datetime
from airflow.providers.amazon.aws.operators.s3 import S3ListOperator
from airflow.providers.amazon.aws.operators.s3 import S3CopyObjectOperator

default_args = {
    'owner': 'airflow',
    'start_date': datetime(2023, 2, 16),
    'email_on_failure': False,
    'email_on_success': False,
    'email_on_retry': False,
    'schedule': "@daily"
}

with DAG(
    dag_id='myFirstDag',
    start_date=datetime(2023, 5, 15),
    default_args=default_args,
    catchup=False
) as dag:

    list_bucket = S3ListOperator(
        task_id='list_files_in_bucket',
        bucket='s3-bukcet-for-airflow-in',
        aws_conn_id='vivek_aws'
    )

    # 动态生成复制任务:为每个文件创建一个独立的复制任务
    copy_files = S3CopyObjectOperator.partial(
        task_id='copy_s3_file',
        source_bucket_name='s3-bukcet-for-airflow-in',
        dest_bucket_name='s3-bukcet-for-airflow-out',
        aws_conn_id='vivek_aws'
    ).expand(
        source_bucket_key=list_bucket.output,
        dest_bucket_key=list_bucket.output
    )

    # 设置任务依赖
    list_bucket >> copy_files

关键修正说明

  • 用S3CopyObjectOperator.partial()定义任务的固定参数,再通过.expand()根据文件列表动态生成多个任务实例
  • 任务依赖直接通过list_bucket >> copy_files设置,Airflow会自动处理动态任务的调度逻辑
  • 移除原代码中在@task内创建Operator的错误逻辑,避免运行阶段无法调度任务的问题
  • 自动生成的任务会在task_id后添加后缀(如copy_s3_file__0),避免文件名含特殊字符导致的task_id无效问题

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 14:12:52