← Cloudflare R2 / r2 / api / workers
Použití vícedílného API R2 z Workers
Podle tohoto průvodce vytvoříte Worker, přes který mohou vaše aplikace provádět multipart uploady. Tento ukázkový Worker může sloužit jako základ pro váš vlastní případ použití: můžete do něj přidat ověřování, nebo třeba další validační logiku při nahrávání jednotlivých částí. Součástí tohoto průvodce je také ukázková aplikace v Pythonu, která do tohoto Workeru nahrává soubory.
Tento návod předpokládá, že máte nastavené Vazba R2 pro váš Worker. Viz Použití R2 z Workers pro pokyny k nastavení R2 bindingu.
Příklad Workeru využívajícího multipart API
Následující ukázkový Worker vystavuje HTTP API, které aplikacím umožňuje používat vícedílné API prostřednictvím Workeru.
V tomto příkladu je každý požadavek směrován podle metody HTTP a parametru požadavku action. Jakmile se váš Worker stane složitějším, zvažte využití bezserverového webového frameworku, jako je Hono ↗ aby za vás zajistil směrování.
Následující ukázkový Worker zahrnuje do odpovědi na každý požadavek veškeré nové informace o stavu vícedílného nahrávání. U požadavku, který vícedílné nahrávání vytváří, uploadId se vrátí. U požadavků nahrávajících část se vrátí číslo části a etag jsou vráceny. Klient si tento stav následně uchovává a do dalších požadavků zahrnuje uploadId a etag a číslo dílu každé části při dokončování multipart uploadu.
Přidejte následující kód do souboru projektu index.ts soubor a nahraďte MY_BUCKET s názvem vašeho bucketu:
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"},
)Po aktualizaci Workeru výše uvedeným kódem spusťte npx wrangler deploy.
Tento Worker nyní můžete použít k provádění multipart uploadů. Požadavky na nahrávání můžete posílat tomuto Workeru přímo ze své stávající aplikace, nebo k nahrávání souborů přes tento Worker použít skript.
Následující část je nepovinná a ukazuje příklad skriptu v Pythonu, který nahraje vybraný soubor z vašeho počítače do vašeho Workeru.
Proveďte multipart upload pomocí svého Workeru (volitelné)
Tato ukázková aplikace nahrává lokální soubor do Workeru po částech. Využívá vestavěný modul Pythonu ThreadPoolExecutor k paralelnímu nahrávání částí do Workeru, což zvyšuje rychlost nahrávání. HTTP požadavky na Worker se odesílají pomocí požadavky ↗ knihovna.
Použití vícedílného API tímto způsobem vám také umožňuje nahrávat pomocí Workeru soubory větší než Limit velikosti těla požadavku u Workers. Nahrávání jednotlivých částí tomuto limitu i nadále podléhá.
Uložte následující kód do souboru s názvem mpuscript.py na vašem lokálním počítači. Změňte worker_endpoint variable tam, kde je nasazen váš worker. Soubor, který chcete nahrát, předejte jako argument při spuštění tohoto skriptu: python3 mpuscript.py myfile. Tím se nahraje soubor myfile z vašeho počítače do bucketu přes 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)Správa stavu
Stavová povaha vícedílných nahrávání se obtížně slučuje s modelem používání Workers, které jsou ze své podstaty bezstavové. Při běžném vícedílném nahrávání celý proces obvykle probíhá v rámci jednoho nepřerušeného běhu klientské aplikace. To se liší od vícedílného nahrávání ve Workeru, kde se nahrávání často dokončuje v rámci více vyvolání daného Workeru. Správa stavu je tak složitější.
K vyřešení tohoto problému se stav přidružený k vícedílnému nahrávání, konkrétně uploadId a které díly byly nahrány, je třeba sledovat někde mimo Worker.
V ukázkovém Workeru a Python aplikaci popsaných v tomto průvodci se stav vícedílného nahrávání sleduje v klientské aplikaci, která odesílá požadavky Workeru, přičemž každý požadavek obsahuje potřebný stav. Sledování stavu vícedílného nahrávání v klientské aplikaci umožňuje maximální flexibilitu a paralelní i neuspořádané nahrávání jednotlivých částí.
Pokud není možné tento stav sledovat na straně klienta, lze zvážit alternativní návrhy. Můžete například sledovat uploadId a které díly byly nahrány, v Durable Object nebo jiné databázi.