如何将Web请求JSON参数传入BigQuery算子的SQL语句
问题
通过Google Cloud Function触发Airflow DAG时,在请求体conf中传入了file_name参数,目前可以在PythonOperator中通过kwargs["params"]读取到该参数,需要在不依赖XCom的前提下,直接将该参数传入BigQueryExecuteQueryOperator的SQL语句中。
触发DAG的参考代码:
endpoint = f"api/v1/dags/{dag_id}/dagRuns" request_url = f"{web_server_url}/{endpoint}" json_data = {"conf": {"file_name": file_name}} response = make_composer2_web_server_request( request_url, method="POST", json=json_data )
PythonOperator读取参数的参考实现:
def printParams(**kwargs): for k, v in kwargs["params"].items(): print(f"{k}:{v}") task = PythonOperator( task_id="print_parameters", python_callable=printParams)
待改造的BigQuery算子代码:
perform_analytics = BigQueryExecuteQueryOperator( task_id="perform_analytics", sql=f""" Select **passFileNameVariableFromJsonDataHere** as FileName FROM my_project.sample.table1 """, destination_dataset_table=f"my_project.sample.table2", write_disposition="WRITE_APPEND", use_legacy_sql=False )
解决方案
完全不需要使用XCom,BigQueryExecuteQueryOperator的sql参数原生支持Jinja模板渲染,可以直接在SQL中通过模板变量读取DAG运行时传入的参数,无额外依赖。
有两种可直接落地的实现方式:
- 方式一:通过params对象读取
和PythonOperator中读取参数的逻辑一致,直接在SQL中使用{{ params.file_name }}作为占位符即可,修改后的算子代码如下:
perform_analytics = BigQueryExecuteQueryOperator( task_id="perform_analytics", sql=""" Select {{ params.file_name }} as FileName FROM my_project.sample.table1 """, destination_dataset_table="my_project.sample.table2", write_disposition="WRITE_APPEND", use_legacy_sql=False )
注意:不要给sql字符串加f前缀,否则会在DAG文件解析阶段就执行字符串替换,此时任务尚未触发,拿不到运行时传入的参数,会直接报错。保留普通字符串格式,交给Airflow在任务执行前完成模板渲染即可。
- 方式二:直接读取dag_run原始配置
也可以跳过params映射,直接在模板中读取DagRun的原始conf配置,占位符写{{ dag_run.conf.get('file_name') }}。这种方式灵活性更高,还支持设置默认值避免参数缺失报错,比如写为{{ dag_run.conf.get('file_name', 'default_file_name') }},就算触发时没传file_name参数,也不会因为模板渲染失败中断任务。
所有标注了template_fields的Airflow算子都支持这种运行时参数读取方式,不需要前置任务中转,也不需要操作XCom,参数会在任务实例执行前的渲染阶段自动替换为实际传入的值。
内容的提问来源于stack exchange,提问作者Ruslan Lomov
相关产品推荐
相关产品推荐

