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

Python SDK下Dataflow作业连接CloudSQL的故障排查求助

Dataflow连接CloudSQL运行异常解决方案

核心问题排查方向

1. 网络连通性配置

  • 若CloudSQL使用公共IP:需为Dataflow工作节点所属VPC配置Cloud NAT,保障节点可访问公网连接CloudSQL
  • 若CloudSQL使用私有IP:需将Dataflow工作节点部署在与CloudSQL相同的VPC网络下,同时完成VPC peering配置保障内网连通
  • 无需额外配置端口,cloud-sql-python-connector会自动处理连接路由逻辑

2. 运行时依赖配置

当前requirements.txt缺失必要依赖,需补充完整,保证worker节点运行时可自动安装所有需要的包:

apache-beam[gcp]==2.40.0
cloud-sql-python-connector[pymysql]==0.4.0
pymysql==1.0.2

注意:requirements.txt中的apache-beam版本需与本地打包使用的版本完全一致,避免版本不兼容报错。

模板创建命令需修正路径分隔符,正确命令如下:

python3 -m template --runner DataflowRunner \
                    --project project_name \
                    --staging_location gs://bucket_name/folder/staging \
                    --temp_location gs://bucket_name/folder/temp \
                    --template_location gs://bucket_name/folder/templates/template-df \
                    --region europe-west1 \
                    --requirements_file requirements.txt

无需提前在Cloud Shell手动安装依赖,只要requirements.txt配置正确,运行时worker会自动拉取安装。

3. 权限配置

为Dataflow运行时默认服务账号(格式为[项目号]-compute@developer.gserviceaccount.com)授予以下权限:

  • Cloud SQL Client角色(读操作必要)
  • Cloud SQL Editor角色(写操作必要)
  • 若使用CloudSQL公共IP,需在CloudSQL实例的授权网络配置中放行Dataflow节点的出口IP段。

4. 代码优化

原代码存在连接重复创建、无超时配置、无异常捕获、无资源释放逻辑的问题,优化后的DoFn代码如下:

import logging
import apache_beam as beam
from google.cloud.sql.connector import connector

class ReadSQLTable(beam.DoFn):
    def __init__(self, hostaddr, driver, username, password, dbname):
        super().__init__()
        self.hostaddr = hostaddr
        self.driver = driver
        self.username = username
        self.password = password
        self.dbname = dbname
        self.conn = None

    def setup(self):
        # 连接逻辑放在setup方法,每个worker仅初始化一次,避免process重复创建连接
        try:
            self.conn = connector.connect(
                self.hostaddr,
                self.driver,
                user=self.username,
                password=self.password,
                db=self.dbname,
                connect_timeout=30
            )
        except Exception as e:
            logging.error(f"CloudSQL连接失败: {str(e)}", exc_info=True)
            raise e

    def process(self, element):
        if not self.conn:
            raise Exception("数据库连接未初始化成功")
        try:
            cursor = self.conn.cursor()
            cursor.execute("SELECT * from table_name")
            result = cursor.fetchall()
            for row in result:
                yield row
        except Exception as e:
            logging.error(f"SQL执行失败: {str(e)}", exc_info=True)
            raise e
        finally:
            cursor.close()

    def teardown(self):
        if self.conn:
            self.conn.close()

5. 日志排查

若作业仍卡住无输出,可到Cloud Logging中筛选对应作业的worker日志,筛选规则如下:

resource.type="dataflow_step"
resource.labels.job_id="[你的Dataflow作业ID]"
severity>=ERROR

即可看到具体的报错信息,用于定位问题。


内容的提问来源于stack exchange,提问作者jbra95

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.05 04:57:04