Базы данных

Построение надёжной очереди задач на PostgreSQL: SKIP LOCKED, партиционирование и мониторинг без Redis

Ruslan Ismailov Опубликовано 14 мин чтения
П

Введение: когда PostgreSQL лучше Redis или брокера сообщений

Большинство команд по умолчанию тянутся к Redis или RabbitMQ, когда речь заходит об очередях задач. Но в 2026 году PostgreSQL — это зрелая, battle-tested платформа с транзакционными гарантиями, которых у Redis нет из коробки. Если у вас уже есть PostgreSQL в стеке, добавление отдельного брокера означает: новую точку отказа, дополнительные операционные расходы, усложнение инфраструктуры и необходимость поддерживать согласованность между двумя хранилищами.

PostgreSQL queue оправдан в следующих сценариях:

  • Задачи тесно связаны с бизнес-данными и требуют атомарных операций с основной БД.
  • Нагрузка умеренная — до нескольких тысяч задач в секунду.
  • Важна гарантия exactly-once или at-least-once семантики с подтверждением через транзакцию.
  • Команда хочет упростить стек и избежать операционной сложности Redis Cluster.

В этой статье мы построим production-ready очередь задач без Redis: с SKIP LOCKED, партиционированием, воркерами на Go и PHP/Laravel, мониторингом через Grafana и деплоем в Kubernetes.

Паттерн Job Queue на PostgreSQL: схема, статусы, индексы

Основа паттерна — таблица заданий с явными статусами и пессимистичной блокировкой. Вот базовая схема:

CREATE TABLE jobs (
    id              BIGSERIAL PRIMARY KEY,
    queue           TEXT NOT NULL DEFAULT 'default',
    payload         JSONB NOT NULL,
    status          TEXT NOT NULL DEFAULT 'pending'
                        CHECK (status IN ('pending','processing','done','failed')),
    attempts        INT NOT NULL DEFAULT 0,
    max_attempts    INT NOT NULL DEFAULT 3,
    scheduled_at    TIMESTAMPTZ NOT NULL DEFAULT NOW(),
    locked_until    TIMESTAMPTZ,
    created_at      TIMESTAMPTZ NOT NULL DEFAULT NOW(),
    updated_at      TIMESTAMPTZ NOT NULL DEFAULT NOW()
) PARTITION BY LIST (queue);

CREATE TABLE jobs_default PARTITION OF jobs FOR VALUES IN ('default');
CREATE TABLE jobs_email   PARTITION OF jobs FOR VALUES IN ('email');
CREATE TABLE jobs_reports PARTITION OF jobs FOR VALUES IN ('reports');

-- Индексы для эффективного polling
CREATE INDEX idx_jobs_default_pending
    ON jobs_default (scheduled_at)
    WHERE status = 'pending';

CREATE INDEX idx_jobs_default_processing
    ON jobs_default (locked_until)
    WHERE status = 'processing';

Поле scheduled_at позволяет реализовать отложенные задачи. locked_until используется для детекта зависших воркеров: если воркер упал, задача автоматически возвращается в очередь по истечении TTL блокировки. payload типа JSONB даёт гибкость без изменения схемы под каждый тип задачи.

SELECT ... FOR UPDATE SKIP LOCKED: ключевой механизм

Главная проблема очередей на SQL — конкурентный доступ нескольких воркеров к одной задаче. Классический подход с SELECT + UPDATE в двух запросах создаёт race condition. SELECT ... FOR UPDATE SKIP LOCKED, появившийся в PostgreSQL 9.5, решает это элегантно и атомарно.

Принцип работы: когда воркер выполняет SELECT ... FOR UPDATE, PostgreSQL блокирует строку. Другие воркеры, выполняющие тот же запрос с SKIP LOCKED, просто пропускают заблокированные строки вместо того, чтобы ждать снятия блокировки. Это превращает PostgreSQL в эффективный распределённый lock-менеджер для очереди.

-- Атомарный захват задачи воркером
WITH next_job AS (
    SELECT id
    FROM jobs
    WHERE queue = 'default'
      AND status = 'pending'
      AND scheduled_at <= NOW()
    ORDER BY scheduled_at
    LIMIT 1
    FOR UPDATE SKIP LOCKED
)
UPDATE jobs
SET
    status       = 'processing',
    attempts     = attempts + 1,
    locked_until = NOW() + INTERVAL '5 minutes',
    updated_at   = NOW()
FROM next_job
WHERE jobs.id = next_job.id
RETURNING jobs.*;

Весь запрос выполняется в одной транзакции. Если воркер не подтвердит выполнение задачи в течение locked_until, отдельный reaper-процесс вернёт задачу в статус pending:

