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

Airflow使用AWSAthenaOperator执行动态查询列表报错如何解决

问题根因

Airflow的DAG结构是在DagBag加载阶段(也就是调度器扫描DAG文件的时候)就已经确定完成的,此时所有DAG定义代码都会被执行,但前序任务还没有实际运行,read_date_information_file.output只是一个XCom的占位引用,不是真实的返回列表。你直接用enumerate遍历这个占位符的时候,会触发Airflow尝试解析这个还不存在的XCom值,最终导致加载超时。
硬编码的列表可以正常运行,是因为DagBag加载阶段就能确定列表长度和内容,可以直接生成固定数量的Athena任务。

解决方案1:使用动态任务映射(推荐,Airflow 2.3+支持)

动态任务映射是Airflow官方推出的专门用来处理「根据前序任务的输出动态生成多个下游任务」的场景,不需要手动遍历列表,Airflow会在运行时根据前序任务返回的列表自动生成对应数量的任务。
修改后的代码示例:

def get_date_information(ti):
    import boto3
    s3 = boto3.client('s3')
    data = s3.get_object(Bucket=output_bucket, Key=key)
    contents = data['Body'].read().decode("utf-8")
    print('Date information is: ', contents)
    events_list = contents.split(',')
    return events_list

with DAG(
    dag_id='adserver_split_job_emr_job_dag',
    default_args={
        'owner': 'adserver_airflow',
        'depends_on_past': False,
        'email': ['airflow@example.com'],
        'email_on_failure': False,
        'email_on_retry': False,
    },
    dagrun_timeout=timedelta(hours=2),
    start_date=datetime(2021, 9, 22, 9),
    schedule_interval='20 * * * *',
    # Airflow 2.4以下版本需要加这个参数开启动态任务映射
    # render_template_as_native_obj=True
) as dag:

    read_date_information_file = PythonOperator(
        task_id="read_date_information_file",
        python_callable=get_date_information
    )

    # partial定义所有任务通用的参数,expand指定需要动态遍历的参数
    run_query = AWSAthenaOperator.partial(
        task_id='run_query',
        output_location=config.ATHENA_OUTPUT_LOCATION,
        database=config.ATHENA_DATABASE_NAME,
        aws_conn_id='aws_default'
    ).expand(query=read_date_information_file.output)

    read_date_information_file >> run_query

Airflow会自动给每个query生成独立的任务,任务ID会自动加上序号后缀,支持单个任务的重试、日志查看等能力。

解决方案2:老版本Airflow兼容方案

如果你的Airflow版本低于2.3,没法使用动态任务映射,可以把Athena查询的执行逻辑封装到一个Python函数里,在函数内部遍历查询列表挨个执行:

def run_athena_queries(ti):
    import boto3
    from pyathena import connect
    query_list = ti.xcom_pull(task_ids='read_date_information_file')
    # 初始化Athena连接
    conn = connect(s3_staging_dir=config.ATHENA_OUTPUT_LOCATION,
                   region_name='你的AWS区域',
                   aws_conn_id='aws_default')
    cursor = conn.cursor()
    for i, event in enumerate(query_list):
        print(f'执行第{i}个查询: {event}')
        cursor.execute(event)
        # 可自行添加等待查询完成、获取结果的逻辑
        print(f'第{i}个查询执行完成')

with DAG(
    dag_id='adserver_split_job_emr_job_dag',
    default_args={
        'owner': 'adserver_airflow',
        'depends_on_past': False,
        'email': ['airflow@example.com'],
        'email_on_failure': False,
        'email_on_retry': False,
    },
    dagrun_timeout=timedelta(hours=2),
    start_date=datetime(2021, 9, 22, 9),
    schedule_interval='20 * * * *',
) as dag:

    read_date_information_file = PythonOperator(
        task_id="read_date_information_file",
        python_callable=get_date_information
    )

    run_all_queries = PythonOperator(
        task_id='run_all_queries',
        python_callable=run_athena_queries
    )

    read_date_information_file >> run_all_queries

注意:该方案的缺点是所有查询都在同一个任务里执行,没法单独重试单个查询,如果需要单个查询的任务级别控制,还是建议升级Airflow版本使用动态任务映射。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 16:36:03