跳转到内容
搜索文档

从 Workers 使用 R2 分片 API

最后更新 查看 MarkdownAgent 设置

按照本指南,您将创建一个 Worker,应用可通过它执行分片上传。 此示例 Worker 可作为您自己用例的基础,您可以在 Worker 中添加身份验证,或在上传每个分片时添加额外验证逻辑。 本指南还包含向此 Worker 上传文件的 Python 应用示例。

本指南假设您已为 Worker 设置 R2 绑定(binding)。有关设置 R2 绑定的说明,请参阅 从 Workers 使用 R2

使用分片 API 的示例 Worker

以下示例 Worker 暴露 HTTP API,使应用可通过 Worker 使用分片 API。

在此示例中,每个请求根据 HTTP 方法和 action 请求参数路由。随着 Worker 变得更复杂,可考虑使用 Hono 等 serverless Web 框架处理路由。

以下示例 Worker 在每次请求的响应中包含分片上传状态的任何新信息。创建分片上传的请求返回 uploadId。上传分片的请求返回分片编号和 etag。客户端跟踪此状态,在后续请求中包含 uploadId,完成分片上传时包含每个分片的 etag 和分片编号。

将以下代码添加到项目的 index.js 文件,并将 MY_BUCKET 替换为您的存储桶名称:

interface Env {
  MY_BUCKET: R2Bucket;
}

export default {
  async fetch(
    request,
    env,
    ctx
  ): Promise<Response> {
    const bucket = env.MY_BUCKET;

    const url = new URL(request.url);
    const key = url.pathname.slice(1);
    const action = url.searchParams.get("action");

    if (action === null) {
      return new Response("Missing action type", { status: 400 });
    }

    // Route the request based on the HTTP method and action type
    switch (request.method) {
      case "POST":
        switch (action) {
          case "mpu-create": {
            const multipartUpload = await bucket.createMultipartUpload(key);
            return new Response(
              JSON.stringify({
                key: multipartUpload.key,
                uploadId: multipartUpload.uploadId,
              })
            );
          }
          case "mpu-complete": {
            const uploadId = url.searchParams.get("uploadId");
            if (uploadId === null) {
              return new Response("Missing uploadId", { status: 400 });
            }

            const multipartUpload = env.MY_BUCKET.resumeMultipartUpload(
              key,
              uploadId
            );

            interface completeBody {
              parts: R2UploadedPart[];
            }
            const completeBody: completeBody = await request.json();
            if (completeBody === null) {
              return new Response("Missing or incomplete body", {
                status: 400,
              });
            }

            // Error handling in case the multipart upload does not exist anymore
            try {
              const object = await multipartUpload.complete(completeBody.parts);
              return new Response(null, {
                headers: {
                  etag: object.httpEtag,
                },
              });
            } catch (error: any) {
              return new Response(error.message, { status: 400 });
            }
          }
          default:
            return new Response(`Unknown action ${action} for POST`, {
              status: 400,
            });
        }
      case "PUT":
        switch (action) {
          case "mpu-uploadpart": {
            const uploadId = url.searchParams.get("uploadId");
            const partNumberString = url.searchParams.get("partNumber");
            if (partNumberString === null || uploadId === null) {
              return new Response("Missing partNumber or uploadId", {
                status: 400,
              });
            }
            if (request.body === null) {
              return new Response("Missing request body", { status: 400 });
            }

            const partNumber = parseInt(partNumberString);
            const multipartUpload = env.MY_BUCKET.resumeMultipartUpload(
              key,
              uploadId
            );
            try {
              const uploadedPart: R2UploadedPart =
                await multipartUpload.uploadPart(partNumber, request.body);
              return new Response(JSON.stringify(uploadedPart));
            } catch (error: any) {
              return new Response(error.message, { status: 400 });
            }
          }
          default:
            return new Response(`Unknown action ${action} for PUT`, {
              status: 400,
            });
        }
      case "GET":
        if (action !== "get") {
          return new Response(`Unknown action ${action} for GET`, {
            status: 400,
          });
        }
        const object = await env.MY_BUCKET.get(key);
        if (object === null) {
          return new Response("Object Not Found", { status: 404 });
        }
        const headers = new Headers();
        object.writeHttpMetadata(headers);
        headers.set("etag", object.httpEtag);
        return new Response(object.body, { headers });
      case "DELETE":
        switch (action) {
          case "mpu-abort": {
            const uploadId = url.searchParams.get("uploadId");
            if (uploadId === null) {
              return new Response("Missing uploadId", { status: 400 });
            }
            const multipartUpload = env.MY_BUCKET.resumeMultipartUpload(
              key,
              uploadId
            );

            try {
              multipartUpload.abort();
            } catch (error: any) {
              return new Response(error.message, { status: 400 });
            }
            return new Response(null, { status: 204 });
          }
          case "delete": {
            await env.MY_BUCKET.delete(key);
            return new Response(null, { status: 204 });
          }
          default:
            return new Response(`Unknown action ${action} for DELETE`, {
              status: 400,
            });
        }
      default:
        return new Response("Method Not Allowed", {
          status: 405,
          headers: { Allow: "PUT, POST, GET, DELETE" },
        });
    }
  },
} satisfies ExportedHandler<Env>;
from workers import WorkerEntrypoint, Response
from urllib.parse import urlparse, parse_qs
import json


