MidПрактика7 min

Two-stage processing

Quick ack + async process. Upload → URL return → async process. PHP Symfony Messenger + Go channels pipeline

Проблема: долгие операции в HTTP-запросе

Пользователь загружает видео. Сервер должен:

  1. Сохранить файл
  2. Транскодировать в 5 разрешений
  3. Сгенерировать превью
  4. Записать метаданные в БД
  5. Отправить уведомление подписчикам

Общее время: 30 секунд - 5 минут. Держать HTTP-соединение всё это время -- плохо:

  • Браузер / прокси таймаутят на 30-60 секундах
  • Retry тяжёлого запроса = повторная загрузка файла
  • Load balancer держит connection -- исчерпаются worker'ы
  • Пользователь смотрит на спиннер 2 минуты

Решение: разделить на два этапа

Stage 1 (синхронный, быстрый): принять данные, зафиксировать факт запроса, вернуть ID / URL статуса. < 500 ms.

Stage 2 (асинхронный): тяжёлая обработка в воркерах.

Client                 API                  Queue          Worker
  │                     │                     │                │
  │── POST /upload ────▶│                     │                │
  │                     │─ Save file          │                │
  │                     │─ INSERT jobs        │                │
  │                     │─ Publish --────────▶│                │
  │◄── 202 Accepted ────│                     │                │
  │   Location: /j/123  │                     │                │
  │                     │                     │─ Deliver ─────▶│
  │                     │                     │                │─ Transcode
  │                     │                     │                │─ Thumbnail
  │                     │                     │                │─ Notify
  │── GET /j/123 ──────▶│                     │                │─ UPDATE status
  │◄── status: done ────│                     │                │

HTTP статус для accepted-но-не-готово: 202 Accepted + Location: /jobs/{id}.

Pattern: Polling

Клиент периодически опрашивает /jobs/{id}:

POST /uploads
-> 202 Accepted
   Location: /jobs/abc123
   Retry-After: 5

GET /jobs/abc123
-> 200 OK
   { "status": "processing", "progress": 45 }

GET /jobs/abc123
-> 200 OK
   { "status": "done", "resultUrl": "..." }

Retry-After -- подсказка клиенту, через сколько секунд опрашивать. Exponential backoff на клиенте, если Retry-After не возвращается.

Альтернативы polling

  • WebSocket / SSE: push уведомления о статусе
  • Webhook: клиент регистрирует URL, сервер POST-ит на него готовый результат
  • Long polling: сервер держит соединение до готовности

Для публичного API обычно начинают с polling + webhook для premium клиентов.

PHP реализация (Symfony Messenger)

Миграция

CREATE TABLE jobs (
    id           UUID PRIMARY KEY,
    kind         TEXT NOT NULL,
    status       TEXT NOT NULL,          -- queued | processing | done | failed
    progress     SMALLINT NOT NULL DEFAULT 0,
    input        JSONB NOT NULL,
    result       JSONB,
    error        TEXT,
    created_at   TIMESTAMPTZ NOT NULL DEFAULT NOW(),
    started_at   TIMESTAMPTZ,
    finished_at  TIMESTAMPTZ
);
CREATE INDEX idx_jobs_status ON jobs (status, created_at);
<?php

declare(strict_types=1);

namespace App\Upload;

use Symfony\Bundle\FrameworkBundle\Controller\AbstractController;
use Symfony\Component\HttpFoundation\{Request, Response, JsonResponse};
use Symfony\Component\Messenger\MessageBusInterface;
use Symfony\Component\Routing\Annotation\Route;

final class UploadController extends AbstractController
{
    public function __construct(
        private readonly MessageBusInterface $bus,
        private readonly JobRepository $jobs,
        private readonly BlobStorage $blob,
    ) {}

    #[Route('/uploads', methods: ['POST'])]
    public function upload(Request $req): Response
    {
        // Stage 1: fast ack
        $file = $req->files->get('file');
        if (!$file) {
            return new JsonResponse(['error' => 'file missing'], 400);
        }

        // Persist raw bytes to object storage (S3/minio)
        $blobKey = $this->blob->put($file->getPathname(), $file->getMimeType());

        // Persist job in DB in queued state
        $jobId = bin2hex(random_bytes(16));
        $this->jobs->create([
            'id'     => $jobId,
            'kind'   => 'video.transcode',
            'status' => 'queued',
            'input'  => [
                'blobKey' => $blobKey,
                'mime'    => $file->getMimeType(),
                'size'    => $file->getSize(),
            ],
        ]);

        // Dispatch async message (ideally inside outbox-style TX with job insert)
        $this->bus->dispatch(new TranscodeVideo($jobId));

        return new JsonResponse(
            ['jobId' => $jobId, 'status' => 'queued'],
            Response::HTTP_ACCEPTED,
            ['Location' => "/jobs/$jobId", 'Retry-After' => '5'],
        );
    }

    #[Route('/jobs/{id}', methods: ['GET'])]
    public function status(string $id): JsonResponse
    {
        $job = $this->jobs->find($id);
        if (!$job) {
            return new JsonResponse(['error' => 'not found'], 404);
        }

        return new JsonResponse([
            'id'       => $job['id'],
            'status'   => $job['status'],
            'progress' => $job['progress'],
            'result'   => $job['result'],
            'error'    => $job['error'],
        ]);
    }
}
### messenger.yaml
framework:
    messenger:
        transports:
            async_heavy:
                dsn: '%env(MESSENGER_TRANSPORT_DSN)%'
                retry_strategy:
                    max_retries: 5
                    delay: 2000          # 2 sec
                    multiplier: 2        # 2, 4, 8, 16, 32 sec
                    max_delay: 60000
                options:
                    exchange:
                        name: app.async
                    queues:
                        app_heavy: ~
            failed:
                dsn: 'doctrine://default?queue_name=failed'

        routing:
            App\Upload\TranscodeVideo: async_heavy

        failure_transport: failed

