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

