Когда очередь на Redis теряет задачи, сначала нужно выяснить, действительно ли сообщение исчезло. Оно может оставаться в pending entries list, ждать повторной попытки, находиться у остановившегося consumer, быть перенесенным в другую очередь или уже выполниться без обновления статуса в основной базе. Реальная потеря тоже возможна: worker удалил задачу до обработки, producer не завершил запись, Redis восстановился из старого snapshot или ключ очереди был удален политикой памяти.

Исправление зависит от используемой структуры. Redis List с LPOP, надежная List-схема с LMOVE, Redis Streams с consumer groups и готовый framework поверх Redis имеют разные подтверждения и механизмы восстановления. Поэтому нельзя начинать с увеличения timeout или перезапуска всех workers. Сначала фиксируется путь одной задачи от бизнес-транзакции до конечного результата.

Сначала остановите неконтролируемое удаление

Если потери продолжаются, временно ограничьте producer или опасный consumer, но не очищайте очередь. Сохраните метрики, конфигурацию persistence и снимок релевантных ключей. Массовый retry без идемпотентности может создать дубли писем, платежей, документов и внешних запросов.

  • зафиксируйте Redis endpoint, database number, key и тип структуры;
  • сохраните идентификатор одной пропавшей задачи и время постановки;
  • запишите producer, consumer group, consumer name и версию приложения;
  • проверьте длину основной, processing, retry и dead-letter очередей;
  • сохраните INFO persistence, memory, stats и replication;
  • не выполняйте FLUSHDB, DEL очереди, XTRIM или массовый XACK;
  • не перезапускайте Redis до проверки AOF/RDB и журнала запуска;
  • не публикуйте пароль Redis, connection string и содержимое персональных задач.

Если payload содержит персональные или платежные данные, выгружайте только технические поля: job_id, тип, состояние, retry count, timestamps и безопасные ссылки на бизнес-объект. Для диагностики не требуется копировать секреты из тела задачи.

Опишите жизненный цикл одной задачи

У надежной задачи есть не только состояние queued. Нужно различать создание бизнес-события, запись в Redis, выдачу worker, начало обработки, внешний побочный эффект, подтверждение и обновление основной базы. Разрыв между любыми двумя шагами выглядит как потеря, хотя причины различны.

business transaction -> enqueue requested -> message stored in Redis -> delivered to consumer -> processing started -> side effect completed -> message acknowledged -> business status updated

Для каждого шага нужен один correlation_id или job_id. Если producer пишет «задача поставлена» до ответа Redis, журнал создает ложную уверенность. Если worker пишет success до фиксации результата, задача выглядит выполненной, хотя транзакция могла откатиться. Логи должны отражать фактически завершенный шаг.

Определите модель очереди

Простой Redis List с LPOP или RPOP

Команда, которая сразу удаляет элемент из списка, дает семантику at-most-once. Если worker получил payload и завершился до выполнения, Redis уже не знает о задаче. Такая схема допустима только для работы, которую можно безболезненно потерять или восстановить из другого источника. Для важных задач она требует отдельного processing-списка.

List с LMOVE или BLMOVE

Надежная List-схема атомарно перемещает элемент из основной очереди в processing. После успешной обработки worker удаляет его из processing. Watchdog возвращает зависшие задачи по правилам visibility timeout. Но один список payload без времени выдачи и счетчика попыток усложняет безопасное восстановление, поэтому часто рядом хранится метаинформация.

Redis Streams с consumer groups

При XREADGROUP выданная запись попадает в Pending Entries List конкретной группы. Она не считается завершенной до XACK. Если consumer умер, запись остается pending и может быть передана здоровому consumer через XCLAIM или XAUTOCLAIM после минимального idle time. Поэтому «очередь пустая» при чтении только новых сообщений не означает, что задач нет: они могут находиться в PEL.

Проверьте, не застряли ли задачи в pending

Для Streams начните с состояния группы и PEL. Смотрите общее число pending, минимальный и максимальный ID, распределение по consumers и idle time. Большой pending при маленькой скорости подтверждений означает, что consumers получают задачи, но не завершают их либо подтверждают не в той группе.

XINFO GROUPS <stream> XINFO CONSUMERS <stream> <group> XPENDING <stream> <group> XPENDING <stream> <group> - + 20 # Пример переназначения только после выбранного idle timeout XAUTOCLAIM <stream> <group> <healthy-consumer> <min-idle-ms> 0-0 COUNT 50

Не запускайте XAUTOCLAIM с произвольно маленьким idle time. Живая задача может выполняться дольше и будет отдана второму worker, что создаст параллельную обработку. Значение выбирают по реальному времени p95/p99 плюс запасу, а длинные операции используют heartbeat или разбиваются на короткие идемпотентные шаги.

