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

如何用Python分块读取Google Drive大文件并逐块插入数据库?

解决方案:分块下载+分块解析入库

当然可以实现边下载边解析插入,但要注意:Google Drive的next_chunk()返回的是原始字节分块,不是结构化的CSV数据块——直接解析单个字节块会出现半行数据、格式错误的问题,因为CSV的行边界不一定和下载块的边界对齐。

最可靠的实现方式是:把下载的字节流写入本地临时文件,再用Pandas按数据块读取这个临时文件,逐块插入数据库。这样既避免了全量加载到内存,又保证CSV解析的完整性。

修改后的代码示例

1. 整合下载+分块入库的函数

import os
import pandas as pd
from googleapiclient.http import MediaIoBaseDownload
import logging

logger = logging.getLogger(__name__)

def download_and_insert_to_snowflake(file_id, file_name, service, tempdir, table_name, snowflake_client):
    """从Google Drive分块下载CSV文件,逐块插入Snowflake数据库"""
    request = service.files().get_media(fileId=file_id)
    temp_file_path = os.path.join(tempdir, f"{file_name}.tmp")
    
    # 分块下载到临时文件
    with open(temp_file_path, 'wb') as temp_fh:
        downloader = MediaIoBaseDownload(temp_fh, request)
        done = False
        try:
            while not done:
                status, done = downloader.next_chunk()
                logger.info(f"下载进度: {status.progress()*100:.2f}%")
            logger.info(f"文件 {file_name} 下载完成,临时文件路径: {temp_file_path}")
        except Exception as exc:
            logger.error(f"文件 {file_name} 下载失败")
            # 清理临时文件
            if os.path.exists(temp_file_path):
                os.remove(temp_file_path)
            raise exc
    
    # 分块读取临时文件并插入数据库
    try:
        # 按chunk_size分块读取CSV,避免内存溢出
        for df_chunk in pd.read_csv(temp_file_path, chunksize=100000):
            df_chunk.to_sql(
                table_name,
                snowflake_client.connection,
                if_exists="append",
                index=False,
                method=sfpd.pd_writer,
            )
            logger.info(f"已插入 {len(df_chunk)} 条数据到表 {table_name}")
        logger.info(f"文件 {file_name} 全量数据插入完成")
    except Exception as exc:
        logger.error(f"文件 {file_name} 插入数据库失败")
        raise exc
    finally:
        # 无论成功失败,都清理临时文件
        if os.path.exists(temp_file_path):
            os.remove(temp_file_path)

2. 关键改动说明

  • 不再把全量文件加载到内存的BytesIO,而是直接写入本地临时文件,大幅降低内存占用
  • 用pd.read_csv(chunksize=...)按行块读取CSV,每块数据直接调用to_sql插入,不需要缓存全量DataFrame
  • 增加了临时文件的清理逻辑,避免磁盘空间浪费

替代方案(无临时文件)

如果不想依赖本地文件,可以用内存缓冲区积累字节,直到凑出完整的CSV行再解析,但实现复杂度高,需要处理行边界、编码、换行符等问题,不如临时文件方案稳定,不推荐在生产环境使用。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 20:39:19