-- Возврат зависших задач (reaper)
UPDATE jobs
SET status = 'pending', updated_at = NOW()
WHERE status = 'processing'
  AND locked_until < NOW()
  AND attempts < max_attempts;

Реализация воркера на Go: polling, обработка, подтверждение

Go идеально подходит для написания воркеров: низкое потребление памяти, горутины для параллельной обработки, встроенный context для graceful shutdown. Используем pgx как драйвер PostgreSQL.

package main

import (
    "context"
    "encoding/json"
    "log"
    "time"

    "github.com/jackc/pgx/v5/pgxpool"
)

type Job struct {
    ID      int64
    Queue   string
    Payload json.RawMessage
}

func acquireJob(ctx context.Context, pool *pgxpool.Pool, queue string) (*Job, error) {
    row := pool.QueryRow(ctx, `
        WITH next_job AS (
            SELECT id FROM jobs
            WHERE queue = $1
              AND status = 'pending'
              AND scheduled_at <= NOW()
            ORDER BY scheduled_at
            LIMIT 1
            FOR UPDATE SKIP LOCKED
        )
        UPDATE jobs
        SET status = 'processing',
            attempts = attempts + 1,
            locked_until = NOW() + INTERVAL '5 minutes',
            updated_at = NOW()
        FROM next_job
        WHERE jobs.id = next_job.id
        RETURNING jobs.id, jobs.queue, jobs.payload
    `, queue)

    var job Job
    err := row.Scan(&job.ID, &job.Queue, &job.Payload)
    if err != nil {
        return nil, err
    }
    return &job, nil
}

func completeJob(ctx context.Context, pool *pgxpool.Pool, id int64) error {
    _, err := pool.Exec(ctx,
        `UPDATE jobs SET status = 'done', updated_at = NOW() WHERE id = $1`, id)
    return err
}

func failJob(ctx context.Context, pool *pgxpool.Pool, id int64) error {
    _, err := pool.Exec(ctx, `
        UPDATE jobs
        SET status = CASE WHEN attempts >= max_attempts THEN 'failed' ELSE 'pending' END,
            updated_at = NOW()
        WHERE id = $1
    `, id)
    return err
}

func runWorker(ctx context.Context, pool *pgxpool.Pool, queue string) {
    for {
        select {
        case <-ctx.Done():
            return
        default:
        }

        job, err := acquireJob(ctx, pool, queue)
        if err != nil {
            // Нет задач — ждём перед следующим poll
            time.Sleep(500 * time.Millisecond)
            continue
        }

        log.Printf("Processing job %d", job.ID)
        if err := processPayload(job.Payload); err != nil {
            log.Printf("Job %d failed: %v", job.ID, err)
            _ = failJob(ctx, pool, job.ID)
            continue
        }
        _ = completeJob(ctx, pool, job.ID)
        log.Printf("Job %d done", job.ID)
    }
}

func processPayload(payload json.RawMessage) error {
    // Бизнес-логика обработки задачи
    return nil
}

func main() {
    ctx, cancel := context.WithCancel(context.Background())
    defer cancel()

    pool, _ := pgxpool.New(ctx, "postgres://user:pass@localhost/db")
    defer pool.Close()

    // Запускаем несколько воркеров параллельно
    for i := 0; i < 5; i++ {
        go runWorker(ctx, pool, "default")
    }

    // Graceful shutdown по сигналу...
    select {}
}

Обратите внимание: при ошибке acquireJob из-за отсутствия задач воркер делает short sleep вместо busy-wait. В продакшне это значение можно вынести в конфигурацию и применять exponential backoff.

Реализация воркера на PHP/Laravel: кастомный драйвер

Laravel имеет встроенную поддержку PostgreSQL через драйвер database, но он использует advisory locks, а не SKIP LOCKED. Напишем кастомный драйвер, который использует правильный механизм блокировки.

<?php

namespace App\Queue;

use Illuminate\Queue\DatabaseQueue;
use Illuminate\Queue\Jobs\DatabaseJob;
use Illuminate\Support\Facades\DB;

class SkipLockedQueue extends DatabaseQueue
{
    public function pop($queue = null)
    {
        $queue = $this->getQueue($queue);

        return $this->getConnection()->transaction(function () use ($queue) {
            $job = $this->getNextAvailableJobWithSkipLocked($queue);

            if ($job !== null) {
                return new DatabaseJob(
                    $this->container,
                    $this,
                    $job,
                    $this->connectionName,
                    $queue
                );
            }
        });
    }