class Default(WorkerEntrypoint):
    async def fetch(self, request):
        bucket = self.env.MY_BUCKET

        url = urlparse(request.url)
        key = url.path[1:]
        params = parse_qs(url.query)
        action = params.get("action", [None])[0]

        if action is None:
            return Response("Missing action type", status=400)

        if request.method == "POST":
            if action == "mpu-create":
                multipart_upload = await bucket.createMultipartUpload(key)
                return Response.json({
                    "key": multipart_upload.key,
                    "uploadId": multipart_upload.uploadId,
                })
            elif action == "mpu-complete":
                upload_id = params.get("uploadId", [None])[0]
                if upload_id is None:
                    return Response("Missing uploadId", status=400)

                multipart_upload = bucket.resumeMultipartUpload(key, upload_id)
                complete_body = await request.json()
                if complete_body is None:
                    return Response("Missing or incomplete body", status=400)

                try:
                    obj = await multipart_upload.complete(complete_body.parts)
                    return Response(None, headers={"etag": obj.httpEtag})
                except Exception as error:
                    return Response(str(error), status=400)
            else:
                return Response(f"Unknown action {action} for POST", status=400)

        elif request.method == "PUT":
            if action == "mpu-uploadpart":
                upload_id = params.get("uploadId", [None])[0]
                part_number_str = params.get("partNumber", [None])[0]
                if part_number_str is None or upload_id is None:
                    return Response("Missing partNumber or uploadId", status=400)
                if request.body is None:
                    return Response("Missing request body", status=400)

                part_number = int(part_number_str)
                multipart_upload = bucket.resumeMultipartUpload(key, upload_id)
                try:
                    uploaded_part = await multipart_upload.uploadPart(part_number, request.body)
                    return Response.json(uploaded_part)
                except Exception as error:
                    return Response(str(error), status=400)
            else:
                return Response(f"Unknown action {action} for PUT", status=400)

        elif request.method == "GET":
            if action != "get":
                return Response(f"Unknown action {action} for GET", status=400)
            obj = await bucket.get(key)
            if obj is None:
                return Response("Object Not Found", status=404)
            body = await obj.text()
            headers = {"etag": obj.httpEtag}
            return Response(body, headers=headers)

        elif request.method == "DELETE":
            if action == "mpu-abort":
                upload_id = params.get("uploadId", [None])[0]
                if upload_id is None:
                    return Response("Missing uploadId", status=400)
                multipart_upload = bucket.resumeMultipartUpload(key, upload_id)
                try:
                    await multipart_upload.abort()
                except Exception as error:
                    return Response(str(error), status=400)
                return Response(None, status=204)
            elif action == "delete":
                await bucket.delete(key)
                return Response(None, status=204)
            else:
                return Response(f"Unknown action {action} for DELETE", status=400)

        else:
            return Response(
                "Method Not Allowed",
                status=405,
                headers={"Allow": "PUT, POST, GET, DELETE"},
            )