Go реализация (channels pipeline)

Для in-process конвейеров Go даёт изящную реализацию через каналы. Для распределённого случая замените каналы на очередь (RabbitMQ, Kafka, NATS).

package upload

import (
	"context"
	"encoding/json"
	"net/http"

	"github.com/google/uuid"
)

type Server struct {
	jobs   JobRepo
	blob   BlobStore
	queue  chan<- Job // producer side of pipeline
}

// POST /uploads -- stage 1: accept file, persist blob + job, enqueue.
func (s *Server) Upload(w http.ResponseWriter, r *http.Request) {
	r.Body = http.MaxBytesReader(w, r.Body, 500<<20) // 500 MB limit

	file, header, err := r.FormFile("file")
	if err != nil {
		http.Error(w, "file missing", http.StatusBadRequest)
		return
	}
	defer file.Close()

	blobKey, err := s.blob.Put(r.Context(), file, header.Header.Get("Content-Type"))
	if err != nil {
		http.Error(w, "storage error", http.StatusInternalServerError)
		return
	}

	jobID := uuid.NewString()
	job := Job{
		ID:     jobID,
		Kind:   "video.transcode",
		Status: "queued",
		Input:  Input{BlobKey: blobKey, Size: header.Size},
	}
	if err := s.jobs.Create(r.Context(), job); err != nil {
		http.Error(w, "db error", http.StatusInternalServerError)
		return
	}

	// Enqueue for stage 2
	select {
	case s.queue <- job:
	case <-r.Context().Done():
		return
	}

	w.Header().Set("Location", "/jobs/"+jobID)
	w.Header().Set("Retry-After", "5")
	w.WriteHeader(http.StatusAccepted)
	_ = json.NewEncoder(w).Encode(map[string]string{"jobId": jobID, "status": "queued"})
}
## Backpressure

Если stage 2 не справляется с потоком:

Решение Плюсы Минусы
Bounded queue (отклонять 429) Простое Потеря запросов
Unbounded + автоскейл worker'ов Без потерь Тратится деньги, лаг
Две очереди (hot/cold) Приоритизация Сложно
Упаковывать по batch Эффективно Сложная реализация

В Go пример выше -- chan Job с буфером 100. При переполнении Upload блокируется. В HTTP handler это плохо -- лучше либо возвращать 429, либо иметь внешнюю очередь (RabbitMQ).

Progress reporting

Long-running jobs должны репортить прогресс. Варианты:

  1. Периодический UPDATE в БД из worker'а: progress = 45
  2. Redis PUBLISH каналу jobs:123:progress, клиент слушает через SSE
  3. Webhook на клиента при milestone (25%, 50%, 75%, 100%)

Транскодер ffmpeg парсит свой stderr и вызывает callback каждые ~5%.

Idempotency на приёме

Клиент нажал "upload" дважды -- не хотим два транскодинга одного файла:

POST /uploads
Idempotency-Key: <client-generated UUID>

Сервер хранит таблицу idempotency_keys (key, job_id). Если ключ уже есть -- возвращаем существующий jobId.

См. 02.design-approaches/9.idempotency.md для деталей.

Retention статусов

Job статус нужен N часов/дней (достаточно для клиента, чтобы забрать результат). Потом удалить.

-- Partition jobs by created_at, drop partitions older than 30 days
DELETE FROM jobs WHERE status IN ('done', 'failed') AND finished_at < NOW() - INTERVAL '30 days';

Failure modes

  • Worker упал во время обработки: job в статусе processing вечно. Watchdog: job в processing > N минут → вернуть в queued.
  • Клиент забросил job: job выполнился впустую. OK, но при дорогих операциях -- cancel token в input, worker проверяет jobs.cancelled = true.
  • Race между retry и done: worker обработал, закоммитил, нода упала до ack → retry. Идемпотентность MUST.
  • Результат слишком большой (> row size): blob storage, в jobs.result -- ссылка.

Мониторинг

  • jobs_queued (gauge), jobs_processing (gauge), jobs_done_total, jobs_failed_total
  • job_duration_seconds (histogram) -- per kind
  • job_wait_seconds -- время от создания до start
  • queue_depth -- длина очереди
  • worker_utilization -- % времени в работе

Алерты: queue_depth > threshold, jobs_wait_p95 > SLO, rate(jobs_failed_total[5m]) / rate(jobs_done_total[5m]) > 0.05.

Когда применять

  • Любая операция > 1 сек в HTTP хэндлере
  • Генерация отчётов, экспорт в CSV
  • Отправка emails (сразу не обязательно)
  • Image/video/pdf processing
  • Импорт/экспорт данных
  • ML inference с батчами

Когда не надо

  • Операция <100 ms -- ненужное усложнение
  • Жёсткий требование sync (например: валидация данных перед возвратом)
  • Нет инфраструктуры очередей -- начните с простого polling

Выводы

  • Two-stage processing: быстрый ack + async работа
  • HTTP 202 Accepted + Location header -- стандарт
  • Stage 1 обязан быть < 500 ms и idempotent
  • Stage 2 retryable, idempotent, с progress reporting
  • Polling -> WebSocket/SSE -> Webhook (по мере роста требований)
  • Watchdog для зависших job'ов обязателен
  • Метрики queue_depth и p95 wait -- ключевые