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

如何在Airflow DAG中实现WITH..AS语句?

在Airflow DAG中实现Athena的WITH..AS查询

首先纠正你查询语句的小问题:原语句的JOIN缺少关联表,应该修改为:

WITH outer_query AS (SELECT id FROM outer)
SELECT * FROM inner JOIN outer_query ON inner.id = outer_query.id

下面提供两种在Airflow中实现该查询的方法:

方法1:使用AthenaOperator(推荐)

Airflow的Amazon Provider提供了AthenaOperator,专门用于执行Athena查询,无需手动编写SDK调用逻辑,代码更简洁。

实现代码:

from airflow import DAG
from airflow.providers.amazon.aws.operators.athena import AthenaOperator
from datetime import datetime

default_args = {
    'owner': 'airflow',
    'start_date': datetime(2024, 1, 1),
    'retries': 1
}

with DAG(
    'athena_with_cte_dag',
    default_args=default_args,
    schedule_interval='@daily',
    catchup=False
) as dag:
    # 定义包含WITH..AS的查询语句
    athena_query = """
        WITH outer_query AS (SELECT id FROM outer)
        SELECT * FROM inner JOIN outer_query ON inner.id = outer_query.id
    """

    # 创建AthenaOperator任务
    run_athena_cte_query = AthenaOperator(
        task_id='run_athena_cte_query',
        query=athena_query,
        database='your_athena_database',  # 替换为你的Athena数据库名
        output_location='s3://your-bucket/path/to/results/',  # 替换为存储查询结果的S3路径
        aws_conn_id='aws_default'  # 替换为你的Airflow AWS连接ID
    )

    run_athena_cte_query

方法2:使用PythonOperator结合boto3

如果需要自定义查询执行后的逻辑(比如处理结果、日志记录等),可以用PythonOperator调用boto3执行Athena查询。

实现代码:

from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime
import boto3

default_args = {
    'owner': 'airflow',
    'start_date': datetime(2024, 1, 1),
    'retries': 1
}

def run_athena_cte_query():
    athena_client = boto3.client('athena', region_name='your-region')  # 替换为你的AWS区域
    query = """
        WITH outer_query AS (SELECT id FROM outer)
        SELECT * FROM inner JOIN outer_query ON inner.id = outer_query.id
    """

    # 启动查询执行
    response = athena_client.start_query_execution(
        QueryString=query,
        QueryExecutionContext={
            'Database': 'your_athena_database'  # 替换为你的Athena数据库名
        },
        ResultConfiguration={
            'OutputLocation': 's3://your-bucket/path/to/results/'  # 替换为S3结果路径
        }
    )

    # 可选:等待查询完成(同步执行逻辑)
    query_execution_id = response['QueryExecutionId']
    while True:
        status = athena_client.get_query_execution(QueryExecutionId=query_execution_id)['QueryExecution']['Status']['State']
        if status in ['SUCCEEDED', 'FAILED', 'CANCELLED']:
            break
        import time
        time.sleep(5)

    if status == 'SUCCEEDED':
        print(f"查询执行成功,ID: {query_execution_id}")
    else:
        error_msg = athena_client.get_query_execution(QueryExecutionId=query_execution_id)['QueryExecution']['Status']['StateChangeReason']
        raise Exception(f"查询失败: {error_msg}")

with DAG(
    'athena_cte_python_dag',
    default_args=default_args,
    schedule_interval='@daily',
    catchup=False
) as dag:
    run_query_task = PythonOperator(
        task_id='run_athena_cte_query',
        python_callable=run_athena_cte_query
    )

    run_query_task

注意事项:

  • 确保Airflow的AWS连接(aws_default或自定义连接)拥有Athena执行权限和S3读写权限
  • 替换代码中的占位符(数据库名、S3路径、AWS区域、连接ID)为实际值
  • 若使用PythonOperator,可根据需求添加更多自定义逻辑,比如读取查询结果到变量中

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.23 01:03:11