如何将Cloud Composer DAG变量传递给HTTPS触发的Cloud Function
我使用Cloud Composer的DAG任务调用HTTPS触发的Cloud Function发送邮件,基础功能可正常运行。但尝试将DAG内定义的变量通过GET请求的URL参数传递给Cloud Function时,出现「Unauthorized」权限错误——尽管已为服务账号配置了Cloud Function调用权限,推测是带参数的URL导致认证异常。最终通过改用POST请求,将变量放入JSON请求体传递,同时使用纯净的触发URL获取认证凭据,成功解决问题。
最初的DAG代码尝试
# -------------------------------------------------------------------------------- # 导入库 # -------------------------------------------------------------------------------- import datetime from airflow.models import DAG from airflow.operators.dummy_operator import DummyOperator from airflow.operators.python_operator import PythonOperator from airflow.contrib.operators.bigquery_operator import BigQueryOperator from airflow.providers.google.cloud.operators.bigquery import BigQueryInsertJobOperator,BigQueryExecuteQueryOperator from airflow.providers.google.common.utils import id_token_credentials as id_token_credential_utils import google.auth.transport.requests from google.auth.transport.requests import AuthorizedSession # -------------------------------------------------------------------------------- # 设置变量 # -------------------------------------------------------------------------------- (...) report_name_url = "report_name_url" end_user = "end_user@email.com" # -------------------------------------------------------------------------------- # 函数定义 # -------------------------------------------------------------------------------- def invoke_cloud_function(): # 将字符串添加到URL的?后,作为参数传递给Cloud Function url = "https://<trigger_url>?report_name_url={}&end_user={}".format(report_name_url, end_user) # 用于获取凭据的请求 request = google.auth.transport.requests.Request() # 如果Cloud Function URL带查询参数,需移除后再传入受众参数 id_token_credentials = id_token_credential_utils.get_default_id_token_credentials(url, request=request) # 使用授权会话访问Cloud Function resp = AuthorizedSession(id_token_credentials).request("GET", url=url) print(resp.status_code) # 应返回200 print(resp.content) # HTTP响应体 # -------------------------------------------------------------------------------- # 定义DAG # -------------------------------------------------------------------------------- with DAG( dag_id, schedule_interval= '0 13 05 * *', # DAG Cron调度器 default_args = default_args) as dag: (...) send_email = PythonOperator( task_id="send_email", python_callable=invoke_cloud_function ) start >> run_stored_procedure >> composer_logging >> send_email >> end
最初的Cloud Function代码
def send_email(request): import ssl from email.message import EmailMessage import smtplib import os report_name_url = request.args.get('report_name_url') report_name = report_name_url.replace("_", " ") end_user = request.args.get('end_user') (...) context = ssl.create_default_context() with smtplib.SMTP_SSL('smtp.gmail.com', 465, context=context) as smtp: smtp.login(sender_email, password) smtp.sendmail(sender_email, receiver_email, em.as_string())
日志错误信息
"(...)Unauthorized</h1> <h2>Your client does not have permission to the requested URL <code>...</code>.</h2> <h2></h2> </body></html> '"
最终解决方案代码
DAG中的调用函数
def invoke_cloud_function_success(): url = "<trigger url>" # URL同时也是目标受众 # 用于获取凭据的请求 request = google.auth.transport.requests.Request() # 如果Cloud Function URL带查询参数,需移除后再传入受众参数 id_token_credentials = id_token_credential_utils.get_default_id_token_credentials(url, request=request) headers = {"Content-Type": "application/json"} body = {"report_name":report_name, "end_user":end_user, "datastudio_link":datastudio_link} # 使用授权会话通过POST请求访问Cloud Function resp = AuthorizedSession(id_token_credentials).post(url=url, json=body, headers=headers) print(resp.status_code) # 应返回200 print(resp.content) # HTTP响应体
Cloud Function中的处理代码
request_json = request.get_json() report_name = list(request_json['report_name']) datastudio_link = list(request_json['datastudio_link']) end_user = list(request_json['end_user'])
内容的提问来源于stack exchange,提问作者Filipe Pereira
相关产品推荐
相关产品推荐

