Skip to content

非公式本サイトは非公式の日本語ドキュメントであり、Cloudflare 公式サイトではありません。最新情報はdevelopers.cloudflare.comをご確認ください。

Workers から R2 のマルチパート API を使う

最終更新 Markdown で表示Agent セットアップ

このガイドに沿って進めると、アプリケーションからマルチパートアップロードを実行できる Worker を作成できます。 このサンプル Worker は、認証の追加や、各パートのアップロード時に検証ロジックを足すなど、自身のユースケースの土台にできます。 このガイドには、この Worker へファイルをアップロードする Python のサンプルアプリケーションもあります。

このガイドは、Worker に R2 バインディング を設定済みであることを前提とします。R2 バインディングの設定手順は Workers から R2 を使う を参照してください。

マルチパート API を使う Worker の例

次のサンプル Worker は、アプリケーションが Worker 経由でマルチパート API を使える HTTP API を公開します。

この例では、HTTP メソッドと action リクエストパラメーターに基づいて各リクエストを振り分けます。Worker が複雑になってきたら、ルーティングは Hono などのサーバーレス Web フレームワークに任せることを検討してください。

次のサンプル Worker は、マルチパートアップロードの新しい状態情報を各リクエストのレスポンスに含めます。マルチパートアップロードを作成するリクエストでは uploadId を返します。パートをアップロードするリクエストでは、パート番号と etag を返します。クライアント側でこの状態を保持し、以降のリクエストに uploadId を含め、マルチパートアップロードの完了時には各パートの etag とパート番号を含めます。

プロジェクトの index.ts に次のコードを追加し、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 への HTTP リクエストには requests ライブラリを使います。

この方法でマルチパート API を使うと、Workers のリクエストボディサイズ制限 を超えるファイルも Worker 経由でアップロードできます。個々のパートのアップロードには、この制限が引き続き適用されます。

次のコードを手元のマシンに mpuscript.py として保存します。worker_endpoint 変数を、デプロイした Worker の URL に変更します。アップロードするファイルを引数にしてスクリプトを実行します: python3 mpuscript.py myfile。これで、手元のファイル myfile が Worker 経由でバケットにアップロードされます。

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 の利用モデルとは相性がよくありません。通常のマルチパートアップロードでは、クライアントアプリケーションの 1 回の連続した実行で完了することが多いです。Worker でのマルチパートアップロードは、複数回の呼び出しにまたがって完了することが多く、状態管理が難しくなります。

これを解決するには、マルチパートアップロードに紐づく状態、つまり uploadId とどのパートがアップロード済みかを、Worker の外で追跡する必要があります。

このガイドのサンプル Worker と Python アプリケーションでは、Worker へリクエストを送るクライアントアプリケーション側でマルチパートアップロードの状態を追跡し、必要な状態を各リクエストに含めます。クライアント側で状態を持つと柔軟性が最大になり、各パートの並列アップロードや順不同のアップロードもできます。

クライアント側で状態を追跡できない場合は、別の設計を検討できます。たとえば、uploadId とアップロード済みパートを Durable Object や別のデータベースで追跡できます。

役に立ちましたか?