Python调用BigQuery Storage Write API默认流写入报流打开失败错误
问题背景
编写Python脚本,计划使用BigQuery Storage Write API的默认流流式写入能力,将Pub/Sub中的数据加载到BigQuery中。参考官方示例适配自身业务逻辑,已按官方文档要求将数据转换为Python客户端支持的ProtoBuf格式,但运行程序时持续抛出如下错误:google.api_core.exceptions.Unknown: None There was a problem opening the stream. Try turning on DEBUG level logs to see the error
已收集的报错信息
开启DEBUG级别日志后,完整运行报错如下:
(venv) {{MY_COMPUTER}} {{FOLDER_NAME}} % python3 default_Stream.py DEBUG:urllib3.connectionpool:Starting new HTTP connection (1): metadata.google.internal.:80 DEBUG:urllib3.connectionpool:Starting new HTTP connection (1): metadata.google.internal.:80 DEBUG:google.cloud.logging_v2.handlers.transports.background_thread:Background thread started. DEBUG:urllib3.connectionpool:Starting new HTTP connection (1): metadata.google.internal.:80 DEBUG:urllib3.connectionpool:Starting new HTTP connection (1): metadata.google.internal.:80 Traceback (most recent call last) : File "default_Stream.py" line 116, in <module> append_rows_default("{{GCLOUD_PROJECT_NAME}}", "{{BIGQUERY_DATASET_NAME}}", "{{BIGQUERY_TABLE}}") File "default_Stream.py", line 95, in append_rows_default response_future_1 = append_rows_stream.send(request) File "{{VIRTUAL_ENVIRONMENT_PATH}}/venv/lib/python3.7/site-packages/google/cloud/bigquery_storage_v1/writer.py", line 234, in send return self._open(request) File "{{VIRTUAL_ENVIRONMENT_PATH}}/venv/lib/python3.7/site-packages/google/cloud/bigquery_storage_v1/writer.py", line 207, in _open raise request_exception google.api_core.exceptions.Unknown: None There was a problem opening the stream. Try turning on DEBUG level logs to see the error. Waiting up to 5 seconds. Sent all pending logs.
现有代码与配置
业务实现脚本
# [START bigquerystorage_append_rows_default] """ This code sample demonstrates how to write records using the low-level generated client for Python. """ from xmlrpc.client import boolean from google.cloud import bigquery_storage_v1 from google.cloud.bigquery_storage_v1 import types from google.cloud.bigquery_storage_v1 import writer from google.protobuf import descriptor_pb2 import logging import google.cloud.logging # 若更新proto定义,执行 protoc --python_out=. 对应proto文件 生成pb2模块 import debezium_record_pb2 def create_row_data(id: int, name: str, role: int, joining_date: int, last_updated: int, is_deleted: boolean): row = debezium_record_pb2.DebeziumRecord() row.column1 = column1 row.column2 = column2 row.column3 = column3 row.column4 = column4 row.column5 = column5 row.column6 = column6 return row.SerializeToString() def append_rows_default(project_id: str, dataset_id: str, table_id: str): """Create a write stream, write some sample data, and commit the stream.""" client = google.cloud.logging.Client() logging.basicConfig(level=logging.DEBUG) client.setup_logging() write_client = bigquery_storage_v1.BigQueryWriteClient() parent = write_client.table_path(project_id, dataset_id, table_id) stream_name = f'{parent}/_default' write_stream = types.WriteStream() # 预留的PENDING流逻辑,当前已注释 #write_stream.type_ = types.WriteStream.Type.PENDING # write_stream = write_client.create_write_stream( # parent=parent, write_stream=write_stream # ) #stream_name = write_stream.name # 构造首次请求模板 request_template = types.AppendRowsRequest() request_template.write_stream = stream_name # 传入Proto序列化Schema proto_schema = types.ProtoSchema() proto_descriptor = descriptor_pb2.DescriptorProto() debezium_record_pb2.DebeziumRecord.DESCRIPTOR.CopyToProto(proto_descriptor) proto_schema.proto_descriptor = proto_descriptor proto_data = types.AppendRowsRequest.ProtoData() proto_data.writer_schema = proto_schema request_template.proto_rows = proto_data # 初始化流式写入对象 append_rows_stream = writer.AppendRowsStream(write_client, request_template) # 构造测试行数据 proto_rows = types.ProtoRows() proto_rows.serialized_rows.append(create_row_data(8, "E", 13, 1643673600000, 1654556118813, False)) # 构造首批写入请求,初始offset必须为0 request = types.AppendRowsRequest() request.offset = 0 proto_data = types.AppendRowsRequest.ProtoData() proto_data.rows = proto_rows request.proto_rows = proto_data logging.basicConfig(level=logging.DEBUG) response_future_1 = append_rows_stream.send(request) logging.basicConfig(level=logging.DEBUG) print(response_future_1.result()) # 关闭流连接 append_rows_stream.close() # 封流与提交逻辑(当前未注释) write_client.finalize_write_stream(name=write_stream.name) batch_commit_write_streams_request = types.BatchCommitWriteStreamsRequest() batch_commit_write_streams_request.parent = parent batch_commit_write_streams_request.write_streams = [write_stream.name] write_client.batch_commit_write_streams(batch_commit_write_streams_request) print(f"Writes to stream: '{write_stream.name}' have been committed.") if __name__ == "__main__": append_rows_default("{{GCLOUD_PROJECT_NAME}}", "{{BIGQUERY_DATASET_NAME}}", "{{BIGQUERY_TABLE}}") # [END bigquerystorage_append_rows_default]
ProtoBuf定义文件
对应生成debezium_record_pb2.py的proto文件内容:
syntax = "proto3"; // 不可包含BigQuery表中不存在的字段 message DebeziumRecord { uint32 column1 = 1; string column2 = 2; uint32 column3 = 3; uint64 column4 = 4; uint64 column5 = 5; bool column6 = 6; }
BigQuery目标表建表语句
CREATE TABLE `{{GCLOUD_PROJECT_NAME}}.{{BIGQUERY_DATASET_NAME}}.{{BIGQUERY_TABLE}}` ( column1 INT64 NOT NULL, column2 STRING, column3 INT64, column4 INT64 NOT NULL, column5 INT64, column6 BOOL );
排查方向与修复建议
按优先级从高到低逐一验证修复:
- 修复代码基础语法错误
create_row_data函数入参定义为id/name/role/joining_date/last_updated/is_deleted,但内部赋值使用了未定义的column1~column6变量,运行时会触发NameError,直接修改内部赋值映射到入参即可。- 检查所有字符串引号为英文半角格式,确认函数调用时的参数引号、括号完全闭合,避免语法解析异常。
- 删除默认流不支持的冗余逻辑
_default默认流是BigQuery预创建的流,写入即生效,不需要手动执行create_write_stream创建流、finalize_write_stream封流、batch_commit_write_streams提交操作,这几个接口仅对用户手动创建的PENDING/BUFFERED类型流生效,对默认流调用会触发资源不存在错误,直接删除这部分冗余代码即可。 - 排查认证与网络连通性
从DEBUG日志可见程序反复尝试连接元数据服务metadata.google.internal,说明当前运行环境的认证/网络存在问题:- 如果是本地环境运行,先执行
gcloud auth application-default login完成应用默认凭证认证,确认使用的账号拥有目标BigQuery表的写入权限。 - 确认运行环境到BigQuery Storage API端点的网络连通性,本地调试时需配置可正常访问Google服务的网络环境,避免gRPC连接建立失败。
- 如果是本地环境运行,先执行
- 修正ProtoBuf字段类型匹配问题
当前传入的时间戳是13位毫秒级数值,大小远超过ProtoBufuint32类型支持的最大值(4294967295,约10位数字),需要把proto文件中column1、column3的类型从uint32改为int64,和BigQuery表的INT64字段完全匹配,避免值溢出导致数据解析失败。 - 开启gRPC层日志定位根因
如果上述修复后仍报流打开错误,在代码最开头添加如下配置开启gRPC全链路trace日志,可以拿到流建立阶段的具体错误码和错误信息,替代当前的通用Unknown报错:import os os.environ['GRPC_TRACE'] = 'all' os.environ['GRPC_VERBOSITY'] = 'DEBUG'
内容的提问来源于stack exchange,提问作者user18697941
相关产品推荐
相关产品推荐