Убедитесь, что XACK выполняется после результата

Самая частая логическая потеря в Streams — ранний XACK. Worker подтверждает запись сразу после чтения, затем падает во время HTTP-запроса, записи файла или транзакции базы. Redis честно удаляет ссылку из PEL и больше не предлагает сообщение. Подтверждение должно происходить после надежной фиксации результата либо состояния, из которого обработку можно продолжить.

  • не подтверждайте запись в общей обертке до вызова handler;
  • проверьте finally-блоки, которые выполняют XACK и при исключении;
  • разделяйте retryable и permanent ошибки;
  • для permanent ошибки сначала запишите задачу в DLQ, затем подтвердите исходную;
  • проверяйте ответ Redis на XACK и правильность имени группы;
  • не объединяйте большой batch одним подтверждением до обработки всех элементов;
  • фиксируйте completed_at и result reference до ack.

Проверьте List-очередь на удаление до обработки

Если используется LPOP/BRPOP без processing-списка, задача действительно исчезает в момент выдачи. Переведите потребление на атомарное перемещение, но не меняйте направление списков наугад: порядок FIFO зависит от стороны добавления и чтения. Зафиксируйте текущую семантику и проверьте ее тестами.

Основная очередь -> атомарное перемещение -> processing Успех -> удаление из processing Временная ошибка -> retry с задержкой Постоянная ошибка -> dead-letter queue Смерть worker -> watchdog возвращает просроченную задачу

Watchdog должен различать медленную и умершую задачу, учитывать количество попыток и не возвращать один payload бесконечно. Без DLQ проблемная задача способна занять workers и создать впечатление, что новые сообщения пропадают.

Найдите разрыв между основной базой и Redis

Часто задача «теряется» еще до Redis. Приложение сохраняет заказ в SQL, затем отдельно отправляет job. Если процесс падает между операциями, заказ существует, а сообщения нет. Обратный порядок тоже плох: worker получает задачу до commit и не находит объект либо видит старые данные. Одна транзакция не может атомарно охватить обычную SQL-базу и отдельный Redis без специального протокола.

Используйте transactional outbox

В одной транзакции с бизнес-изменением запишите событие в outbox-таблицу. Отдельный publisher читает незавершенные строки, отправляет сообщения в Redis и отмечает публикацию. Если он падает после отправки, событие будет отправлено повторно, поэтому job_id и consumer должны быть идемпотентными. Такой подход заменяет риск потери контролируемым риском дубля.

BEGIN UPDATE orders ... INSERT INTO outbox(event_id, type, payload_ref, created_at) ... COMMIT publisher: outbox -> Redis -> mark published consumer: deduplicate by event_id -> process -> ack

Не храните в outbox лишние персональные данные. Достаточно event_id, типа, версии схемы и ссылки на бизнес-объект либо минимального неизменяемого payload.

Проверьте идемпотентность consumer

Надежные очереди обычно дают at-least-once: задача может прийти повторно после тайм-аута, failover или падения между побочным эффектом и ack. Нельзя одновременно гарантировать отсутствие потерь и считать любой повтор ошибкой транспорта. Consumer должен безопасно распознавать уже выполненную операцию.

  • используйте стабильный idempotency key из бизнес-события;
  • создавайте уникальное ограничение для необратимого результата;
  • сохраняйте статус started/completed с owner и временем;
  • передавайте idempotency key внешнему API, если оно это поддерживает;
  • не отмечайте локально completed раньше ответа и фиксации внешнего результата;
  • при неизвестном результате сначала выполняйте сверку, а не повтор;
  • делайте повторный handler возвращающим прежний результат.

Проверьте persistence Redis

Если Redis используется как источник важных задач, режим кеша без persistence не подходит. RDB хранит периодические снимки и при аварии может потерять записи после последнего snapshot. AOF журналирует операции; политика fsync определяет компромисс между скоростью и допустимым окном потери. При стандартном everysec авария может затронуть последние секунды, а режим always надежнее, но дороже по latency.

CONFIG GET appendonly CONFIG GET appendfsync CONFIG GET save INFO persistence INFO replication # Только диагностика: не меняйте CONFIG на production без плана

Проверьте фактический конфигурационный файл и параметры запуска, а не только CONFIG GET: временное CONFIG SET может исчезнуть после перезапуска. Убедитесь, что persistence-файлы находятся на постоянном volume, диск не заполнен, последние BGSAVE/BGREWRITEAOF завершились успешно и резервная копия действительно восстанавливается.

Репликация и failover не заменяют durability