    protected function getNextAvailableJobWithSkipLocked($queue)
    {
        $job = DB::selectOne("
            WITH next_job AS (
                SELECT id FROM jobs
                WHERE queue = ?
                  AND status = 'pending'
                  AND scheduled_at <= NOW()
                ORDER BY scheduled_at
                LIMIT 1
                FOR UPDATE SKIP LOCKED
            )
            UPDATE jobs
            SET status = 'processing',
                attempts = attempts + 1,
                locked_until = NOW() + INTERVAL '5 minutes',
                updated_at = NOW()
            FROM next_job
            WHERE jobs.id = next_job.id
            RETURNING jobs.*
        ", [$queue]);

        return $job ? (object) $job : null;
    }
}

Регистрируем драйвер в AppServiceProvider:

<?php

// В AppServiceProvider::boot()
Queue::extend('pgsql_skip_locked', function () {
    return new SkipLockedQueueConnector();
});

В config/queue.php добавляем соединение:

'pgsql_skip' => [
    'driver'    => 'pgsql_skip_locked',
    'table'     => 'jobs',
    'queue'     => 'default',
    'retry_after' => 300,
],

Теперь php artisan queue:work --queue=default --connection=pgsql_skip использует нативный SKIP LOCKED PostgreSQL вместо advisory locks.

Партиционирование таблицы очереди: масштабирование и очистка

Таблица очереди без очистки растёт бесконечно. Партиционирование по queue (LIST partitioning, как в схеме выше) решает две задачи: изоляцию нагрузки между очередями и эффективную очистку через DROP PARTITION вместо DELETE.

Для очистки завершённых задач добавляем партиционирование по времени — архивную схему:

CREATE TABLE jobs_archive (
    LIKE jobs INCLUDING ALL
) PARTITION BY RANGE (created_at);

CREATE TABLE jobs_archive_2025_q1
    PARTITION OF jobs_archive
    FOR VALUES FROM ('2025-01-01') TO ('2025-04-01');

CREATE TABLE jobs_archive_2025_q2
    PARTITION OF jobs_archive
    FOR VALUES FROM ('2025-04-01') TO ('2025-07-01');

-- Переносим завершённые задачи в архив (cron-job)
INSERT INTO jobs_archive
SELECT * FROM jobs
WHERE status IN ('done', 'failed')
  AND updated_at < NOW() - INTERVAL '7 days';

DELETE FROM jobs
WHERE status IN ('done', 'failed')
  AND updated_at < NOW() - INTERVAL '7 days';

Удаление старой партиции архива происходит мгновенно и не создаёт нагрузки:

-- Удаление квартала без table lock
DROP TABLE jobs_archive_2025_q1;

Для highload-сценариев можно разбить основную таблицу по queue + диапазону id или использовать pg_partman для автоматического создания партиций.

Мониторинг: метрики, Prometheus, Grafana

Без мониторинга очередь — чёрный ящик. Ключевые метрики для PostgreSQL queue:

  • Queue depth — количество задач в статусе pending по каждой очереди.
  • Processing time — среднее и p95 время обработки задачи.
  • Failed rate — доля задач в статусе failed за последние N минут.
  • Stuck jobs — задачи в processing с истёкшим locked_until.

SQL-запросы для сбора метрик (exporter на Go или pg_stat_statements):

-- Глубина очереди по типам
SELECT queue, status, COUNT(*) as count
FROM jobs
GROUP BY queue, status;

-- Среднее время обработки (последние 10 минут)
SELECT queue,
       AVG(EXTRACT(EPOCH FROM (updated_at - created_at))) AS avg_processing_sec,
       PERCENTILE_CONT(0.95) WITHIN GROUP (
           ORDER BY EXTRACT(EPOCH FROM (updated_at - created_at))
       ) AS p95_processing_sec
FROM jobs
WHERE status = 'done'
  AND updated_at > NOW() - INTERVAL '10 minutes'
GROUP BY queue;

-- Зависшие задачи
SELECT COUNT(*) AS stuck_count
FROM jobs
WHERE status = 'processing'
  AND locked_until < NOW();

Собираем метрики через кастомный Prometheus exporter и визуализируем в Grafana. Пример конфига для дашборда: панель с графиком Queue Depth over Time, алерт при pending > 1000 дольше 5 минут, heatmap распределения времени обработки по очередям.

Для интеграции с Prometheus используем postgres_exporter с кастомными запросами через queries.yaml:

pg_job_queue_depth:
  query: |
    SELECT queue, status, COUNT(*) as count
    FROM jobs GROUP BY queue, status
  metrics:
    - queue:
        usage: LABEL
    - status:
        usage: LABEL
    - count:
        usage: GAUGE
        description: Number of jobs by queue and status

Сравнение с Redis Queues: плюсы и минусы

PostgreSQL queue — не серебряная пуля. Вот честное сравнение:

Преимущества PostgreSQL:

