按照本指南,您将创建一个 Worker,应用可通过它执行分片上传。 此示例 Worker 可作为您自己用例的基础,您可以在 Worker 中添加身份验证,或在上传每个分片时添加额外验证逻辑。 本指南还包含向此 Worker 上传文件的 Python 应用示例。
本指南假设您已为 Worker 设置 R2 绑定(binding)。有关设置 R2 绑定的说明,请参阅 从 Workers 使用 R2。
以下示例 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。它使用 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 和已上传的分片。