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

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默认绑定的是你配置的自定义服务账号,所以能成功触发。

修复步骤

  1. 替换token获取方式:不要调用gcloud命令,直接用Google官方身份验证库获取当前服务账号的身份令牌,确保使用的是Composer配置的自定义服务账号。
  2. 修正异常逻辑:原代码中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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 09:20:55