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
相关产品推荐
相关产品推荐

