open atlas
↑ К треку
NestJS с нуля до senior NEST · 08 · 03

Паттерны сообщений и семантика доставки

Два паттерна сообщений (send/запрос-ответ против emit/событие) и гарантия доставки, которая реально едет в прод: at-least-once. Exactly-once доставка — миф; effectively-once получается только идемпотентным консьюмером. Дедуп по ключу, DLQ для poison-сообщений.

NEST Senior ◷ 18 min
Уровень
ОсновыJuniorMiddleSenior

В отчёте об инциденте была одна строка, от которой у дежурного инженера ёкнуло: «клиент списан дважды за один заказ, с разницей в 4 минуты». В коде приложения не было ретрая. Никто не нажимал «оплатить» дважды. Случилось более тихое: консьюмер платежей обработал событие payment.charge, вызвал Stripe, а затем — в те 30 миллисекунд между списанием с карты и коммитом подтверждения обратно брокеру — под был ротирован рутинным деплоем. Подтверждение так и не доехало. Брокер сделал ровно то, что обещал: увидел неподтверждённое сообщение и передоставил его. Новый под подхватил его, увидел совершенно валидное событие списания и списал с карты снова. Брокер не был сломан. Код не был «багнутым» ни в одной строке, на которую можно ткнуть пальцем. Вся система вела себя корректно — а клиент всё равно был списан дважды. Этот урок о том, почему такое в итоге случается с каждым, и о трёх решениях, которые это останавливают: какой паттерн сообщения ты выбираешь, какая гарантия доставки у тебя реально есть и почему идемпотентность не опциональна.

Два паттерна: вопрос против факта

Почему выбор паттерна важен ещё до того, как выбираешь broker? Потому что неверный вариант тихо связывает сервисы во времени — и именно эта связь превращает сбой одного downstream в каскадный отказ. Клиент Nest-микросервиса говорит в двух регистрах, и выбор неверного — это первая структурная ошибка. client.send(pattern, data) — это запрос-ответ: он возвращает Observable ответа, и вызывающий ждёт, пока обработчик (@MessagePattern) вычислит и вернёт ответ. client.emit(pattern, data) — это событие: это fire-and-forget — обработчик (@EventPattern) реагирует, ничего не возвращает, а отправитель не ждёт.

// REQUEST-RESPONSE: you need an answer NOW, and you block on it
const total: number = await firstValueFrom(
  this.client.send('cart.total', { cartId }),   // ↔ @MessagePattern('cart.total')
);

// EVENT-BASED: you are announcing a fact; nobody is waited on
this.client.emit('order.placed', { orderId, total }); // ↔ @EventPattern('order.placed')

Сеньорское правило — про связанность во времени. Используй send только когда тебе по-настоящему нужен ответ, чтобы продолжить — запрос, результат которого ты используешь на следующей строке. Используй emit для фактов «это произошло», на которые другие сервисы реагируют в своём темпе. Ловушка — тянуться к send везде, потому что он ощущается как вызов функции: цепочка send → send → send через сервисы тихо пересобирает синхронный монолит поверх шины сообщений. Теперь у тебя вся операционная цена распределённой системы — сетевая задержка, складывающаяся на каждом хопе, и каскадный отказ, когда любое звено лежит, — и при этом ноль развязки, которая оправдывала разделение сервисов в первую очередь. emit — это то, что позволяет сервису заказов завершиться, даже когда сервис email лежит.

Семантика доставки: три гарантии, одна из которых — ложь

Каждый брокер предлагает гарантию доставки, и их ровно три, которые надо знать.

At-most-once — это fire-and-forget без передоставки: сообщение отправлено, и если оно потерялось в пути или консьюмер умер до обработки, оно просто исчезло. Быстро, дёшево, с потерями. Приемлемо для тика метрики; неприемлемо для платежа.

At-least-once — это рабочая лошадка и дефолт для Kafka, NATS JetStream и RabbitMQ: брокер передоставляет сообщение, пока консьюмер его не подтвердит. Ничто никогда не теряется — но одно и то же сообщение может быть доставлено больше одного раза, потому что передоставка срабатывает всякий раз, когда брокер не уверен, что консьюмер закончил.

Exactly-once — это та, которую просят, и та, которую ты почти никогда реально не получаешь. Настоящая сквозная доставка exactly-once через сеть в общем случае невозможна. Что продакшен-системы реально строят — это at-least-once доставка + идемпотентная обработка, что даёт effectively-once: сообщение может прийти дважды, но второе прибытие не имеет дополнительного эффекта. Сформулируй точно, потому что неточная версия — это то, что вызывает двойное списание: ты не получаешь exactly-once доставку; ты инженеришь exactly-once эффект.

