このガイドに沿って進めると、アプリケーションからマルチパートアップロードを実行できる Worker を作成できます。 このサンプル Worker は、認証の追加や、各パートのアップロード時に検証ロジックを足すなど、自身のユースケースの土台にできます。 このガイドには、この Worker へファイルをアップロードする Python のサンプルアプリケーションもあります。
このガイドは、Worker に R2 バインディング を設定済みであることを前提とします。R2 バインディングの設定手順は Workers から R2 を使う を参照してください。
次のサンプル 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 へアップロードします。パートのアップロードを並列化するために 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 や別のデータベースで追跡できます。