Cloud Composer 3 DAG调用同VPC HTTP云函数遇403 Forbidden错误求助
问题描述
我在Cloud Composer 3中运行DAG,触发同一共享VPC内的HTTP云函数,两者使用相同的自定义服务账号。已按照官方文档配置了主机项目和服务项目的必要权限,但仍遇到错误:HTTP request failed: 403 Client Error: Forbidden for URL。
有趣的是,在Compute Engine中使用以下curl命令触发该云函数却能成功:
curl -m 70 -X POST <Function endpoint> \ -H "Authorization: bearer $(gcloud auth print-identity-token)" \ -H "Content-Type: application/json" \ -d '{ }'
以下是我在DAG中使用的Python代码:
from datetime import datetime, timedelta from airflow import DAG from airflow.operators.python import PythonOperator import requests import subprocess import logging default_args = { 'owner': 'airflow', 'depends_on_past': False, 'start_date': datetime(2024, 7, 2), 'email_on_failure': False, 'email_on_retry': False, 'retries': 1, 'retry_delay': timedelta(minutes=5), } dag = DAG( 'run_cloud_function_daily', default_args=default_args, description='Run HTTP-based Cloud Function daily ', schedule='0 19 * * *', # Cron expression for daily at 7 PM ) # Define the URL of your HTTP-based Cloud Function endpoint cloud_function_url = '' # Function to fetch authorization token using gcloud and call the Cloud Function def call_cloud_function_with_auth_token(): try: # Fetch authorization token using gcloud command token_command = 'gcloud auth print-identity-token' result = subprocess.run(token_command, capture_output=True, text=True, shell=True) authorization_token = result.stdout.strip() # Get the output of the command (token) if authorization_token: print('Token Successfully retrieved',authorization_token) # Prepare headers with authorization token headers = { "Content-Type": "application/json", "Authorization": f"bearer {authorization_token}" } # Prepare payload if needed payload = {} # Replace with your JSON payload if needed # Make HTTP POST request to Cloud Function endpoint response = requests.post(cloud_function_url, headers=headers) # Check response status response.raise_for_status() # Raise an exception for HTTP errors (4xx, 5xx) except subprocess.CalledProcessError as e: logging.error(f"Failed to fetch authorization token: {e}") raise logging.info("HTTP request successful:", response.text) except requests.exceptions.RequestException as e: logging.error(f"HTTP request failed: {e}") raise print(f"Details: {e.response.text if e.response else 'No response details'}") # Define the PythonOperator to execute the function run_cloud_function_task = PythonOperator( task_id='call_cloud_function_with_auth_token', python_callable=call_cloud_function_with_auth_token, dag=dag, ) # Set task dependencies run_cloud_function_task
解决方案
问题根源
DAG中用subprocess.run('gcloud auth print-identity-token')获取的token,并非来自你配置的自定义服务账号,而是Composer默认的工作节点服务账号(比如composer-worker@<project-id>.iam.gserviceaccount.com),这个账号没有调用目标云函数的权限,导致403错误。而你在Compute Engine上测试时,VM默认绑定的是你配置的自定义服务账号,所以能成功触发。
修复步骤
- 替换token获取方式:不要调用gcloud命令,直接用Google官方身份验证库获取当前服务账号的身份令牌,确保使用的是Composer配置的自定义服务账号。
- 修正异常逻辑:原代码中
raise语句后的日志打印永远不会执行,需要调整异常处理的结构。
修改后的完整代码:
from datetime import datetime, timedelta from airflow import DAG from airflow.operators.python import PythonOperator import requests import logging from google.auth import default from google.auth.transport.requests import Request default_args = { 'owner': 'airflow', 'depends_on_past': False, 'start_date': datetime(2024, 7, 2), 'email_on_failure': False, 'email_on_retry': False, 'retries': 1, 'retry_delay': timedelta(minutes=5), } dag = DAG( 'run_cloud_function_daily', default_args=default_args, description='Run HTTP-based Cloud Function daily ', schedule='0 19 * * *', # Cron expression for daily at 7 PM ) # 替换为你的云函数端点 cloud_function_url = 'YOUR_CLOUD_FUNCTION_ENDPOINT' def call_cloud_function_with_auth_token(): try: # 获取当前服务账号的身份令牌 credentials, project_id = default(scopes=['https://www.googleapis.com/auth/cloud-platform']) if not credentials.valid: credentials.refresh(Request()) authorization_token = credentials.id_token if authorization_token: logging.info('Token Successfully retrieved') headers = { "Content-Type": "application/json", "Authorization": f"bearer {authorization_token}" } payload = {} response = requests.post(cloud_function_url, headers=headers, json=payload) response.raise_for_status() logging.info(f"HTTP request successful: {response.text}") except requests.exceptions.RequestException as e: error_details = e.response.text if e.response else 'No response details' logging.error(f"HTTP request failed: {e}, Details: {error_details}") raise except Exception as e: logging.error(f"Unexpected error: {e}") raise run_cloud_function_task = PythonOperator( task_id='call_cloud_function_with_auth_token', python_callable=call_cloud_function_with_auth_token, dag=dag, ) run_cloud_function_task
额外检查项
- 确认自定义服务账号已被授予
roles/cloudfunctions.invoker权限,作用于目标云函数。 - 检查共享VPC的防火墙规则,允许Composer工作节点所在的子网访问云函数所在的VPC网络。
- 验证云函数的触发器设置为“允许内部流量和已认证的外部流量”,或者仅允许内部流量(如果两者在同一VPC)。
内容的提问来源于stack exchange,提问作者Pankaj Goyal
相关产品推荐
相关产品推荐

