Построение надёжной очереди задач на PostgreSQL: SKIP LOCKED, партиционирование и мониторинг без Redis
Введение: когда 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. Подробнее обо мне →