Вместе эти три категории охватывают каждый broker, с которым ты встретишься: at-most-once там, где потеря приемлема, at-least-once там, где нет, и effectively-once там, где сторона обработки поглощает дубли идемпотентностью. Без этого слоя идемпотентности at-least-once превращается в «иногда дважды» — именно так произошёл инцидент из Hook.

Почему это работает

Почему сквозная доставка exactly-once фактически невозможна, и что вместо этого делают реальные системы? Представь, что консьюмер закончил работу и шлёт ack обратно брокеру. Сеть может уронить этот ack по дороге. Теперь брокер застрял с вопросом, на который не может ответить: консьюмер обработал сообщение и ack потерялся, или консьюмер умер до обработки? Эти два случая неразличимы со стороны брокера — единственный сигнал у него «ack не пришёл». Он должен выбрать дефолтное поведение для этой неоднозначности. Если он предположит «обработано» и пойдёт дальше, он рискует потерять сообщение, которое на деле не было обработано. Если он предположит «не обработано» и передоставит, он рискует продублировать сообщение, которое было. Нет третьего варианта, который магически знает правду, потому что информация, нужная чтобы различить случаи, потерялась вместе с ack. Поэтому зрелые системы перестают пытаться сделать доставку exactly-once и вместо этого выбирают at-least-once (никогда не терять) и делают обработку идемпотентной (дедуп по стабильному ключу), чтобы дубль доставки не давал дубля эффекта — effectively-once, а не exactly-once-доставка.

Подтверждения и окно передоставки

Дубль рождается в зазоре между «работа сделана» и «ack записан». Консьюмер делает работу, затем подтверждает: Kafka коммитит консьюмерский offset; NATS JetStream и RabbitMQ шлют явный ack. Если консьюмер падает или ack теряется до того, как брокер записал это подтверждение, брокер передоставляет — и сообщение обрабатывается второй раз.

Это значит, что порядок операций и есть выбор того, какой отказ ты принимаешь, и бесплатного варианта нет:

// ACK BEFORE processing → at-most-once in effect: a crash here LOSES the message
@EventPattern('payment.charge')
async handleBad(@Payload() data, @Ctx() ctx: KafkaContext) {
  await ctx.getConsumer().commitOffsets([/* ... */]); // committed first
  await this.charge(data); // if the pod dies here, the charge never happens, never retried
}

// ACK AFTER processing → at-least-once: a crash here DUPLICATES the message on redelivery
@EventPattern('payment.charge')
async handleGood(@Payload() data, @Ctx() ctx: KafkaContext) {
  await this.charge(data);                              // work first
  await ctx.getConsumer().commitOffsets([/* ... */]);  // crash before this → redelivered → charged twice
}

Подтверждение до обработки теряет сообщения при падении; подтверждение после обработки дублирует их при падении. Нельзя получить ни того, ни другого. Зрелые системы выбирают ack-after (никогда не терять) и платят за это идемпотентностью (никогда не применять дважды). В чём и весь смысл следующего раздела.

Идемпотентность: фикс, а не приятная мелочь

Поскольку at-least-once будет доставлять дубли, консьюмер обязан быть идемпотентным: обработка одного и того же сообщения дважды должна давать то же конечное состояние, что и обработка один раз. Это достигается не надеждой; это инженерится из трёх движущихся частей:

  1. Стабильный ключ идемпотентности, который идентифицирует логическую операцию — несомый на сообщении, одинаковый между передоставками (например, paymentId или ключ от клиента). Не случайный id на попытку, иначе каждая передоставка выглядит новой.
  2. Дедуп-запись уже обработанных ключей — таблица processed_events, проверяемая первой и записываемая в той же транзакции, что и эффект.
  3. Побочный эффект, безопасный для повтора — защищённый так, чтобы второе прибытие было no-op.
@EventPattern('payment.charge')
async handleCharge(@Payload() data: { paymentId: string; amount: number }) {
  await this.dataSource.transaction(async (manager) => {
    // claim the key; if it already exists, this insert affects 0 rows → already processed
    const res = await manager.query(
      `INSERT INTO processed_events (key) VALUES ($1) ON CONFLICT (key) DO NOTHING`,
      [data.paymentId],
    );
    if (res.rowCount === 0) return; // duplicate redelivery → skip the charge entirely
    await this.charge(data); // first time only; commits atomically with the dedup row
  });
}

Этот ON CONFLICT (key) DO NOTHING — несущая строка: дедуп-строка и списание коммитятся в одной транзакции, так что передоставка либо находит ключ уже на месте (и пропускает), либо застолбляет его и списывает — но не оба сразу. Именно этого недоставало в инциденте с двойным списанием. Фикс, уехавший в тот день, был ровно этим: уникальный ключ идемпотентности на списании плюс проверка processed-events до вызова Stripe.

