如何使用Python获取gzip数据并写入BigQuery
故障原因
你当前实现的核心问题是将完整gzip文件全量加载到内存后解压转DataFrame:30MB的gzip压缩包解压后CSV体积通常在300MB~1GB区间,转成pandas DataFrame的内存峰值还会再翻2~3倍,很容易触发运行环境的内存上限,导致读取环节崩溃。
BigQuery原生支持直接加载gzip压缩的CSV文件,完全不需要经过本地解压、转DataFrame的中间环节,没有本地内存瓶颈,支持最大单文件4GB的gzip加载,是这类场景的最优方案。
最优实现方案:直接调用BigQuery Load Job加载gzip源文件
实现逻辑
- 拿到Branch返回的gzip临时下载URL后,不需要把文件内容下载到本地内存
- 配置BigQuery Load任务,指定源格式为CSV、压缩格式为Gzip,直接拉取gzip源数据写入目标表
- 等待Load任务执行完成即可,全程无大内存占用
参考代码
提前安装依赖:pip install google-cloud-bigquery requests schedule
from google.cloud import bigquery import requests import json import schedule import threading from datetime import date, timedelta # 提前初始化BQ客户端,确保运行环境已配置好服务账号访问凭证 bq_client = bigquery.Client(project=projectId) def getandwritebq(): dateYesterday = (date.today() - timedelta(days=1)).strftime("%Y-%m-%d") url = "https://api2.branch.io/v3/export" header = {'Content-Type': 'application/json'} dat = json.dumps({ "branch_key": 'xxxxx', "branch_secret":"xxxxx", "export_date":f"{dateYesterday}" }) resp = requests.post(url, headers=header, data=dat) print(f'{resp.status_code}') print('branch:getting list of web') listItems = ['eo_click','eo_commerce_event','eo_custom_event','eo_impression','eo_install','eo_open','eo_reinstall','eo_user_lifecycle_event'] bqlist=['branch_eo_click','branch_eo_commerce_event','branch_eo_custom_event','branch_eo_impression','branch_eo_install','branch_eo_open','branch_eo_reinstall','branch_eo_user_lifecycle_event'] # 统一Load任务配置 job_config = bigquery.LoadJobConfig( source_format=bigquery.SourceFormat.CSV, compression=bigquery.Compression.GZIP, autodetect=True, # 自动识别表结构,有固定schema可替换为自定义schema配置 write_disposition=bigquery.WriteDisposition.WRITE_APPEND, skip_leading_rows=1 # CSV带表头时开启,无表头则删除该配置 ) for bq,i in zip(bqlist,listItems): tempurl=json.loads(resp.text)[i][0] print(tempurl) table_id = f"{projectId}.{tableName}.{bq}" # 直接提交Load任务,从gzip源URL拉取数据写入BQ load_job = bq_client.load_table_from_uri( tempurl, table_id, job_config=job_config ) # 阻塞等待任务完成 load_job.result() print(f'branch:{bq} push to bigquery, loaded rows: {load_job.output_rows}') def run_threaded(job_func): job_thread = threading.Thread(target=job_func) job_thread.start() if __name__ == "__main__": schedule.every().day.at('06:00').do(run_threaded,lambda:getandwritebq()) while True: schedule.run_pending()
备选方案:必须保留DataFrame处理逻辑的优化方式
如果有数据清洗、字段转换的需求必须经过pandas处理,不要全量加载文件,做分块读写即可:
- 给
requests.get加上stream=True参数,开启流式下载,不要一次性把全量响应内容读到内存 - 调用
pd.read_csv时指定chunksize参数(建议设为10000~50000行),分块读取解压后的CSV内容,每块处理完成后立刻写入BigQuery,写完就释放当前块的内存 - 写入时给
pandas_gbq.to_gbq指定chunksize参数,避免单次写入攒出过大的内存缓存
注意事项
- gzip是不支持分片读取的压缩格式,BigQuery加载gzip压缩文件的单文件上限为4GB,如果后续Branch导出的单gzip文件超过这个阈值,需要在Branch导出配置中开启文件切分,输出多个小体积gzip文件再逐个加载
- 如果BigQuery Load Job提示无法访问公网gzip源地址,可以先将gzip文件流式上传到GCS临时存储桶,再从GCS路径加载数据,全程不需要将文件落地到本地内存
内容的提问来源于stack exchange,提问作者Hal
相关产品推荐
相关产品推荐