Репликация Redis обычно асинхронна. Primary может подтвердить запись producer и завершиться до передачи replica; при failover новая primary не содержит эту задачу. Оцените допустимое окно, задержку реплики и настройки подтверждения записи. Даже при высокой доступности нужна идемпотентность, reconciliation и источник, из которого можно восстановить критичные события.

  • сопоставьте время потери с failover, restart и network partition;
  • проверьте master_repl_offset и отставание replica;
  • изучите журналы Sentinel, Cluster или управляющей платформы;
  • проверьте, куда подключались producers во время переключения;
  • учтите retry клиента после разрыва соединения;
  • не считайте read replica подходящим endpoint для постановки задач.

Проверьте maxmemory, eviction и TTL

Очередь может быть удалена Redis как обычный ключ, если выбранная eviction policy разрешает вытеснение и память закончилась. TTL, установленный общей оберткой кеша, также способен удалить stream или list целиком. Для инфраструктуры очередей обычно требуется отдельный Redis или политика, которая не вытесняет критичные ключи.

INFO memory INFO stats CONFIG GET maxmemory CONFIG GET maxmemory-policy TTL <queue-key> TYPE <queue-key> # В INFO stats проверьте evicted_keys и expired_keys

Не меняйте maxmemory-policy без оценки всех данных на экземпляре. Если кеш и очередь живут вместе, noeviction может вместо потери ключей начать отклонять новые записи. Producer обязан проверять ошибку записи и не сообщать об успешной постановке. Лучшее решение часто состоит в разделении кеша и критичной очереди.

Проверьте trimming Redis Streams

XTRIM и MAXLEN ограничивают размер stream, но агрессивное удаление может убрать payload раньше завершения обработки. В старых и смешанных версиях Redis особенно важно понимать связь stream и PEL: ссылка pending может остаться, а сама запись уже отсутствовать. Новые возможности управления ссылками нельзя использовать как универсальный совет без проверки версии сервера и клиентов.

  • сравните XLEN, длину PEL и возраст самой старой pending-задачи;
  • задайте retention больше максимального времени обработки и восстановления;
  • не trim по маленькому MAXLEN только ради снижения памяти;
  • учтите несколько consumer groups с разной скоростью;
  • проверьте версию Redis перед применением новых опций trimming;
  • храните DLQ и аудит отдельно от короткого рабочего stream.

Разберите ошибки worker

Задача может исчезать только визуально: handler ловит исключение, пишет короткое сообщение и подтверждает запись. Другой вариант — процесс завершается по OOM, timeout контейнера, SIGKILL или deploy, а recovery pending не настроен. Сопоставьте job_id с журналами приложения, orchestrator и операционной системы.

  • OOMKilled и лимиты памяти контейнера;
  • graceful shutdown и время на завершение текущей задачи;
  • автоматический ack в finally или middleware;
  • необработанная ошибка десериализации payload;
  • неверная версия схемы сообщения после deploy;
  • тайм-аут supervisor короче нормального выполнения;
  • исключение до записи job_id в журнал;
  • одинаковое consumer name у нескольких экземпляров, если framework требует уникальность;
  • потеря подключения и неправильная логика reconnect.

При graceful shutdown worker прекращает получать новые сообщения, завершает или безопасно оставляет текущую задачу pending и только затем выходит. Deployment без этого механизма регулярно создает зависшие сообщения и дубли.

Настройте retry и dead-letter queue

Бесконечный немедленный retry перегружает Redis и внешний сервис, а отсутствие retry превращает временный сбой в потерю. Для каждого типа ошибки задайте количество попыток, backoff с jitter и максимальное время. После лимита задача попадает в DLQ вместе с безопасной причиной и исходным job_id.

  • сетевой timeout — повтор после проверки неизвестного результата;
  • HTTP 429 — backoff с учетом Retry-After;
  • временный 5xx — ограниченные повторы;
  • ошибка валидации payload — DLQ без бесконечного retry;
  • отсутствующий бизнес-объект — короткий retry только при возможной задержке реплики;
  • неизвестная версия схемы — остановка и сигнал разработчику;
  • необратимый внешний эффект — сверка по idempotency key перед повтором.

Соберите метрики, которые отличают потерю от задержки

Одна длина очереди недостаточна. При быстром чтении и медленной обработке основная очередь равна нулю, а PEL растет. Измеряйте возраст старейшей задачи, время от enqueue до start и от start до completion, pending, retries, DLQ и расхождение между бизнес-событиями и завершенными jobs.

  • число поставленных, начатых, завершенных и подтвержденных задач;
  • queue lag и возраст старейшего сообщения;
  • pending по группам и consumers;
  • idle time и delivery count pending-сообщений;
  • retry rate и размер DLQ;
  • ошибки enqueue и reconnect;
  • evicted_keys, expired_keys и rejected writes;
  • состояние AOF/RDB и время последнего успешного сохранения;
  • replication lag и события failover;
  • reconciliation: события в базе без завершенной задачи.