Poison-сообщения и dead-letter queue

У at-least-once есть второе, более острое лезвие. Допустим, сообщение никогда не может пройти — кривой payload, рассинхрон схемы, баг, который бросает на этой конкретной записи. Передоставка брокера, которая спасает тебя от временных падений, теперь работает против тебя: она передоставляет poison-сообщение вечно. Хуже того, в упорядоченном потоке poison сидит в голове и блокирует каждое сообщение за ним — head-of-line blocking — сжигая CPU в бесконечном цикле ретрая.

Ответ — это dead-letter queue (DLQ): после N неудачных попыток (обычно 3–5 ретраев) перестань ретраить inline и направь сообщение в отдельную очередь для разбора и ручной или автоматической обработки, чтобы основной поток продолжал течь. Счётчик ретраев — это ручка: слишком мало — и ты отправляешь в DLQ временные сбои, которые бы прошли; слишком много — и poison-сообщение стопорит партицию на минуты.

Ещё одна гарантия, с которой надо быть честным: порядок — только в пределах партиции. Kafka гарантирует порядок внутри партиции (выбираемой ключом сообщения), JetStream — внутри стрима субъекта, но между партициями/ключами глобального порядка нет. Если два события обязаны примениться по порядку, они должны разделять ключ партиции; иначе «мы обработаем их по порядку» — это гарантия, которой у тебя на деле нет.

Выбери лучший вариант

Сообщение «списать с клиента» не должно списывать дважды, хотя брокер может передоставить его после рестарта консьюмера. Какой подход реально держит?

Викторина

Когда стоит использовать client.send() против client.emit() между Nest-микросервисами?

Викторина

Событие payment.charge передоставлено после рестарта пода консьюмера, и клиент списан дважды. Какова корневая причина и фикс?

Вспомните перед уходом
  1. 01
    Сопоставь запрос-ответ (send) с событийным (emit) обменом сообщениями и объясни сеньорское правило выбора — включая failure mode злоупотребления send.
  2. 02
    Объясни три семантики доставки, почему exactly-once доставка фактически невозможна и как at-least-once + идемпотентность + DLQ вместе дают безопасный консьюмер.
Итог

Nest-микросервисы говорят в двух паттернах: client.send(pattern, data) — это запрос-ответ — он возвращает Observable ответа, и вызывающий ждёт обработчик @MessagePattern, связывая сервисы во времени — тогда как client.emit(pattern, data) — это fire-and-forget, обрабатываемый @EventPattern без ответа и без ожидания. Используй send только когда нужен ответ сейчас; используй emit для фактов «это произошло», потому что цепочка вызовов send пересобирает синхронный монолит со складывающейся задержкой и каскадным отказом. По доставке: at-most-once с потерями, at-least-once (дефолт для Kafka, NATS JetStream, RabbitMQ) никогда не теряет, но МОЖЕТ дублировать, а exactly-once доставка — миф: когда ack потерян, брокер не может отличить «обработано» от «не обработано», поэтому обязан рисковать потерей или дублем. Дубль рождается в окне между выполнением работы и записью ack/offset: ack-до-обработки теряет при падении, ack-после дублирует при падении, поэтому выбираешь ack-после и платишь идемпотентностью. Идемпотентный консьюмер несёт стабильный ключ идемпотентности, столбит его в дедуп-таблице processed-events через INSERT … ON CONFLICT DO NOTHING в той же транзакции, что и эффект, и делает побочный эффект безопасным для повтора — так что передоставка находит ключ и пропускает, давая effectively-once. Именно этого недоставало в инциденте с двойным списанием. Poison-сообщение, которое всегда падает, передоставлялось бы вечно и head-of-line-блокировало бы свою партицию, поэтому после ~3–5 ретраев направляй его в dead-letter queue (DLQ — отдельная очередь для неудачных сообщений); а порядок держится только в пределах ключа партиции, так что события, которые должны быть упорядочены, обязаны разделять ключ. Теперь, когда встретишь отчёт о двойном списании, сразу знаешь, что искать: нет ключа идемпотентности и есть зазор ack-after — и что именно шипить, чтобы закрыть его.

Практика

Начни сверху. Задачи идут от простого к сложному: вспомнить факт, применить к случаю, затем senior-уровень. Открой, попробуй, потом открой ответ.

вспомнитьприменитьуглубить0 из 5 завершено

Что-то непонятно?

Задай вопрос по этому уроку. Вопросы анонимны и попадают напрямую автору — урок станет лучше.

хоткеи развернуть
поиск
K
пред. пьеса
k
след. пьеса
j
тиры
t
это меню
?
sources2
expand
  1. 01
  2. 02

Trademarks belong to their respective owners. Editorial reference only.