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

使用Tornado将二进制文件流式传输至Google Storage的问题求助

解决Tornado流式请求边接收边上传至Google Cloud Storage的问题

我之前刚处理过类似的场景,你遇到的核心问题在于:GCS的upload_from_file方法默认是从类文件对象读取到EOF才会完成上传,但Tornado的流式请求是分块接收数据的,这就导致你没法在数据还没接收完的时候启动上传流程。要实现边接收边上传,得换用GCS的**可恢复上传(Resumable Uploads)**方案,它天生支持分块传输,完美适配流式场景。

下面是具体的实现步骤和代码示例:

核心思路

  1. 在Tornado请求的初始化阶段(prepare方法),先和GCS建立一个可恢复上传会话,拿到专属的上传URL。
  2. 每次Tornado收到请求数据块(data_received方法),就直接把这个数据块发送到GCS的上传会话中。
  3. 当所有数据接收完成(finish方法),通知GCS上传结束,完成文件的最终写入。

完整代码示例

import tornado.web
from google.cloud import storage
from google.resumable_media.requests import ResumableUpload
from google.resumable_media.common import DataCorruption
import asyncio

class StreamUploadHandler(tornado.web.RequestHandler):
    @tornado.web.stream_request_body
    def initialize(self):
        # 初始化GCS客户端
        self.storage_client = storage.Client()
        self.bucket = self.storage_client.bucket("your-bucket-name")
        self.upload_session = None
        self.bytes_uploaded = 0

    async def prepare(self):
        # 创建可恢复上传会话
        destination_blob_name = self.get_argument("file_name", default="default_file.bin")
        blob = self.bucket.blob(destination_blob_name)
        
        # 初始化可恢复上传请求
        upload = ResumableUpload(
            blob.self_link,
            blob.content_type or "application/octet-stream",
        )
        # 发起会话请求,获取上传URL
        response = await asyncio.to_thread(
            upload.initiate,
            self.storage_client._credentials,
            stream=None,  # 因为是流式上传,初始不需要数据
            content_length=None,  # 未知总长度的话可以设为None
        )
        self.upload_session = upload
        self.bytes_uploaded = 0

    async def data_received(self, chunk):
        if not self.upload_session:
            return
        
        try:
            # 发送当前数据块到GCS
            await asyncio.to_thread(
                self.upload_session.transmit,
                chunk,
                self.bytes_uploaded,
            )
            self.bytes_uploaded += len(chunk)
        except DataCorruption as e:
            # 处理数据损坏异常,可根据需求重试
            self.set_status(500)
            self.finish(f"Upload failed due to data corruption: {str(e)}")

    async def finish(self):
        if not self.upload_session:
            return
        
        try:
            # 完成上传会话
            await asyncio.to_thread(
                self.upload_session.finalize,
                self.storage_client._credentials,
            )
            self.set_status(200)
            self.finish("Upload completed successfully")
        except Exception as e:
            self.set_status(500)
            self.finish(f"Failed to finalize upload: {str(e)}")

# Tornado应用配置
def make_app():
    return tornado.web.Application([
        (r"/upload", StreamUploadHandler),
    ])

if __name__ == "__main__":
    app = make_app()
    app.listen(8888)
    tornado.ioloop.IOLoop.current().start()

关键细节说明

  • 可恢复上传的优势:它允许你分块发送数据,不需要预先知道文件总大小(适合流式请求场景),而且上传中断后可以从断点续传,稳定性更高。
  • 异步处理:因为GCS的SDK方法是同步的,所以我们用asyncio.to_thread把同步调用包装成异步操作,避免阻塞Tornado的IO循环。
  • 数据块大小:GCS要求除了最后一块数据,其他块的大小至少为256KB。如果你的Tornado接收的chunk太小,可以先缓存到内存,凑够大小再上传(比如维护一个缓冲区,达到256KB再调用transmit)。
  • 异常处理:代码里加了基础的异常捕获,你可以根据业务需求扩展重试逻辑或者错误反馈。

额外优化建议

  • 如果你的请求有Content-Length头,可以在初始化上传会话时把这个值传给content_length,GCS会做更严格的校验。
  • 对于大文件,可以考虑在data_received中使用缓冲区,避免频繁调用transmit产生过多的网络请求。
  • 可以添加日志记录,跟踪每个数据块的上传状态,方便排查问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 06:51:39