  • Транзакционность: задача и бизнес-данные изменяются атомарно в одной транзакции.
  • Нет дополнительного компонента в инфраструктуре — меньше точек отказа.
  • Полноценный SQL для аналитики, мониторинга, debugging.
  • Гарантии долговечности (WAL) из коробки без дополнительной настройки Redis AOF/RDB.
  • Партиционирование и индексы для эффективного управления большими очередями.

Ограничения PostgreSQL:

  • Пропускная способность ниже Redis: при нагрузке >10k задач/сек PostgreSQL начинает проигрывать.
  • Polling создаёт нагрузку на БД; при агрессивном polling это видно на CPU.
  • Нет pub/sub и fan-out из коробки — для этого нужен отдельный механизм (LISTEN/NOTIFY).
  • VACUUM нагрузка: интенсивные UPDATE статусов создают dead tuples.

Вывод: если у вас highload с тысячами задач в секунду — Redis/Kafka оправданы. Для большинства бизнес-приложений PostgreSQL queue — это более простое, надёжное и дешёвое решение.

Деплой воркеров в Kubernetes: Deployment vs Job

В Kubernetes воркеры очереди деплоятся как Deployment (не как Job), поскольку они должны работать непрерывно, а не выполнять одноразовую задачу. Пример манифеста для Go-воркера:

apiVersion: apps/v1
kind: Deployment
metadata:
  name: queue-worker
  namespace: production
spec:
  replicas: 3
  selector:
    matchLabels:
      app: queue-worker
  template:
    metadata:
      labels:
        app: queue-worker
      annotations:
        prometheus.io/scrape: "true"
        prometheus.io/port: "9090"
    spec:
      containers:
      - name: worker
        image: myapp/queue-worker:1.2.0
        env:
        - name: DATABASE_URL
          valueFrom:
            secretKeyRef:
              name: db-secret
              key: url
        - name: QUEUE_NAME
          value: default
        - name: WORKER_CONCURRENCY
          value: "5"
        resources:
          requests:
            cpu: 100m
            memory: 64Mi
          limits:
            cpu: 500m
            memory: 256Mi
        livenessProbe:
          httpGet:
            path: /healthz
            port: 9090
          initialDelaySeconds: 5
          periodSeconds: 10
      terminationGracePeriodSeconds: 60

Ключевые аспекты деплоя в Kubernetes:

  • terminationGracePeriodSeconds: даём воркеру время завершить текущую задачу перед остановкой пода. Воркер должен перехватывать SIGTERM и останавливать polling после завершения текущей итерации.
  • HPA по кастомной метрике: настраиваем горизонтальное масштабирование по глубине очереди через Prometheus Adapter — при росте pending задач автоматически добавляются реплики.
  • PodDisruptionBudget: гарантируем минимум 2 реплики при rolling update.
  • Kubernetes Job используем для reaper-процесса: запускаем его по расписанию через CronJob каждые 5 минут для возврата зависших задач.
apiVersion: batch/v1
kind: CronJob
metadata:
  name: queue-reaper
spec:
  schedule: "*/5 * * * *"
  jobTemplate:
    spec:
      template:
        spec:
          containers:
          - name: reaper
            image: myapp/queue-reaper:1.0.0
            env:
            - name: DATABASE_URL
              valueFrom:
                secretKeyRef:
                  name: db-secret
                  key: url
          restartPolicy: OnFailure

Заключение

Построение надёжной очереди задач на PostgreSQL — это вполне достижимая цель без привлечения дополнительных технологий. Комбинация SELECT ... FOR UPDATE SKIP LOCKED, грамотной схемы с партиционированием, воркеров на Go или Laravel и мониторинга через Prometheus/Grafana даёт production-ready решение для большинства бизнес-задач.

PostgreSQL queue в 2026 году — это осознанный архитектурный выбор в пользу простоты, надёжности и транзакционности. Начните с одной очереди, измерьте нагрузку, и только если вы упрётесь в потолок производительности — рассмотрите Redis или Kafka. Для 80% проектов этот потолок так и не наступит.

Технологии

Теги

Руслан Исмаилов

Senior Web / Backend разработчик. Senior web/backend разработчик с 9-летним опытом. Стек: PHP, Laravel, PostgreSQL, Redis, Docker, Kubernetes, REST, микросервисы, CI/CD. Подробнее обо мне →