INTEGRITY Документация

Использование R2 multipart API в Workers

Следуя этому руководству, вы создадите Worker, через который ваши приложения смогут выполнять составную загрузку. Этот пример worker может служить основой для вашего собственного сценария использования, где вы можете добавить аутентификацию к worker или даже дополнительную логику валидации при загрузке каждой части. Это руководство также содержит пример приложения на Python, которое загружает файлы через этот worker.

В этом руководстве предполагается, что вы уже настроили Привязка R2 для вашего Worker. См. Использование R2 в Workers для инструкций по настройке R2 binding.

Пример Worker с использованием multipart API

В следующем примере Worker предоставляет HTTP API, который позволяет приложениям использовать составной (multipart) API через Worker.

В этом примере каждый запрос маршрутизируется на основе HTTP-метода и параметра запроса action. По мере усложнения вашего Worker рассмотрите возможность использования serverless веб-фреймворка, например Hono чтобы обработать маршрутизацию за вас.

В следующем примере 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.

Следующий раздел необязателен: в нём приведён пример скрипта на Python, который загружает выбранный файл с вашего компьютера в Worker.

Составная загрузка с помощью Worker (необязательно)

Это пример приложения, которое загружает локальный файл в Worker несколькими частями. Оно использует встроенный в Python ThreadPoolExecutor чтобы распараллелить загрузку частей в Worker, что увеличивает скорость загрузки. HTTP-запросы к Worker выполняются с запросы библиотека.

Такой способ применения составного API также позволяет использовать Worker для загрузки файлов размером больше Ограничение размера тела запроса Workers. Загрузка отдельных частей по-прежнему подчиняется этому ограничению.

Сохраните следующий код в файл с именем mpuscript.py на вашем локальном компьютере. Измените worker_endpoint variable туда, где развёрнут ваш worker. При запуске этого скрипта передайте файл, который нужно загрузить, в качестве аргумента: 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, которые изначально не хранят состояние. При обычной составной загрузке она, как правило, выполняется в рамках одного непрерывного запуска клиентского приложения. В Worker всё иначе: составная загрузка часто завершается за несколько отдельных вызовов этого Worker, из-за чего управление состоянием усложняется.

Чтобы решить эту проблему, состояние, связанное с составной загрузкой, а именно uploadId и какие части уже загружены, необходимо отслеживать где-то за пределами Worker.

В примере Worker и Python-приложения, описанном в этом руководстве, состояние составной загрузки отслеживается в клиентском приложении, которое отправляет запросы к Worker, причём необходимое состояние содержится в каждом запросе. Отслеживание состояния составной загрузки на стороне клиентского приложения обеспечивает максимальную гибкость и позволяет выполнять параллельную и не упорядоченную по времени загрузку частей.

Если отслеживать это состояние на стороне клиента невозможно, можно рассмотреть альтернативные варианты реализации. Например, можно отслеживать uploadId и какие части уже загружены, в Durable Object или другой базе данных.