使用上述代码更新 Worker 后,运行 npx wrangler deploy

您现在可以使用此 Worker 执行分片上传。您可以从现有应用向此 Worker 发送请求执行上传,或使用脚本通过此 Worker 上传文件。

下一节为可选内容,展示将机器上选定文件上传到 Worker 的 Python 脚本示例。

使用 Worker 执行分片上传(可选)

此示例应用将本地文件分多个部分上传到 Worker。它使用 Python 内置的 ThreadPoolExecutor 并行上传分片到 Worker,提高上传速度。向 Worker 的 HTTP 请求使用 requests 库。

以这种方式使用分片 API 还允许您使用 Worker 上传大于 Workers 请求体大小限制 的文件。单个分片的上传仍受此限制。

将以下代码保存为本地机器上的 mpuscript.py 文件。将 worker_endpoint variable 更改为 Worker 的部署地址。运行脚本时将要上传的文件作为参数传入:python3 mpuscript.py myfile。这将通过 Worker 将机器上的 myfile 文件上传到存储桶。

import math
import os
import requests
from requests.adapters import HTTPAdapter, Retry
import sys
import concurrent.futures

# Take the file to upload as an argument
filename = sys.argv[1]
# The endpoint for our worker, change this to wherever you deploy your worker
worker_endpoint = "https://myworker.myzone.workers.dev/"
# Configure the part size to be 10MB. 5MB is the minimum part size, except for the last part
partsize = 10 * 1024 * 1024


def upload_file(worker_endpoint, filename, partsize):
    url = f"{worker_endpoint}{filename}"

    # Create the multipart upload
    uploadId = requests.post(url, params={"action": "mpu-create"}).json()["uploadId"]

    part_count = math.ceil(os.stat(filename).st_size / partsize)
    # Create an executor for up to 25 concurrent uploads.
    executor = concurrent.futures.ThreadPoolExecutor(25)
    # Submit a task to the executor to upload each part
    futures = [
        executor.submit(upload_part, filename, partsize, url, uploadId, index)
        for index in range(part_count)
    ]
    concurrent.futures.wait(futures)
    # get the parts from the futures
    uploaded_parts = [future.result() for future in futures]

    # complete the multipart upload
    response = requests.post(
        url,
        params={"action": "mpu-complete", "uploadId": uploadId},
        json={"parts": uploaded_parts},
    )
    if response.status_code == 200:
        print("🎉 successfully completed multipart upload")
    else:
        print(response.text)


def upload_part(filename, partsize, url, uploadId, index):
    # Open the file in rb mode, which treats it as raw bytes rather than attempting to parse utf-8
    with open(filename, "rb") as file:
        file.seek(partsize * index)
        part = file.read(partsize)

    # Retry policy for when uploading a part fails
    s = requests.Session()
    retries = Retry(total=3, status_forcelist=[400, 500, 502, 503, 504])
    s.mount("https://", HTTPAdapter(max_retries=retries))

    return s.put(
        url,
        params={
            "action": "mpu-uploadpart",
            "uploadId": uploadId,
            "partNumber": str(index + 1),
        },
        data=part,
    ).json()


upload_file(worker_endpoint, filename, partsize)

状态管理

分片上传的有状态特性不易映射到 Workers 的使用模型,Workers 本质上是无状态的。在正常的分片上传中,分片上传通常在客户端应用的一次连续执行中完成。这与 Worker 中的分片上传不同,后者通常会在 Worker 的多次调用中完成。这使状态管理更具挑战性。

为克服此问题,与分片上传关联的状态(即 uploadId 和已上传的分片)需要在 Worker 外部某处跟踪。

在本指南描述的示例 Worker 和 Python 应用中,分片上传状态在发送请求到 Worker 的客户端应用中跟踪,必要状态包含在每次请求中。在客户端应用中跟踪分片状态可实现最大灵活性,并允许并行和无序上传每个分片。

当无法在客户端跟踪此状态时,可考虑替代设计。例如,您可以在 Durable Object 或其他数据库中跟踪 uploadId 和已上传的分片。

这篇文档对您有帮助吗?