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

如何将Composer 3中的Airflow 2通过Private Service Connect连接到Cloud SQL?

解决Composer v3(Airflow 2)跨项目访问PSC暴露的Cloud SQL PostgreSQL问题

核心前提梳理

你的场景核心是Composer v3租户项目中的Airflow资源(调度器/工作节点)要访问客户项目中通过Private Service Connect(PSC)暴露在共享VPC的Cloud SQL实例,需要同时打通网络路径和配置跨项目权限。


第一步:网络层面打通租户项目到客户项目的PSC端点

客户项目侧配置

  1. 确认Cloud SQL实例的PSC服务端点已创建:进入Cloud SQL实例→Connections→Private Service Connect,确保已生成服务端点,且允许列表包含租户项目ID(添加租户项目为授权服务项目)。
  2. 确保共享VPC的子网配置允许跨项目访问:共享VPC的宿主项目需将租户项目设为服务项目,且子网的防火墙规则允许TCP 5432端口(PostgreSQL默认端口)的入站/出站流量。

租户项目侧配置

  1. 创建Composer v3环境时,指定共享VPC的子网(需确保租户项目有该子网的使用权限)。
  2. 在租户项目中创建PSC连接端点,指向客户项目的Cloud SQL PSC服务端点:
    • 进入VPC网络→Private Service Connect→连接端点,新建端点,选择共享VPC子网,指定客户项目的PSC服务端点URL。

第二步:跨项目IAM权限配置

  1. 租户项目的Airflow服务账号(默认格式:composer-{环境ID}@tenant-project.iam.gserviceaccount.com)需在客户项目中被授予以下角色:

    • roles/cloudsql.client:允许访问Cloud SQL实例
    • roles/iam.serviceAccountTokenCreator:若使用Cloud SQL Python Connector,需该角色生成认证令牌
    • 若使用IAM数据库认证,需额外授予roles/cloudsql.instanceUser
  2. 客户项目的Cloud SQL实例IAM配置中,添加租户项目的Airflow服务账号为授权用户。


第三步:推荐的两种连接实现方案

方案1:Cloud SQL Python Connector(解决之前无限挂起问题)

挂起通常是网络未打通或权限不足导致Connector无法获取PSC端点信息,以下是修正后的代码示例:

首先确保Composer环境安装依赖包:在Composer环境的PyPI包中添加cloud-sql-python-connector和psycopg2-binary。

from airflow import DAG
from airflow.operators.python import PythonOperator
from google.cloud.sql.connector import Connector, IPTypes
import sqlalchemy
from datetime import datetime

def test_cloud_sql_connection():
    # 替换为你的客户项目Cloud SQL信息
    INSTANCE_CONNECTION_NAME = "客户项目ID:区域:实例名称"
    DB_USER = "postgres"
    DB_PASS = "你的数据库密码"  # 若用IAM认证可省略,添加enable_iam_auth=True
    DB_NAME = "目标数据库名"

    # 初始化Connector,指定私有IP访问(对应PSC)
    connector = Connector()

    def get_db_connection():
        conn = connector.connect(
            INSTANCE_CONNECTION_NAME,
            "psycopg2",
            user=DB_USER,
            password=DB_PASS,
            db=DB_NAME,
            ip_type=IPTypes.PRIVATE,  # 关键:强制使用私有IP访问PSC端点
            # enable_iam_auth=True  # 若使用IAM数据库认证,取消注释
        )
        return conn

    # 创建SQLAlchemy引擎并测试连接
    engine = sqlalchemy.create_engine(
        "postgresql+psycopg2://",
        creator=get_db_connection,
    )

    with engine.connect() as conn:
        result = conn.execute(sqlalchemy.text("SELECT version();"))
        print(f"连接成功,PostgreSQL版本:{result.fetchone()[0]}")
    
    connector.close()

with DAG(
    dag_id="cloud_sql_psc_python_connector",
    start_date=datetime(2024, 1, 1),
    schedule_interval=None,
    catchup=False,
) as dag:
    connect_task = PythonOperator(
        task_id="test_psc_connection",
        python_callable=test_cloud_sql_connection,
    )

方案2:Cloud SQL Proxy(解决连接被意外关闭问题)

之前的错误是因为Proxy未指定使用私有IP访问PSC端点,推荐用KubernetesPodOperator在Composer的K8s集群中运行Proxy:

from airflow import DAG
from airflow.providers.cncf.kubernetes.operators.kubernetes_pod import KubernetesPodOperator
from datetime import datetime

with DAG(
    dag_id="cloud_sql_psc_proxy",
    start_date=datetime(2024, 1, 1),
    schedule_interval=None,
    catchup=False,
) as dag:
    # 运行Cloud SQL Proxy,指定私有IP访问
    run_proxy = KubernetesPodOperator(
        task_id="run_cloud_sql_proxy",
        name="cloud-sql-proxy",
        image="gcr.io/cloudsql-docker/gce-proxy:1.33.0",
        cmds=[
            "/cloud_sql_proxy",
            "-instances=客户项目ID:区域:实例名称=tcp:5432",
            "-ip_address_types=PRIVATE",  # 关键:强制使用私有IP访问PSC
        ],
        volume_mounts=[
            {
                "name": "gcp-service-account",
                "mountPath": "/var/secrets/google",
                "readOnly": True,
            }
        ],
        volumes=[
            {
                "name": "gcp-service-account",
                "secret": {
                    "secretName": "airflow-service-account-secret",  # 租户项目中存储Airflow服务账号密钥的Secret
                }
            }
        ],
        env_vars={"GOOGLE_APPLICATION_CREDENTIALS": "/var/secrets/google/key.json"},
        namespace="default",
        get_logs=True,
        is_delete_operator_pod=True,
    )

    # 测试数据库连接
    test_connection = KubernetesPodOperator(
        task_id="test_sql_connection",
        name="test-sql-client",
        image="python:3.9-slim",
        cmds=["pip", "install", "-q", "psycopg2-binary", "&&", "python", "-c"],
        arguments=[
            """
            import psycopg2
            try:
                conn = psycopg2.connect(
                    host="localhost",
                    port=5432,
                    user="postgres",
                    password="你的数据库密码",
                    dbname="目标数据库名"
                )
                cur = conn.cursor()
                cur.execute("SELECT version();")
                print(f"连接成功:{cur.fetchone()[0]}")
                cur.close()
                conn.close()
            except Exception as e:
                print(f"连接失败:{str(e)}")
                raise
            """
        ],
        depends_on_past=False,
        get_logs=True,
        is_delete_operator_pod=True,
        namespace="default",
    )

    run_proxy >> test_connection

问题排查步骤

如果上述方案仍有问题,按以下顺序排查:

  1. 网络连通性:在租户项目的Composer工作节点中,尝试ping Cloud SQL的PSC端点IP,或用telnet {PSC_IP} 5432测试端口是否可达。
  2. 权限验证:在租户项目中用Airflow服务账号执行gcloud sql instances describe 客户项目ID:区域:实例名称,确认能获取实例信息。
  3. 日志分析:查看Composer工作节点的日志,搜索cloud-sql-connector或cloud_sql_proxy相关错误,定位具体问题(如网络超时、权限拒接)。

内容的提问来源于stack exchange,提问作者Henry Xiloj Herrera

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 09:59:57