title: Статья 07 — Асинхронная обработка: очереди и потоки source: подготовлено 2026-09-17, TASK-45.47; каркас — brief.md §07; развёртка knowledge-base.md §8 (мосты в §7, §9, §10); источники — DDIA 2-е изд. гл. 11–13 (TASK-45.39), статья Яндекса (TASK-45.5), эталоны classic-designs.md §2–4
Статья 07: асинхронная обработка — очереди и потоки
Седьмая статья серии — про то, что происходит вне request-path: всё, что не нужно для ответа пользователю, — забота очередей и потоковой обработки. Критерий «умею, если» из брифа: объясняю, что выношу из request-path в очередь и что делаю при потере/дубле сообщения. Это ровно два вопроса, и на секции их задают именно так: сначала решение «что → очередь», затем гарантии «что будет, если сообщение потеряется или придёт дважды». Обратите внимание: критерий про ваше обоснование выбора, а не про знание Kafka — поэтому §2 (критерий выноса) и §6–7 (потеря и дубль) здесь главные, а словарь (§3–4) — только опора для них.
Место раздела в каркасе: статья 03 на лестнице масштабирования назвала очередь «шагом 9» — ответом на пики нагрузки, роняющие синхронный путь (03-scaling.md §2), и обещала развернуть это здесь; статья 02 уже посчитала fan-out в 700k вставок/с (02-estimations.md §6); статья 04 сослалась сюда на outbox/CDC для синхронизации проекций; статья 06 — на событийную инвалидацию кэша через шину. Всё это сходится в один раздел: очередь — механизм развязки, буферизации и ретраев; потоковая обработка — как за очередью что-то делается. Соседи с другой стороны: поведение консьюмера при отказе и повторном вызове — раздел 08 (таймауты, ретраи, идемпотентность уже анонсированы там §3), гарантии между системами — раздел 09 (outbox вместо 2PC, консенсус координации), мониторинг лага — раздел 12 (golden signals). Глубина — ../knowledge-base.md §8; полные разборы — ../classic-designs.md §2–4; DDIA — гл. 11–13.
1. Зачем очередь: что покупаем и чем платим
Очередь — не «ещё один блок на схеме», а сознательный обмен: четыре выигрыша против трёх плат. Выигрыши (канон КБ §8):
- Развязка сервисов. Производитель и потребитель не знают друг о друге: не делят транзакцию, не синхронизируют деплои, не держат общую БД. Пиксель архитектуры: событие → очередь → N независимых групп потребителей (аналитика, уведомления, индексация). Плата — событийный контракт: схема события становится API, эволюция схемы — отдельная забота (DDIA гл. 5: обратная совместимость).
- Буферизация пиков (load leveling). Пиковая волна (рекламная кампания, «горячий автор», суточный всплеск загрузок) встаёт в очередь и размазывается по времени. Потребители работают на средней, а не на пиковой нагрузке — ферма воркеров меньше и стабильнее. Классическая фраза из эталона медиа-хостинга: «задержка обработки вместо роста фермы».
- Ретраи. Упавший потребитель не теряет работу: сообщение остаётся неподтверждённым и переходит к следующей попытке. В синхронном мире один зависший dependency съедает потоки; в асинхронном — просто откладывается.
- Broadcast и порядок. Одна запись → многие читатели (паб-саб) или строгий порядок для одного ключа (партиция) — то, что синхронно делать либо дорого, либо невозможно (DDIA гл. 12).
Плата — три вещи, которые нужно называть самому, не дожидаясь вопроса:
- Лаг. Между событием и его обработкой — задержка от миллисекунд (горячий контур) до минут (батчи). Всё, что ушло в очередь, перестаёт быть read-your-writes для других потребителей.
- Дубли. Надёжная доставка вообще не бывает без дублей (at-least-once, §5) — каждый потребитель обязан быть идемпотентным. Это не «мелочь», а постоянная цена, зашитая в каждый консьюмер.
- Сложность отладки. Локально не воспроизводится: порядок, гонки, ретраи, «а где это сообщение?» — большая часть инцидентов живёт в асинхронной части.
Формулировка-канва всего раздела: «в request-path — минимум, без которого пользователь не получит ответ; всё остальное — асинхронно через очередь, и тогда каждая асинхронная зависимость обязана отвечать на два вопроса: что при потере и что при дубле».
2. Что выношу из request-path в очередь (критерий)
Первый вопрос брифа. Правильный ответ — не список, а критерий решения, и уже на его основе — примеры. Критерий двушаговый:
- Нужен ли результат для ответа пользователю? Если действие не входит в критический путь ответа — кандидат на вынос. «Пользователь опубликовал пост»: лента подписчиков, счётчики, уведомления, поиск — ничего из этого не нужно для ответа «пост создан» → всё фоном. «Пользователь платит»: подтверждение платежа — часть ответа → остаётся синхронным.
- Допустима ли задержка и готова ли система к eventual? Всё, что вынесено, видно не сразу. Кому-то (другим пользователям) — норма, окно в секунды; самому автору — «read-your-writes» нужен для своих данных и делается не очередью, а синхронным попаданием в свой кэш/свою ленту (§13, эталон ленты). Если задержка недопустима — она не выносится, а решается кэшем/репликой (06, 05).
Четыре класса «выношу» — по ситуациям из брифа и эталонов:
| Что | Пример | Почему | Где в серии |
|---|---|---|---|
| Тяжёлая работа (минуты CPU, внешние вызовы) | транскодинг видео, генерация отчёта, отправка рассылки | не помещается в таймаут ответа (сек–десятки сек) | медиа-хостинг, shortener |
| Fan-out на N читателей | пост → ленты 300 подписчиков, 700k вставок/с у селебрити | запись дёшева, доставка дорога — не держать в API | лента (§11) |
| Побочные эффекты записи | счётчики просмотров, аналитика, индекс поиска | ответ не зависит от результата; точность ±% приемлема | медиа, shortener |
| События для других сервисов | OrderPlaced для уведомлений/аналитики/антифрода, инвалидация кэша |
развязка + независимое масштабирование | 06 §5, весь раздел |
Правило-ограничитель из эталона медиа: пост откладывать до готовности медиа не будем — трейдофф «свежесть поста > полная комплектность»: в ответе полный факт, фон догоняет. И наоборот — антипример, встречающийся на секциях: «а давайте всё сделаем асинхронно» — вынос без этого критерия меняет read-your-writes и предсказуемые таймауты на буфер, который при пиках всё равно переполняется (§9).
Готовый ответ на первый вопрос брифа: «из request-path выношу всё, что не нужно для ответа пользователю: тяжёлую работу, fan-out, побочные эффекты и события для других сервисов; критерий — "нужен результат в этом ответе?" и "допустимо ли окно eventual?"; всё вынесенное обязано переживать потерю и дубль — вот как» — и переход ко второму вопросу.
3. Два вида брокеров и словарь Kafka
Выбор очереди на секции — одна фраза, но за ней два класса, которые нужно различать (DDIA гл. 12):
- Классические брокеры (RabbitMQ/AMQP): сообщение живёт, пока не подтверждено (ack), после — удаляется (delete-after-delivery). Брокер сам решает, кому отдать сообщение. Плюс: точная доставка «одному из», приоритеты, TTL сообщений. Минус: нельзя перечитать выданное — нет истории, нет реплея; «брокер как конечная точка», а не как журнал.
- Log-based брокеры (Kafka): сообщения не удаляются — пишутся в append-only партиции и живут до истечения retention (по времени или размеру). Потребитель сам хранит позицию (offset) и может перечитать с любого места. Плюс: история, реплей, независимые потребители с разной скоростью. Минус: порядок — только внутри партиции, а не в топике (§8); «удалить одно сообщение» нельзя в принципе.
Выбор по сути: «нужна работа "отдал и забыл" с точным контролем — RabbitMQ; нужен журнал изменений, многие независимые читатели, реплей, потоковая обработка — Kafka». В наших системах дефолт — Kafka.
Словарь Kafka, который на секции обязан быть на языке (КБ §8):
- Topic — именованный поток событий. Партиции — единицы параллелизма и порядка: порядок гарантируется внутри партиции, параллелизм — числом партиций. Ключ маршрутизирует:
key % partitions(или по хешу) — одинаковый ключ всегда в одной партиции. - Consumer group — группа консьюмеров, делящих топик: каждая партиция отдаётся одному консьюмеру группы; консьюмеров может быть и больше партиций — лишние просто простаивают. Разные группы — разные читатели, каждая получает всё (fan-out, §4).
- Offsets — позиция консьюмера в партиции; коммитится (автокоммит или вручную после обработки — §6) и позволяет реплей с любого места. Retention — сколько хранить (например, 7 дней или 1 ТиБ); лог можно перечитать целиком.
- ISR (in-sync replicas) и
acks=all— подтверждение записи брокером после репликации; min.insync.replicas=2 — минимум живых реплик для приёма записи. Это и есть «сообщение не потеряно» на стороне брокера (§6).
Произносимая формула словаря: «Kafka: топик из партиций — порядок внутри партиции, параллелизм по партициям; консьюмеры в группе делят партиции, разные группы читают всё; offset — позиция и возможность реплея; репликация партиций и acks=all — защита от потери» — 20 секунд, закрывает и выбор, и устройство.
4. Паб-саб: как устроен broadcast
Вопрос брифа «паб-саб». Общая модель: производитель (publisher) не знает потребителей — событие уходит в брокер, дальше маршрутизация. Два родительских паттерна доставки (DDIA гл. 12):
- Load balancing (конкуренция): сообщение получает один из группы потребителей — «работа» распределяется. В Kafka — одна consumer group на топик: партиция → один консьюмер, группа потребляет все партиции параллельно; в RabbitMQ — конкурирующие очереди (competing consumers).
- Fan-out (broadcast): каждый независимый читатель получает все сообщения. В Kafka — несколько consumer groups на один топик (аналитика и индексация читают одно и то же независимо, со своей скоростью и своим offset); в RabbitMQ — fanout/topic-обменники + отдельные очереди на подписчика.
Разница, которую стоит произнести явно: «одна группа — делим работу, разные группы — каждый получает всё». Это и есть ответ на «как сделать паб-саб на Kafka»: подписка = своя consumer group.
Прикладной момент: «fan-out» в бизнес-смысле (раскладка постов по лентам) и в смысле доставки — разные вещи, не путать. В ленте новостей fan-out-воркеры — это одна consumer group, делящая работу по раскладке (load balancing: партиция → один воркер); а вот индексация поиска и аналитика — другие группы, читающие тот же топик целиком (каждая получает всё). Один топик → три независимых пайплайна: один — распределение работы, два — подписка на всё. Это готовая схема «как из одного события живут многие проекции» — DDIA гл. 13 называет это интеграцией из единого журнала.
5. Семантики доставки: at-most-once, at-least-once, effectively-once
Три семантики — ядро ответа на второй вопрос брифа (КБ §8, DDIA гл. 12):
| Семантика | Что гарантирует | Что может случиться | Когда годится |
|---|---|---|---|
| At-most-once | доставим не более одного раза (fire-and-forget) | может потерять | метрики, телеметрия, где потеря лога не страшна |
| At-least-once | доставим, пока брокер жив; будет повторён, пока не подтверждён | может дублировать | дефолт надёжных систем: платёжный контур, лента, медиа |
| Effectively-once | at-least-once + дедупликация на стороне потребителя | дублей «не видно» | всё, где дубль стоит денег (списания, заказы) |
Ключевая фраза, которую на секции нужно произнести дословно: «exactly-once между системами не существует — есть at-least-once плюс идемпотентный потребитель; effectively-once — это не свойство брокера, а правило приложения». Почему: между производителем и потребителем всегда есть две независимые системы (БД и брокер, брокер и потребитель) с собственными транзакциями — единственное «exactly-once» строится глубоко внутри одного процесса (Kafka Streams + transactions) и не переносится через границу систем. На секции это отличает знающего: отвечающий «у Kafka есть exactly-once» — заваливает вопрос; отвечающий «effectively-once = at-least-once + дедуп» — проходит.
Логика выбора между тремя — по цене потери и цене дубля: телеметрию можно терять (at-most-once); ленту нельзя терять и можно дублировать (at-least-once + идемпотентная вставка); деньги нельзя ни терять, ни дублировать (effectively-once, §7). На схеме обязательно назвать семантику для каждой очереди: интервьюер это проверяет.
6. Потеря сообщения: по звеньям пути
Второй вопрос брифа. Ответ строится прогоном по пути сообщения и по звеньям: где теряется и что делаем на каждом. Три звена — три защиты, и это нужно выдать одним потоком.
Звено 1 — производитель: событие может не дойти до брокера. Публикация события — второй ход рядом с записью в БД; «записал данные и отправил событие» — две операции, между ними гонка: упали между — данные есть, события нет. Защита — transactional outbox (§10): событие пишется в outbox-таблицу в той же транзакции, что и данные, отдельный публикатор доставляет его в Kafka. Единственный способ сделать «данные записаны» и «событие опубликовано» атомарными без распределённой транзакции. Плюс на стороне producer: acks=all + ретраи + идемпотентный продюсер (возврат дублей при ретрае, §7).
Звено 2 — брокер: сообщение может физически потеряться. Репликация партиций: RF=3 (три копии), min.insync.replicas=2 — запись принимается, только когда её увидели две реплики; лидер-фейловер прозрачен (об этом — таблица отказов статьи 08, строка «Очередь (Kafka)»: under-replicated партиции, реплики в разных брокерах/ДЦ). Произносимая фраза: «на брокере сообщение не теряю тройной репликацией партиций с подтверждением от двух реплик».
Звено 3 — потребитель: обработал и упал до коммита offset. Пока offset не закоммичен, сообщение будет перечитано (at-least-once — это здесь). Правило: коммит offset после обработки (а не до), обработка идемпотентная (§7), ретраи с backoff, DLQ для ядовитых (§9), мониторинг lag. Единственная «могила» потери на этом звене — коммит до обработки: это категорическое «нет».
И финальный слой — реплей: благодаря retention и offsets любое «потерялось/испортилось» лечится перечтением лога с нужного места (пересборка проекций, восстановление индекса). Формулировка: «потеря закрывается по звеньям: outbox на производителе, репликация на брокере, коммит-после-обработки и DLQ на потребителе; плюс реплей как страховка» — это и есть полноценный ответ «что делаю при потере сообщения».
7. Дубль сообщения: идемпотентность потребителя
Вторая половина вопроса брифа. Дубль — норма (at-least-once), значит потребитель обязан делать одно и то же при повторном получении. Три уровня защиты, от простого к строгому:
- Естественная идемпотентность операции.
UPDATE ... SET status = 'paid' WHERE id = X— повтор безопасен;INSERT ... ON CONFLICT (id) DO NOTHING— безопасен. Пока операция идемпотентна по построению — дубль невидим. Годится для вставки постов, апдейтов статусов, дедупликации по естественному ключу. - Дедуп по уникальному ключу события. У каждого события — свой ID (
message_id/ дедуп-ключ). Потребитель в той же локальной транзакции, что и бизнес-эффект, вставляет ID в таблицу обработанных (уникальный индекс); дубль отклоняется constraint'ом. Это паттерн Kafka Streams и самая ходовая реализация effectively-once: «сообщение обработано» и «ключ записан» — один коммит. - Idempotency key — когда у операции нет естественного ключа и она не идемпотентна (списание денег): клиент/производитель генерирует
operation_id, сервер хранит результат по ключу и на повтор отдаёт сохранённый ответ вместо повторного исполнения (мост в статью 08 §5 — там это общий инструмент ретраев).
Фраза-мостик из статьи 08, которую здесь повторяем дословно: «транспорт дублирует, приложение дедуплицирует: at-least-once превращаю в effectively-once идемпотентным ключом». На секции: вопрос «а если сообщение придёт дважды?» — самый вероятный из раздела; ответ — «идемпотентный потребитель: либо операция идемпотентна по построению, либо дедуп по message_id в локальной транзакции, либо idempotency key для неидемпотентных» — с примерами на месте (§11).
Отдельным штрихом — случай «не дубль, а потеря на реплике»: at-least-once не спасает от расщепления сети (брокер подтвердил, но реплики разошлись) — для денежных контуров добавляют сверку/аудит-поток. Упоминание одной фразой — маркер глубины.
8. Порядок: партиция, ключ, причинность
Третий скрытый вопрос раздела — «а порядок?». Порядок — самый дорогой ресурс асинхронности, и его осознанно ограничивают:
- Порядок существует только внутри партиции. Все сообщения с одним ключом — в одной партиции и обрабатываются по порядку (Kafka). Поэтому ключ партиции = сущность домена, порядок в которой важен:
dialog_id(мессенджер — порядок сообщений в диалоге),user_id(операции одного пользователя),order_id(статусы заказа). Консьюмер одной группы на партицию — порядок сохраняется на всём пути. - Глобального порядка нет и не строим. DDIA гл. 13 прямо: глобальный total ordering недостижим и не нужен; упорядочиваем по причинности — где события связаны причиной, там же один ключ. Фраза из эталона мессенджера: «строгость нужна только внутри диалога — глобальный порядок сообщений не нужен и не строится».
- Плата за подтверждённый порядок — параллелизм: одна партиция = один поток обработки. Горячий ключ (диалог на 10k сообщений/с) упирается в один консьюмер — лечение как у горячего ключа кэша (06 §6): размазать по времени, буферизовать, поднять ключ на уровень выше (партиционирование по более грубой сущности:
dialog_id→dialog_id % Nс сохранением порядка внутри поддиалога, если бизнес позволяет).
Произносимая формула: «порядок — внутри партиции; ключ — доменная сущность, порядок в которой критичен; параллелизм покупаю числом партиций, а не беспорядком; глобальный порядок не строю, строю причинность».
9. Backpressure и перегрузка: лаг, DLQ
Очередь — буфер, а не бездонная яма; вопрос «а если потребители не успевают?» — обязательный на секции. Словарь:
- Лаг (consumer lag) — отставание потребителя от хвоста партиции: главная метрика очереди (§12) и главный индикатор перегрузки. Лаг растёт → первым делом говорим о нём.
- Load leveling — то, зачем очередь вообще стоит (§1): пики встают в буфер, потребители работают на средней нагрузке. Пока лаг конечен и убывает — всё работает, это и есть назначение.
- Backpressure и разгрузка перегрузки, когда лаг не убывает: (1) competing consumers — добавить консьюмеров в группу, но не больше числа партиций (иначе лишние простаивают: партиция → один консьюмер); (2) шардировать топик на больше партиций (с выбором нового ключа — порядок при этом может переразложиться, §8); (3) rate limit на входе / load shedding — не давать очереди расти до retention; (4) приоритетные очереди — отдельные топики для срочного (платежи) и фонового (аналитика); (5) dead man's switch — алерт на аномальное затишье: поток сообщений давно молчит — пайплайн мог застрять, а не нагрузка кончилась.
- Долгие операции потребителя — не держать транзакцию во время внешнего вызова: читать, ставить в работу, коммитить по факту (или выделять под это отдельную очередь с большим таймаутом).
- Ядовитые сообщения (poison messages) — сообщение, которое при каждой попытке валит потребителя (битый JSON, невалидные данные): без защиты оно блокирует партицию, потому что порядок. Защита — DLQ (dead-letter queue): после N неудачных попыток сообщение уходит в отдельный топик/очередь, годный к ручному разбору, а основной поток продолжает. Плюс ретраи с экспоненциальным backoff + jitter (мост в 08 §1: ретраи без jitter = retry storm).
Формулировка перегрузки: «лаг — моя главная метрика: пока он конечен, очередь работает как буфер пиков; когда не убывает — добавляю консьюмеров до числа партиций, шардирую, режу вход rate limit'ом, ядовитое — в DLQ» — 20 секунд, закрывает «а если не успевают?».
10. Outbox и CDC: как событие уходит без потери
Два механизма, которые в этой серии уже анонсировали статья 04 (синхронизация проекций, §4) и статья 06 (событийная инвалидация кэша, §5) — здесь они раскрываются, потому что это ответы на «как гарантировать публикацию».
Transactional outbox. БД и брокер — две системы; «записал данные» и «опубликовал событие» нельзя сделать атомарно обычным вызовом брокера (двойная запись = гонка и потеря). Решение: в той же транзакции, что и данные, пишем событие в таблицу outbox; отдельный публикатор (воркер, опрашивающий таблицу, или CDC) доставляет его в Kafka и помечает прочитанным. Событие не потеряется: оно либо в БД вместе с данными, либо уже в логе. Идемпотентность публикации — по id события. Это стандартный ответ на «как надёжно опубликовать событие» (KB §7; мост в 09: outbox вместо 2PC).
CDC (change data capture). Репликационный журнал БД (binlog/WAL) → поток изменений: Debezium/Kafka Connect читает лог БД и публикует факты изменений (строка вставлена/обновлена/удалена). БД не знает о подписчиках; историческая загрузка — initial snapshot + смещение; log compaction — хранить только последнее значение по ключу (навсегда), а не окно по времени — таблица-состояние поверх лога. CDC — то же «одно событие на запись», только источник — движок БД, а не наш код.
Различение, которое на секции превращает «знаю термины» в «понимаю»: CDC — факты изменений («что произошло со строкой»), event sourcing — намерения домена («что произошло в бизнесе»: OrderPlaced, MoneyTransfered). Первое — зеркало БД; второе — прикладные события. Наши системы строятся на прикладных событиях через outbox; CDC — инструмент подпитки проекций и инвалидации кэша по данным БД.
11. Batch vs stream: окна, event time, унификация
Последний пункт брифа — и здесь DDIA гл. 11–13 даёт парную картину:
- Batch — ограниченный вход (bounded), обрабатываем за один проход: сводка за день, пересборка индекса, обучение модели. Свойство — детерминированность и перезапускаемость: всегда можно пересчитать заново (reprocess) → производные данные заново строятся из лога. Запаздывает: изменение входа доходит до выхода не раньше конца окна.
- Stream — неограниченный вход (unbounded), обрабатываем по мере поступления: лента, алерты, антифрод. Свежесть в секунды, но точность — вопрос окон (§11 ниже).
Три механизма потоковой обработки, которые нужно уметь назвать (DDIA гл. 12): окна — tumbling (фиксированные, непересекающиеся), sliding (скользящие), session (по паузам активности); event time vs processing time — время в событии (часы клиента) против времени обработки: часы клиентов плывут, события приходят с опозданием → watermarks — границы доверия, после которых окно можно закрывать и считать завершённым. Фраза: «считаю по event time с watermark, иначе агрегаты зависят от скорости доставки».
И главный тренд-мостик: batch = поток с концом; batch и stream унифицируются (Spark Structured Streaming, Flink, Beam — одна модель, batch как частный случай). Поэтому «batch vs stream» на секции — не выбор, а шкала: свежесть против стоимости пересчёта. Практическое правило: дешёвый ответ «использую оба: стрим для горячего (лента, алерты), батч для тяжёлого пересчёта (индекс, рекомендации); и то и другое читает один лог и при эволюции схемы просто пересобирается» (DDIA гл. 13: все производные — проекции одного журнала).
12. Отказ очереди и метрики (мост в раздел 08)
Статья 08 разбирает отказ очереди в своей таблице отказов; здесь — три факта, которые нужно знать до того, как туда переходить:
- Асинхронность сама по себе — механизм отказоустойчивости: отказ воркера = рост очереди, а не падение сервиса; пользователю «просто дольше». Именно поэтому пайплайны (транскодинг, fan-out) проектируют асинхронными не только ради развязки.
- Отказ брокера (нода/партиция) лечится репликацией: RF=3,
min.insync.replicas=2, лидер-фейловер прозрачен (05 §5 — тот же механизм, что у БД). Полную остановку кластера переживают стороны по-разному: продюсеры буферизуют на диске, потребители видят лаг; DLQ/ретраи — по месту. - Метрики очереди = лаг и скорость обработки (golden signals раздела 12):
consumer lagна партицию (алерт на порог), старость самых старых сообщений, under-replicated партиции, число DLQ-сообщений. Фраза для этапа 6: «очередь видна по лагу и DLQ, а не по сайзенду» — и это же связка с канарейкой: regression доставки видна в лаге до полного выката (эталон ленты: canary на fan-out-воркерах с контролем lag).
13. Прикладное: как это звучит в эталонах
Четыре готовые формулировки из ../classic-designs.md — заготовки ответов «что выносите в очередь?» и «гарантии?»:
URL shortener (Яндекс, §1). Счётчики переходов и актуализация времени доступа — не в горячем пути: «не нужно для ответа пользователю → из горячего пути — в память воркера, батчами во внешний мир» (причём очередь — на базе существующего KV, без отдельного брокера: KISS из эталона). Железное правило: data plane никогда не пишет в СУБД напрямую, только через очередь.
Лента новостей (§2). PostCreated → Kafka (ключ = author_id); fan-out-воркеры одной группой раскладывают пост в кэши лент 300 подписчиков; селебрити — не пушим, merge на чтении (fan-out on read), иначе 700k вставок/с. Очередь сглаживает пик, lag — главная метрика (алерт на порог > X с); отказ воркера = лента стареет, graceful degradation (08). Фраза: «запись поста дёшева, доставка дорога — доставка это и есть дизайн, очередь + шардированные воркеры».
Мессенджер (§3). MessageSaved → Kafka, партиции по dialog_id — порядок внутри партиции = порядок диалога; последовательность seq назначает хранилище при вставке; at-least-once + дедуп по (dialog_id, seq) на клиенте; «строгость только внутри диалога — глобальный порядок не строится» (§8). Метрика: lag per партиция.
Медиа-хостинг (§4). Upload → событие в очередь → транскодеры → CDN; at-least-once + идемпотентные воркеры + DLQ для битых видео; очередь сглаживает суточный пик загрузок — «задержка обработки вместо роста фермы» (SLA «ready ≤ 30 мин» определяет размер фермы, не наоборот); пайплайн — не одна очередь, а DAG специализированных воркеров с очередями между этапами. Пост не откладывается до готовности медиа — плейсхолдер, фон догоняет (§2).
Сборка ответа на вопрос «что выносите в очередь»: «уведомления и тяжёлую работу — всегда; fan-out — лента (группой воркеров); счётчики/аналитику — батчами из горячего пути; события для проекций — через outbox; гарантии: at-least-once + идемпотентные консьюмеры + DLQ, потеря закрыта по звеньям» — 30 секунд, четыре класса из §2 и гарантии из §6–7.
14. Ошибки этого раздела
Подмножество топ-10 ошибок (../methodology.md §4) и анти-паттернов источников:
| Ошибка | Как выглядит | Противоядие |
|---|---|---|
| Очередь «на всякий случай» | «Сразу микросервисы на Kafka», коробки без боли (03 §2) | Вынос — по критерию §2: «нужен результат в ответе?» |
| Dual write без outbox | «Записал в БД и отправил в Kafka» двумя вызовами | Outbox в той же транзакции (§10) |
| «У нас exactly-once» | Упор на брокера вместо приложения | Effectively-once = at-least-once + дедуп (§5, §7) |
| Коммит offset до обработки | Потеря при падении потребителя | Коммит после обработки + ретраи (§6) |
| Потребитель не идемпотентен | Дубль списывает дважды / раскладывает пост дважды | Дедуп по message_id в локальной транзакции (§7) |
| Порядок «глобально» или «никак» | Держимся за total order или теряем порядок диалога | Порядок внутри партиции, ключ = доменная сущность (§8) |
| Консьюмеров больше партиций | Лишние простаивают, лаг не убывает | Число консьюмеров ≤ партиций; топик шардируем (§9) |
| Нет DLQ | Одно битое сообщение блокирует партицию | DLQ после N попыток + backoff/jitter (§9) |
| Очередь как бесконечный буфер | Лаг растёт до retention, брокер умирает по диску | Rate limit на входе, лимит лага, приоритетные топики (§9) |
| Нет метрик очереди | «Что-то медленно» вместо «лаг на партиции X» | consumer lag + DLQ-count + under-replicated (§12) |
| Batch там, где нужна свежесть | Лента с суточной задержкой | Стрим для горячего; батч — для пересчёта (§11) |
15. Что назвать на секции (чек-лист раздела)
Обязательные фразы и действия. Полный чек-лист — ../methodology.md §7; канон — ../knowledge-base.md §8.
На этапе 3–4 (схема и детали): - При первом блоке очереди на схеме — критерий: «выношу всё, что не нужно для ответа пользователю: тяжёлую работу, fan-out, счётчики, события для проекций». - Выбор брокера одной фразой: «Kafka — журнал с реплеем и независимыми читателями; RabbitMQ — точная доставка "одному из", без истории». - Словарь: «топик из партиций; порядок — внутри партиции; консьюмеры группы делят партиции, разные группы читают всё». - Семантика каждой очереди на схеме: «at-least-once + идемпотентный консьюмер; телеметрия — at-most-once». - Порядок: «ключ партиции — сущность домена (dialog_id, author_id); глобальный порядок не строю».
На этапе 6 (эксплуатация): - «Потеря закрыта по звеньям: outbox на производителе, RF=3 и min.insync=2 на брокере, коммит после обработки + DLQ на потребителе, реплей — страховка». - «Дубль закрыт идемпотентностью: операция по построению, дедуп по message_id в локальной транзакции, idempotency key для неидемпотентного». - Метрики: «consumer lag на партицию с алертом, DLQ-count, under-replicated партиции; отказ воркера — это рост очереди, а не падение сервиса». - «Асинхронность — механизм отказоустойчивости: пайплайн деградирует в лаг, а не в 5xx» (мост в 08). - Связка с соседями: outbox и гарантии публикации → 09 (outbox вместо 2PC), ретраи/backoff → 08, лаг и канарейки → 12.
Сквозные формулы раздела: - «Не нужно для ответа пользователю → из request-path в очередь». - «Транспорт дублирует, приложение дедуплицирует» (общая с 08 §5). - «Effectively-once — не свойство брокера, а правило приложения». - «Порядок — внутри партиции; причинность, а не глобальный порядок». - «Лаг — главная метрика очереди».
16. Самопроверка «умею, если»
Критерий раздела из брифа: объясняю, что выношу из request-path в очередь и что делаю при потере/дубле сообщения. Проверка — прогон двух вопросов подряд, без материалов, на материале эталонов (лента / мессенджер / медиа). Вопрос 1 «что выношу»: критерий §2 («нужен результат в ответе?» + «допустимо ли окно?») и четыре класса (тяжёлая работа, fan-out, счётчики, события) — с конкретным примером из ленты (доставка постов) и медиа (транскодинг). Вопрос 2 «потеря/дубль»: прогон по звеньям §6 (outbox → репликация → коммит-после-обработки → DLQ → реплей) и дедуп §7 (идемпотентная операция / message_id в локальной транзакции / idempotency key). Плюс два «вдогонку», которые почти наверняка последуют: «а если не успевают?» — лаг, competing consumers до числа партиций, rate limit, DLQ (§9); «а порядок?» — партиция и ключ, глобальный порядок не строим (§8). Контрольное время: весь блок — 4–6 минут этапа 4–6. Протокол тренировок и рубрика 0–3 — ../methodology.md §6; сверка с эталонами — ../classic-designs.md §2–4. Если оба вопроса брифа звучат как прогон по звеньям, а на «дважды придёт» рука тянется к дедупу, а не к «ну, настроим exactly-once» — раздел готов.
Связанные документы
00-overview.md— обзорная статья по всем 12 разделам брифа (TASK-45.40), раздел 0702-estimations.md— расчёт fan-out 700k вставок/с и «пост — событие в очереди» (§6)03-scaling.md— очередь как шаг 9 лестницы масштабирования (§2), схема с очередью и воркерами (§6)04-data-models-storage.md— outbox/CDC как событийная синхронизация проекций; идемпотентность потребителя — мост сюда06-caching.md— событийная инвалидация кэша через шину (CDC/binlog или outbox) §508-fault-tolerance.md— идемпотентность и idempotency key (§5), строка «Очередь (Kafka)» в таблице отказов (§10), асинхронность как отказоустойчивость (§12)09-consistency.md— outbox вместо 2PC, гарантии между системами, сага через шину10-components-patterns.md— rate limiter как защита входа в очередь, WebSocket-доставка событий11-reference-designs.md— эталоны, где очереди живут целиком (fan-out, delivery, транскодинг)12-operations-observability.md— метрики очереди (lag, DLQ), canary на воркерах с контролем лага../brief.md— бриф: раздел 07 с критерием «умею, если» (два вопроса: что выношу / что при потере и дубле)../knowledge-base.md— §8 очереди и потоковая обработка (словарь, семантики, паттерны, CDC), §7 outbox и идемпотентность, §9 ретраи/backoff, §10 метрики../methodology.md— §2 карточки этапов 3–6, §4 топ-10 ошибок, §6 протокол тренировок, §7 одностраничная карточка на секцию../classic-designs.md— §2 лента (fan-out, ключ author_id, lag), §3 мессенджер (dialog_id, дедуп, порядок), §4 медиа (пайплайн транскодинга, DLQ, «задержка вместо фермы»), §1 shortener (счётчики через очередь)../materials/ddia-2ed/ch-12-stream-processing.md— источник: брокеры, группы потребителей, CDC, окна, exactly-once (TASK-45.39)../materials/ddia-2ed/ch-13-philosophy-streaming.md— источник: причинность вместо total order, интеграция из единого журнала (TASK-45.39)../materials/ddia-2ed/ch-11-batch-processing.md— источник: batch, производные данные, унификация batch/stream (TASK-45.39)- Соседние статьи:
06-caching.md(событийная инвалидация через шину) ·08-fault-tolerance.md(ретраи, DLQ, таблица отказов) ·09-consistency.md(outbox/гарантии) ·11-reference-designs.md(эталоны целиком)