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

Python实现GCS大文件分块流式传输时in_memory_file无法输出分块

问题核心原因

你的代码无法实现预期流式传输,是以下几个逻辑错误导致的:

  • ThreadPoolExecutor的with上下文管理器会阻塞当前线程,直到所有提交的任务全部执行完成才会退出代码块。你把响应初始化放在with块内部,必须等全部分块下载完成才会开始返回响应,完全无法实现边下载边返回的效果。
  • ex.map()返回的是BytesIO内存对象的迭代器,Flask接收迭代器作为响应体时,不会自动读取BytesIO内部存储的二进制数据,只会直接序列化对象本身,客户端拿不到实际文件内容。
  • 分块计算逻辑split_byte_size完全错误,循环范围和字节终止位计算不符合Range请求规范,会生成无效、重复的字节区间。
  • 多线程并发下载分块无法保证返回顺序和文件原有字节顺序一致,会直接导致CSV文件内容错乱损坏。
  • 信号注册逻辑错误:signal.signal(signal.SIGTERM, shutdown_handler)写在了shutdown_handler函数内部,永远不会被执行,无法正常捕获SIGTERM信号。
  • 状态码使用错误:206状态码专门用于响应客户端主动携带Range请求头的部分内容请求,你要返回完整文件的流式响应应该使用200状态码。
  • structlog初始化配置里残留了无意义的`enter code here`占位符,会直接导致代码启动报错。
修复方案

直接用生成器实现流式响应,按顺序逐块下载GCS文件内容,每下载完一个分块立刻将字节内容yield给响应流,既保证内容顺序正确,又能实现第一块下载完成就开始返回的流式效果,同时修复所有已知逻辑错误。

修复后的完整代码如下:

import os
import signal
import sys
import io
from types import FrameType
from google.cloud import storage
from google.oauth2 import service_account
from flask import Flask, Response
import structlog

app = Flask(__name__)

# 初始化结构化日志
def getJSONLogger() -> structlog.stdlib.BoundLogger:
    structlog.configure(
        processors=[
            structlog.stdlib.add_log_level,
            structlog.stdlib.PositionalArgumentsFormatter(),
            structlog.processors.TimeStamper("iso"),
            structlog.processors.JSONRenderer(),
        ],
        wrapper_class=structlog.stdlib.BoundLogger,
    )
    return structlog.get_logger()

logger = getJSONLogger()

# 进程终止信号处理
def shutdown_handler(signal: int, frame: FrameType) -> None:
    logger.info("Signal received, safely shutting down.")
    print("Exiting process.", flush=True)
    sys.exit(0)

# 计算正确的分块字节区间
def split_byte_range(total_size: int, chunk_count: int = 50) -> list:
    byte_ranges = []
    chunk_size = total_size // chunk_count
    start = 0
    while start < total_size:
        end = min(start + chunk_size - 1, total_size - 1)
        byte_ranges.append((start, end))
        start = end + 1
    return byte_ranges

# 初始化GCS客户端
project = 'XYZ'
service_account_credentials_path = 'key.json'
credentials = service_account.Credentials.from_service_account_file(service_account_credentials_path) 
storage_client = storage.Client(project=project, credentials=credentials)

@app.route("/chunk_data")
def chunk_data():
    bucket_name = 'cloudrundemofile'
    source_blob_name = 'demofile.csv'
    
    bucket = storage_client.get_bucket(bucket_name)
    blob = bucket.get_blob(source_blob_name)
    # 生成分块区间
    byte_ranges = split_byte_range(blob.size)

    # 流式生成器:逐块下载,逐块返回
    def generate_chunks():
        for start, end in byte_ranges:
            in_memory_file = io.BytesIO()
            blob.download_to_file(in_memory_file, start=start, end=end)
            # 读取内存对象中的二进制内容
            in_memory_file.seek(0)
            chunk_data = in_memory_file.read()
            # 释放内存
            in_memory_file.close()
            yield chunk_data

    # 返回流式响应
    resp = Response(generate_chunks(), 200, mimetype='text/csv')
    resp.headers["Content-Length"] = blob.size
    return resp

if __name__ == "__main__":
    signal.signal(signal.SIGINT, shutdown_handler)
    signal.signal(signal.SIGTERM, shutdown_handler)
    app.run(host="0.0.0.0", port=8080)
else:
    signal.signal(signal.SIGTERM, shutdown_handler)
优化说明
  • 去掉了无意义的多线程下载逻辑,按顺序下载分块从根源上避免了CSV内容乱序问题,同时减少了线程调度开销。
  • 每个分块传输完成后立刻关闭BytesIO对象释放内存,不会出现全量文件驻留内存的问题,内存占用稳定控制在单块大小级别。
  • 响应头添加了正确的Content-Length字段,客户端可以正确识别文件总大小,显示下载进度。
  • 移除了代码里未使用的冗余导入,避免不必要的依赖加载。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 19:45:31