INTEGRITY Dokumentace

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.