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

使用Apache Beam JDBC连接Cloud SQL失败求助

问题描述

我尝试用Python SDK的io.jdbc模块(ReadFromJdbc类)连接Cloud SQL,结合官方文档和Cloud MySQL JDBC连接指南编写了代码,但运行时遇到gRPC连接错误,具体情况如下:

代码实现

import apache_beam as beam
import apache_beam.io.jdbc as jdbc
import typing
import apache_beam.coders as coders
import os  # 原代码遗漏该导入,需补充

from apache_beam.options.pipeline_options import PipelineOptions

pipeline_options = {
    'project': 'project-name',
    'runner': 'DataflowRunner',
    'region': 'europe-central2',
    'staging_location':"gs://temp",
    'temp_location':"gs://temp",
    'template_location':"gs://templates/temp_name"
}
pipeline_options = PipelineOptions.from_dictionary(pipeline_options)


serviceAccount = r'path\to\serviceaccount.json'
os.environ['GOOGLE_APPLICATION_CREDENTIALS'] = serviceAccount

ExampleRow = typing.NamedTuple('ExampleRow',
                               [('id', int), ('migration', str)])
coders.registry.register_coder(ExampleRow, coders.RowCoder)


with beam.Pipeline(options=pipeline_options) as p:
    res = (
        p
        | "Read database list" >> jdbc.ReadFromJdbc(
            table_name='table',
            driver_class_name='com.mysql.jdbc.Driver',
            jdbc_url='jdbc:mysql:///<DATABASE_NAME>?cloudSqlInstance=<INSTANCE_CONNECTION_NAME>&socketFactory=com.google.cloud.sql.mysql.SocketFactory&user=<MYSQL_USER_NAME>&password=<MYSQL_USER_PASSWORD>',
            username='user',
            password='pass',
            query = "select id, migration from db.table;",
            fetch_size=1,
            classpath=["com.google.cloud.sql:mysql-socket-factory-connector-j-8:1.7.2"],
            expansion_service = 'host:6666'
        )
        | "Print results" >> beam.io.WriteToText(r'gs://output/out.csv')
    )

错误信息

grpc._channel._InactiveRpcError: <_InactiveRpcError of RPC that terminated with:
        status = StatusCode.UNAVAILABLE
        details = "failed to connect to all addresses; last error: UNAVAILABLE: ipv4:127.0.0.1:6666: WSA Error"
        debug_error_string = "UNKNOWN:failed to connect to all addresses; last error: UNAVAILABLE: ipv4:127.0.0.1:6666: WSA Error {grpc_status:14, created_time:"2022-12-08T15:43:05.445755053+00:00"}"

已尝试操作

  • 按照文档在WLS2 Python环境中配置了expansion service
  • 将expansion_service替换为WLS2的实际IP(能ping通且可访问该IP上的web服务),但错误依旧

想确认是否存在操作失误?


解决方案

1. 修正Expansion Service的监听配置

WLS2和主机属于隔离网络环境,默认启动的expansion service如果只监听127.0.0.1,主机的Python进程无法访问。启动服务时需指定监听所有网络接口:

java -jar beam-sdks-java-expansion-service-<你的beam版本>.jar --port 0.0.0.0:6666

这样服务会绑定WLS2的所有IP,主机才能通过WLS2的实际IP连接到服务。

2. 修复JDBC配置冲突与错误

  • 重复认证信息:JDBC URL中已经包含user和password参数,同时又单独传入username和password字段,会导致参数冲突,建议删除其中一组(推荐保留URL内的参数,或移除URL中的认证信息改用单独字段传递)。
  • 驱动类不匹配:你使用的com.mysql.jdbc.Driver是旧版驱动类,搭配mysql-socket-factory-connector-j-8(适配Connector/J 8.0),应改为com.mysql.cj.jdbc.Driver。
  • 不合理的fetch_size:fetch_size=1会强制逐行读取数据,性能极低,建议调整为1000左右的合理值。

3. 注意DataflowRunner的网络限制

如果用DataflowRunner运行管道:

  • 本地WLS2的expansion service属于内网地址,Dataflow的云worker节点无法访问。这种情况需将expansion service部署到云环境(如GCE虚拟机),或使用Beam官方提供的公共expansion service(部分组件支持)。
  • 测试阶段可先切换为DirectRunner验证逻辑,避开Dataflow的网络隔离问题。

4. 检查版本兼容性

确保Beam SDK版本与mysql-socket-factory版本兼容,建议使用最新稳定版Beam SDK,避免版本不匹配引发的隐性问题。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 09:25:25