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

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
);
排查方向与修复建议

按优先级从高到低逐一验证修复:

  • 修复代码基础语法错误
    1. create_row_data函数入参定义为id/name/role/joining_date/last_updated/is_deleted,但内部赋值使用了未定义的column1~column6变量,运行时会触发NameError,直接修改内部赋值映射到入参即可。
    2. 检查所有字符串引号为英文半角格式,确认函数调用时的参数引号、括号完全闭合,避免语法解析异常。
  • 删除默认流不支持的冗余逻辑
    _default默认流是BigQuery预创建的流,写入即生效,不需要手动执行create_write_stream创建流、finalize_write_stream封流、batch_commit_write_streams提交操作,这几个接口仅对用户手动创建的PENDING/BUFFERED类型流生效,对默认流调用会触发资源不存在错误,直接删除这部分冗余代码即可。
  • 排查认证与网络连通性
    从DEBUG日志可见程序反复尝试连接元数据服务metadata.google.internal,说明当前运行环境的认证/网络存在问题:
    1. 如果是本地环境运行,先执行gcloud auth application-default login完成应用默认凭证认证,确认使用的账号拥有目标BigQuery表的写入权限。
    2. 确认运行环境到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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.28 22:12:18