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

如何在Dataflow作业中通过Cloud SQL Proxy连接GCP Cloud SQL(Java/Python)

在Dataflow作业中通过Cloud SQL Proxy连接Cloud SQL(Java/Python实现)

Java 实现方式

Dataflow Java SDK可通过初始化脚本,让每个工作节点启动时自动下载并运行Cloud SQL Proxy。

1. 编写启动代理的Shell脚本

创建start_cloud_sql_proxy.sh脚本,内容如下:

#!/bin/bash
# 下载Cloud SQL Proxy二进制文件
wget https://dl.google.com/cloudsql/cloud_sql_proxy.linux.amd64 -O cloud_sql_proxy
chmod +x cloud_sql_proxy
# 启动代理,替换为你的实例连接名
./cloud_sql_proxy -instances=PROJECT_ID:REGION:INSTANCE_NAME=tcp:5432 &

将该脚本上传到GCS存储桶中。

2. 在Dataflow作业中配置初始化动作

在代码中通过DataflowPipelineOptions指定初始化脚本的GCS路径:

import org.apache.beam.runners.dataflow.options.DataflowPipelineOptions;
import org.apache.beam.sdk.options.PipelineOptions;
import org.apache.beam.sdk.options.PipelineOptionsFactory;
import org.apache.beam.sdk.Pipeline;

public class CloudSqlDataflowJob {
    public static void main(String[] args) {
        PipelineOptions options = PipelineOptionsFactory.create();
        DataflowPipelineOptions dataflowOpts = options.as(DataflowPipelineOptions.class);
        
        // 设置初始化脚本的GCS路径
        dataflowOpts.setSetupFile("gs://your-bucket/path/to/start_cloud_sql_proxy.sh");
        // 配置其他Dataflow参数(项目、区域等)
        dataflowOpts.setProject("your-project-id");
        dataflowOpts.setRegion("your-region");
        
        Pipeline pipeline = Pipeline.create(dataflowOpts);
        
        // 后续编写Dataflow业务逻辑...
        // 连接数据库示例(PostgreSQL为例)
        String jdbcUrl = "jdbc:postgresql://localhost:5432/your-db?user=db-user&password=db-pass";
        // 执行数据库操作...
    }
}

3. 连接数据库

代理启动后,直接通过localhost:5432(PostgreSQL默认端口)连接Cloud SQL,MySQL则改为3306。


Python 实现方式

Python SDK支持通过setup_command直接指定启动命令,或通过自定义setup.py完成初始化。

方法1:直接通过setup_command配置启动命令

在代码中直接设置节点初始化时执行的命令:

import apache_beam as beam
from apache_beam.options.pipeline_options import PipelineOptions, SetupOptions

def process_data(element):
    # 连接数据库示例(psycopg2操作PostgreSQL)
    import psycopg2
    conn = psycopg2.connect(
        host='localhost',
        port=5432,
        dbname='your-db',
        user='db-user',
        password='db-pass'
    )
    # 执行数据库读写操作...
    conn.close()
    return element

if __name__ == "__main__":
    options = PipelineOptions()
    setup_opts = options.view_as(SetupOptions)
    
    # 设置启动Cloud SQL Proxy的命令,替换实例连接名
    setup_opts.setup_command = """
    wget https://dl.google.com/cloudsql/cloud_sql_proxy.linux.amd64 -O cloud_sql_proxy && 
    chmod +x cloud_sql_proxy && 
    ./cloud_sql_proxy -instances=PROJECT_ID:REGION:INSTANCE_NAME=tcp:5432 &
    """
    
    # 配置Dataflow基础参数
    options.view_as(beam.options.pipeline_options.GoogleCloudOptions).project = "your-project-id"
    options.view_as(beam.options.pipeline_options.GoogleCloudOptions).region = "your-region"
    
    with beam.Pipeline(options=options) as p:
        (p
         | beam.Create([1, 2, 3])
         | beam.Map(process_data)
        )

方法2:使用自定义setup.py

如果需要依赖安装+代理启动的组合逻辑,可编写setup.py:

from setuptools import setup

setup(
    name="dataflow-cloudsql-setup",
    version="0.0.1",
    install_requires=[
        "apache-beam[gcp]",
        "psycopg2-binary"
    ],
    scripts=["start_proxy.sh"]
)

start_proxy.sh内容与Java部分的脚本一致,上传到GCS后,在代码中指定:

setup_opts.setup_file = "gs://your-bucket/path/to/setup.py"

通用注意事项

  • 确保Dataflow工作节点的服务账号拥有Cloud SQL Client角色,具备访问Cloud SQL实例的权限。
  • 代理启动命令后的&用于让进程后台运行,避免阻塞节点初始化流程。
  • 根据数据库类型调整端口:PostgreSQL用5432,MySQL用3306。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 20:14:59