如何将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端点
客户项目侧配置
- 确认Cloud SQL实例的PSC服务端点已创建:进入Cloud SQL实例→
Connections→Private Service Connect,确保已生成服务端点,且允许列表包含租户项目ID(添加租户项目为授权服务项目)。 - 确保共享VPC的子网配置允许跨项目访问:共享VPC的宿主项目需将租户项目设为服务项目,且子网的防火墙规则允许TCP 5432端口(PostgreSQL默认端口)的入站/出站流量。
租户项目侧配置
- 创建Composer v3环境时,指定共享VPC的子网(需确保租户项目有该子网的使用权限)。
- 在租户项目中创建PSC连接端点,指向客户项目的Cloud SQL PSC服务端点:
- 进入VPC网络→
Private Service Connect→连接端点,新建端点,选择共享VPC子网,指定客户项目的PSC服务端点URL。
- 进入VPC网络→
第二步:跨项目IAM权限配置
租户项目的Airflow服务账号(默认格式:
composer-{环境ID}@tenant-project.iam.gserviceaccount.com)需在客户项目中被授予以下角色:roles/cloudsql.client:允许访问Cloud SQL实例roles/iam.serviceAccountTokenCreator:若使用Cloud SQL Python Connector,需该角色生成认证令牌- 若使用IAM数据库认证,需额外授予
roles/cloudsql.instanceUser
客户项目的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
问题排查步骤
如果上述方案仍有问题,按以下顺序排查:
- 网络连通性:在租户项目的Composer工作节点中,尝试ping Cloud SQL的PSC端点IP,或用
telnet {PSC_IP} 5432测试端口是否可达。 - 权限验证:在租户项目中用Airflow服务账号执行
gcloud sql instances describe 客户项目ID:区域:实例名称,确认能获取实例信息。 - 日志分析:查看Composer工作节点的日志,搜索
cloud-sql-connector或cloud_sql_proxy相关错误,定位具体问题(如网络超时、权限拒接)。
内容的提问来源于stack exchange,提问作者Henry Xiloj Herrera
相关产品推荐
相关产品推荐

