Проблема: долгие операции в HTTP-запросе
Пользователь загружает видео. Сервер должен:
- Сохранить файл
- Транскодировать в 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'],
]);
}
}
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"})
}
Если stage 2 не справляется с потоком:
| Решение | Плюсы | Минусы |
|---|---|---|
| Bounded queue (отклонять 429) | Простое | Потеря запросов |
| Unbounded + автоскейл worker'ов | Без потерь | Тратится деньги, лаг |
| Две очереди (hot/cold) | Приоритизация | Сложно |
| Упаковывать по batch | Эффективно | Сложная реализация |
В Go пример выше -- chan Job с буфером 100. При переполнении Upload блокируется. В HTTP handler это плохо -- лучше либо возвращать 429, либо иметь внешнюю очередь (RabbitMQ).
Progress reporting
Long-running jobs должны репортить прогресс. Варианты:
- Периодический UPDATE в БД из worker'а:
progress = 45 - Redis PUBLISH каналу
jobs:123:progress, клиент слушает через SSE - 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_totaljob_duration_seconds(histogram) -- per kindjob_wait_seconds-- время от создания до startqueue_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 -- ключевые