Пошаговое безопасное исправление

  1. Выберите один пропавший job_id и восстановите его хронологию.
  2. Определите структуру очереди и фактический момент удаления или XACK.
  3. Проверьте processing/Pending Entries List, retry и DLQ.
  4. Устраните раннее подтверждение или неатомарное получение.
  5. Добавьте recovery зависших задач с реалистичным visibility timeout.
  6. Сделайте handler идемпотентным и защитите побочный эффект.
  7. Закройте dual-write через transactional outbox и reconciliation.
  8. Настройте persistence, постоянный volume и проверяемые backups по требуемой durability.
  9. Исключите eviction, TTL и агрессивный trim для критичных сообщений.
  10. Разверните изменение постепенно и наблюдайте метрики потерь и дублей.

Матрица проверок

Тестировать нужно не только нормальное выполнение, но и падение в каждом промежутке. Используйте тестовую среду и безопасные побочные эффекты. Не выключайте production Redis ради проверки.

  • worker падает сразу после получения сообщения;
  • worker падает после внешнего эффекта, но до ack;
  • producer падает после commit бизнес-транзакции;
  • publisher падает после отправки outbox-события, но до отметки published;
  • одно сообщение доставляется дважды;
  • задача выполняется дольше visibility timeout;
  • Redis временно недоступен при enqueue;
  • происходит controlled failover;
  • диск persistence заполнен;
  • payload имеет неизвестную версию или поврежден;
  • очередь достигает лимита памяти;
  • deploy останавливает worker с активной задачей.

Типичные ошибки

  • считать LLEN или XLEN единственным показателем очереди;
  • читать Streams только с символом > и не проверять pending;
  • выполнять XACK до фиксации результата;
  • использовать LPOP для важной работы без processing-очереди;
  • перезапускать workers вместо возврата зависших задач;
  • ставить слишком короткий visibility timeout;
  • повторять необратимую операцию без idempotency key;
  • считать репликацию гарантией отсутствия потерь;
  • хранить очередь в Redis без persistence как обычный кеш;
  • допускать eviction или TTL для ключа очереди;
  • агрессивно обрезать stream при незавершенных consumers;
  • писать SQL и Redis двумя независимыми операциями без outbox;
  • удалять ошибочные задачи вместо DLQ и анализа.

Как проверить результат после исправления

Надежность подтверждается сопоставлением счетчиков: каждое бизнес-событие имеет job_id, а каждая задача приходит к одному логическому результату, даже если физически доставлялась повторно. В тесте падения сообщение восстанавливается, а выполненный внешний эффект не дублируется. После restart и failover допустимое окно потери соответствует выбранной persistence-политике.

  • pending не растет бесконтрольно и старые записи восстанавливаются;
  • ack происходит только после сохранения результата;
  • повторная доставка возвращает прежний результат идемпотентно;
  • DLQ содержит объяснимые permanent ошибки;
  • enqueue failure виден producer и не маскируется как успех;
  • outbox reconciliation не находит необъяснимых пропусков;
  • очередь не имеет случайного TTL и не вытесняется;
  • AOF/RDB и backup проходят тест восстановления;
  • deploy корректно завершает активных consumers;
  • алерт срабатывает по возрасту, pending и расхождению счетчиков.

Как предотвратить повторение

Проектируйте очередь с допущением, что процессы падают, сеть разрывается, сообщения повторяются, а Redis может переключить primary. Для важных операций предпочтительна семантика at-least-once с идемпотентным consumer и возможностью восстановить событие из долговечного источника. У каждой очереди должны быть владелец, SLO по задержке и документированный сценарий восстановления.

  • используйте Streams consumer groups или надежную processing-схему;
  • вводите job_id и версию payload с момента создания;
  • добавляйте outbox для связи с бизнес-транзакцией;
  • настраивайте retries, DLQ и reconciliation;
  • проверяйте persistence и backups восстановлением, а не наличием файлов;
  • отделяйте критичные очереди от вытесняемого кеша;
  • тестируйте crash, deploy и failover регулярно;
  • наблюдайте не только длину, но и возраст и состояния задач.

Когда нужна помощь с Redis-очередью

Если очередь на Redis теряет задачи, я могу восстановить путь сообщений по job_id, проверить List или Streams, PEL, момент ack, retries, DLQ, persistence, eviction и failover. Затем настрою надежное получение, идемпотентную обработку, recovery зависших задач и transactional outbox без очистки рабочей очереди. Для первичной оценки достаточно обезличенных логов, типа ключа, схемы consumer и результатов INFO без пароля и содержимого секретных payload.