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

使用Python向BigQuery插入数据时SSLError的解决方法

问题描述

维护YouTube频道数据管道时,尝试将直播聊天数据插入BigQuery表,此前可正常插入,现在相同代码在Cloud Function和Google Colab中均抛出SSLError:

SSLError: HTTPSConnectionPool(host='bigquery.googleapis.com', port=443): Max retries exceeded with url: /bigquery/v2/projects/triple-voyage-377203/datasets/test_youtube_data/tables/stream_data/insertAll?prettyPrint=false (Caused by SSLError(SSLEOFError(8, 'EOF occurred in violation of protocol (_ssl.c:2396)')))

已确认数据结构与BigQuery表匹配,示例数据如下:

{'stream_id': 'https://www.youtube.com/watch?v=rBi1_Ggs39U', 'time_in_seconds': -1837, 'action_type': 'add_chat_item', 'message': ':_MIZUKILazy::_MIZUKILazy::_MIZUKILazy:', 'emotes_id': 'UCjv4bfP_67WLuPheS-Z8Ekg/fd89YfbtJ8_j8wTi2Ki4Ag', 'emotes_name': ':_MIZUKILazy:', 'emotes_is_custom_emoji': True, 'message_id': 'ChwKGkNQdjB1c2FucElBREZiM0R3Z1FkQ0NzQzdn', 'timestamp': 1690097427605763, 'time_text': '-30:37', 'author_name': '倉鼠', 'author_images_url': 'https://yt4.ggpht.com/ytc/AOPolaSpe3eNbZY01DHEOdApwgJOwcIuhGj2VIouEvBaHA', 'author_images_id': 'source', 'author_id': 'UCbpvt7VyEpZwJNfwpf09aFw', 'message_type': 'text_message'}

核心插入代码片段:

from google.cloud import bigquery

def insert_to_bigquery(table_id, data):
    client = bigquery.Client() 
    dataset_id = 'test_youtube_data'
    
    table_ref = client.dataset(dataset_id).table(table_id)
    table = client.get_table(table_ref)
    errors = client.insert_rows_json(
        table=table, 
        json_rows=data,
        ignore_unknown_values=True, 
        skip_invalid_rows=True) 
    # 日志输出逻辑
解决方案

1. 升级BigQuery客户端库

旧版本的google-cloud-bigquery可能存在SSL握手相关的bug,执行升级命令:

pip install --upgrade google-cloud-bigquery

如果是Cloud Function环境,需在requirements.txt中指定最新版本:

google-cloud-bigquery>=3.11.0
functions-framework>=3.0.0

2. 拆分插入批次

单次插入过多行可能导致连接超时或SSL异常,将数据拆分为小批次插入:

def insert_to_bigquery(table_id, data, batch_size=500):
    client = bigquery.Client() 
    dataset_id = 'test_youtube_data'
    
    table_ref = client.dataset(dataset_id).table(table_id)
    table = client.get_table(table_ref)
    
    # 按批次拆分数据并插入
    for i in range(0, len(data), batch_size):
        batch = data[i:i+batch_size]
        errors = client.insert_rows_json(
            table=table, 
            json_rows=batch,
            ignore_unknown_values=True, 
            skip_invalid_rows=True) 
        if errors:
            print(f'批次 {i//batch_size} 插入错误: {errors}')
        else:
            print(f'批次 {i//batch_size} 成功插入 {len(batch)} 行')

3. 自定义客户端SSL配置

调整SSL加密套件或降低安全级别(仅调试用,生产环境谨慎使用):

import ssl
from google.cloud import bigquery
from requests.adapters import HTTPAdapter
from urllib3.poolmanager import PoolManager

class SSLAdapter(HTTPAdapter):
    def init_poolmanager(self, *args, **kwargs):
        context = ssl.create_default_context()
        context.set_ciphers('DEFAULT@SECLEVEL=1')
        kwargs['ssl_context'] = context
        return super().init_poolmanager(*args, **kwargs)

def insert_to_bigquery(table_id, data):
    client = bigquery.Client()
    # 给客户端挂载自定义SSL适配器
    client._http.mount('https://', SSLAdapter())
    
    dataset_id = 'test_youtube_data'
    table_ref = client.dataset(dataset_id).table(table_id)
    table = client.get_table(table_ref)
    errors = client.insert_rows_json(
        table=table, 
        json_rows=data,
        ignore_unknown_values=True, 
        skip_invalid_rows=True) 
    # 日志输出逻辑

4. 检查网络环境

  • Colab/本地环境:确认网络未被代理、防火墙拦截BigQuery的443端口
  • Cloud Function环境:若使用VPC连接器,检查配置是否允许出站访问BigQuery服务

5. 改用批量加载API替代行插入

如果行插入API持续报错,可将数据转为DataFrame后用批量加载API导入:

import pandas as pd
from google.cloud import bigquery

def insert_via_load_job(table_id, data):
    client = bigquery.Client()
    dataset_id = 'test_youtube_data'
    table_ref = client.dataset(dataset_id).table(table_id)
    
    df = pd.DataFrame(data)
    job_config = bigquery.LoadJobConfig(
        write_disposition=bigquery.WriteDisposition.WRITE_APPEND,
        ignore_unknown_values=True
    )
    job = client.load_table_from_dataframe(df, table_ref, job_config=job_config)
    job.result()  # 等待加载任务完成
    print(f'成功加载 {job.output_rows} 行数据')

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 06:25:19