Airflow:PythonOperator函数内使用{{ ds }}及传递BigQuery查询结果问题
问题与解决方案
问题描述
目前通过PythonOperator调用自定义函数、借助XCom实现任务间数据传递是可行的,但需要在BigQuery查询语句中使用执行日期宏{{ ds }}时,因语句嵌入函数内部无法被Jinja解析,导致查询试图访问名为some_table__{{ ds }}的表而失败。同时希望了解是否存在无需自定义函数即可传递BigQuery查询结果的方案,要求支持回溯运行、不能使用当前日期。
一、解决PythonOperator中的{{ ds }}解析问题
PythonOperator的自定义函数内部字符串不会被Jinja自动渲染,但可以通过context参数直接获取执行日期,手动替换到SQL语句中:
修改原get_some_values函数:
def get_some_values(**context): # 从context中提取执行日期ds execution_date = context['ds'] hook = BigQueryHook(use_legacy_sql=False) conn = hook.get_conn() cursor = conn.cursor() # 用字符串格式化替换表名中的日期占位符 cursor.execute( f"SELECT value1, value2, value3 FROM some_dataset.some_table__{execution_date}" ) results = cursor.fetchone() # 将结果存入XCom if results is not None: for i, result in enumerate(results): context['ti'].xcom_push(f'value{i+1}', result)
此方式下,回溯运行时会自动使用对应历史执行日期,完全符合需求。
二、无需自定义函数的替代方案
直接使用Airflow原生Operator组合实现,全程无需编写自定义函数:
1. 用BigQueryOperator执行查询并推送结果到XCom
BigQueryOperator默认会将查询结果推送到XCom(key为return_value),只需确保开启do_xcom_push=True:
from airflow.providers.google.cloud.operators.bigquery import BigQueryOperator get_values = BigQueryOperator( task_id='get_values', sql=""" SELECT value1, value2, value3 FROM some_dataset.some_table__{{ ds }} """, use_legacy_sql=False, do_xcom_push=True, dag=dag )
这里的SQL语句在Operator层面会被Jinja正常解析,{{ ds }}会自动替换为执行日期,回溯运行时也会匹配对应历史日期。
2. 用SlackWebhookOperator直接从XCom取值发送消息
通过Jinja模板在message参数中直接拉取前一个任务的XCom结果:
from airflow.providers.slack.operators.slack_webhook import SlackWebhookOperator send_slack = SlackWebhookOperator( task_id='send_slack_message', http_conn_id=SLACK_CONN_ID, webhook_token=slack_webhook_token, channel='#some_channel', # 用Jinja语法从XCom拉取结果,[0][0]代表第一行第一列的值 message="""values returned: {{ ti.xcom_pull(task_ids='get_values')[0][0] }}, {{ ti.xcom_pull(task_ids='get_values')[0][1] }}, {{ ti.xcom_pull(task_ids='get_values')[0][2] }} """, username='airflow', dag=dag )
完整DAG示例
from airflow import DAG from airflow.providers.google.cloud.operators.bigquery import BigQueryOperator from airflow.providers.slack.operators.slack_webhook import SlackWebhookOperator from datetime import datetime default_args = { 'owner': 'airflow', 'start_date': datetime(2023, 1, 1) } dag = DAG( 'test', default_args=default_args, schedule_interval='0 12 * * *', catchup=False, ) get_values = BigQueryOperator( task_id='get_values', sql=""" SELECT value1, value2, value3 FROM some_dataset.some_table__{{ ds }} """, use_legacy_sql=False, do_xcom_push=True, dag=dag ) send_slack = SlackWebhookOperator( task_id='send_slack_message', http_conn_id=SLACK_CONN_ID, webhook_token=slack_webhook_token, channel='#some_channel', message="""values returned: {{ ti.xcom_pull(task_ids='get_values')[0][0] }}, {{ ti.xcom_pull(task_ids='get_values')[0][1] }}, {{ ti.xcom_pull(task_ids='get_values')[0][2] }} """, username='airflow', dag=dag ) get_values >> send_slack
内容的提问来源于stack exchange,提问作者Chrisvdberge
相关产品推荐
相关产品推荐

