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
相关产品推荐
相关产品推荐

