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

Транспорты микросервисов: TCP, Redis, NATS, Kafka

Микросервис Nest слушает на сменном transport, а не на HTTP. send/@MessagePattern — запрос–ответ; emit/@EventPattern — fire-and-forget. Транспорт выбирают по доставке: TCP/Redis теряют сообщения, RabbitMQ/Kafka передоставляют — что требует идемпотентных consumer'ов.

NEST Senior ◷ 18 min
Уровень
ОсновыJuniorMiddleSenior
Уже знаешь этот юнит? Пройди быструю проверку за минуту →

Поток «заказ оформлен» работал на демо: сервис заказов звал emit('order.placed', order) поверх Redis, сервис биллинга списывал с карты, почтовый сервис отправлял чек. Потом под биллинга катился во время деплоя, когда пришло четырнадцать заказов. Redis pub/sub выкинул их на пол — ни один подписчик не слушал, доставлять было некуда — и четырнадцать клиентов получили списание на ноль и отгруженный товар. Фикс был не в баге кода; он был в транспорте. Redis pub/sub — это at-most-once без персистентности, а «отреагируй на это позже» требует broker, который держит сообщение, пока consumer не вернётся. Этот выбор — TCP vs Redis vs RabbitMQ vs Kafka, send vs emit — и есть весь урок.

Микросервис — это app на transport, а не на HTTP

Микросервис Nest — это тот же AppModule, что ты уже пишешь, забутстраппленный так, чтобы слушать на transport, а не на HTTP-сервере. Ты вызываешь createMicroservice с transport и options; transport сменный, и хендлеры твоих контроллеров остаются байт в байт теми же, когда ты переезжаешь с TCP на Kafka.

import { NestFactory } from '@nestjs/core';
import { Transport, MicroserviceOptions } from '@nestjs/microservices';
import { AppModule } from './app.module';

async function bootstrap() {
  const app = await NestFactory.createMicroservice<MicroserviceOptions>(AppModule, {
    transport: Transport.TCP, // меняй на REDIS / NATS / RMQ / KAFKA, хендлеры не трогаются
    options: { host: '0.0.0.0', port: 3001 },
  });
  await app.listen();
}
bootstrap();

Хендлеры бывают двух видов, и это различие — важнейшая идея урока. @MessagePattern(pattern) — это запрос–ответ: возвращаемое значение хендлера (или Observable) сериализуется обратно вызывающему, а transport сопоставляет ответ с запросом через correlation id. @EventPattern(event) — это fire-and-forget: он реагирует на то, что случилось, и не возвращает ничего — канала ответа нет вообще.

import { Controller } from '@nestjs/common';
import { MessagePattern, EventPattern, Payload } from '@nestjs/microservices';

@Controller()
export class OrdersController {
  // запрос–ответ: вызывающий ждёт это возвращаемое значение
  @MessagePattern({ cmd: 'order.total' })
  total(@Payload() ids: number[]): number {
    return ids.reduce((sum, id) => sum + priceOf(id), 0);
  }

  // fire-and-forget: «это случилось»; ответа никто не ждёт
  @EventPattern('order.placed')
  async onOrderPlaced(@Payload() order: OrderDto): Promise<void> {
    await this.billing.charge(order); // никакое значение никуда не отправляется
  }
}

На стороне вызывающего ClientProxy зеркалит это разделение: send() для @MessagePattern, emit() для @EventPattern. send() возвращает холодный Observable — ничего не отправляется, пока ты не подпишешься, и стрим эмитит ответ — так что это для «мне нужен ответ». emit() публикует событие и не ждёт его — это для «это случилось, реагируй, если тебе важно».

import { Injectable, Inject } from '@nestjs/common';
import { ClientProxy } from '@nestjs/microservices';
import { firstValueFrom } from 'rxjs';

@Injectable()
export class OrdersService {
  constructor(@Inject('ORDERS') private readonly client: ClientProxy) {}

  // send: возвращает Observable ответа; дождись его ради значения
  total(ids: number[]) {
    return firstValueFrom(this.client.send<number>({ cmd: 'order.total' }, ids));
  }

  // emit: опубликовал и пошёл дальше; ответа не ждём
  placed(order: OrderDto) {
    this.client.emit('order.placed', order); // fire-and-forget
  }
}
Почему это работает

Почему send() — это холодный Observable? Потому что запрос фактически не отправляется, пока на него кто-то не подпишется, — так что вызвать this.client.send(...) и не подписаться значит не отправить ничего. На этом спотыкаются те, кто пришёл от promise-based HTTP-клиентов: вызов выглядит так, будто сработал, но без subscribe() (или firstValueFrom/await lastValueFrom) сообщение так и не уходит. emit() устроен иначе — Nest подписывается внутри, так что событие публикуется, — но ответа всё равно не возвращает.

Выбор transport — это выбор о доставке

Каждый transport говорит на одном и том же API хендлеров, так что решение никогда не про синтаксис — оно про семантику доставки: персистит ли broker сообщение, может ли восстановившийся consumer его переиграть, гарантирован ли порядок, какую пропускную способность он держит. TCP и Redis pub/sub просты и быстры, но теряют; RabbitMQ и Kafka долговечны, но заставляют платить за это идемпотентностью. Когда видишь постмортем «события потерялись при деплое» — именно с этой таблицы начинается разбор.

TransportМодельДолговечность / replay?ДоставкаВыбирай, когда
TCPТочка-точкаНетAt-most-onceПростой внутренний запрос–ответ, без broker’а в эксплуатации
RedisPub/subНетAt-most-onceДешёвый fan-out, где потерять сообщение приемлемо
NATSPub/sub (subjects)Core нет; JetStream даAt-most-once (core)Лёгкий, низколатентный обмен; JetStream добавляет персистентность
RabbitMQОчереди + routingДа (ack, requeue)At-least-onceWork-очереди, ack на сообщение, сложный routing
KafkaПартиционированный логДа (replay по offset)At-least-onceВысоконагруженные потоки событий, replay, порядок на партицию

Два факта из этой таблицы несут сеньорский вес. Первый: at-least-once (RabbitMQ, Kafka) означает передоставку, а значит, твой consumer может отработать дважды на одном сообщении — поэтому он обязан быть идемпотентным, иначе ты спишешь с карты дважды. Второй: Kafka упорядочивает сообщения только внутри партиции, не по всему топику: события одного orderId остаются упорядоченными, только если попадают в одну партицию (ключуй по orderId), а consumer group раскидывает партиции по инстансам ради параллелизма, сохраняя порядок в пределах партиции. Запрос–ответ в Kafka тоже есть, но тяжелее — клиент должен вызвать subscribeToResponseOf(topic), чтобы ответы приходили в выделенный reply-топик и сопоставлялись назад по correlation id.

// Kafka-микросервис: список broker'ов и consumer group
const app = await NestFactory.createMicroservice<MicroserviceOptions>(AppModule, {
  transport: Transport.KAFKA,
  options: {
    client: { brokers: ['localhost:9092'] },
    consumer: { groupId: 'billing-consumer' }, // group => партиции делятся между подами
  },
});

Гибридные приложения: HTTP и transport разом

HTTP редко выбрасывают. Гибридное приложение — это обычное HTTP-приложение, которое ещё слушает как микросервис: ты делаешь connectMicroservice один или несколько transport’ов, затем startAllMicroservices() перед listen(). Один и тот же процесс обслуживает REST и реагирует на события broker’а.

const app = await NestFactory.create(AppModule);
app.connectMicroservice<MicroserviceOptions>({
  transport: Transport.KAFKA,
  options: { client: { brokers: ['localhost:9092'] }, consumer: { groupId: 'orders' } },
});
await app.startAllMicroservices(); // начинаем потреблять события
await app.listen(3000);            // и продолжаем обслуживать HTTP
Выбери лучший вариант

Событие «заказ оформлен» раздаётся в биллинг и почту. Жёсткое требование: если под consumer'а лежит во время деплоя, события должны пережить это и быть переигрываемыми по возвращении — а события одного заказа должны оставаться упорядоченными. Какой transport?

Викторина

Producer зовёт client.emit('order.placed', order) поверх Redis pub/sub, пока единственный под consumer'а посреди рестарта. Что станет с этим событием?

Викторина

Ты переносишь поток «заказ оформлен» на Kafka. Биллинг теперь изредка списывает с карты дважды. Какова корневая причина и верный фикс?

Вспомните перед уходом
  1. 01
    Объясни разницу между @MessagePattern + send() и @EventPattern + emit(), включая то, что получает вызывающий и когда использовать каждый.
  2. 02
    Тебе нужно выбрать transport для события, которое должно пережить простой consumer'а и быть переигрываемым, с порядком на заказ. Пройдись по TCP, Redis, RabbitMQ и Kafka и обоснуй выбор.
Итог

Микросервис Nest — это твой обычный AppModule, забутстраппленный через createMicroservice на сменный transport вместо HTTP-сервера; transport меняется, а хендлеры — нет. Ключевое смысловое разделение — запрос–ответ против fire-and-forget: @MessagePattern с ClientProxy.send() возвращает холодный Observable сопоставленного ответа (используй, когда нужен ответ), тогда как @EventPattern с ClientProxy.emit() публикует событие и не ждёт ничего назад (используй, чтобы объявить, что что-то случилось, и дать любому числу consumer’ов отреагировать). Выбор transport — это выбор о семантике доставки: TCP — простой точка-точка, а Redis pub/sub — дешёвый fan-out, но оба at-most-once без персистентности, так что сообщение, опубликованное без текущего слушателя, теряется; NATS core тоже at-most-once, а JetStream добавляет персистентность; RabbitMQ даёт долговечные очереди, ack на сообщение и routing при at-least-once; а Kafka — это партиционированный append-only лог с consumer group, порядком на партицию и replay по offset ради высокой пропускной способности. Сеньорские следствия идут напрямую: at-least-once (RabbitMQ, Kafka) означает передоставку, так что consumer’ы должны быть идемпотентны, иначе спишут дважды; порядок Kafka держится только внутри партиции, так что ключуй по orderId; а гибридное приложение — connectMicroservice плюс startAllMicroservices перед listen — даёт одному процессу обслуживать HTTP и реагировать на события broker’а разом. Теперь, когда встретишь инцидент «сообщения пропали при деплое», сразу знаешь — нужно смотреть на колонку «Долговечность» этой таблицы и выбирать transport по ней.

Практика

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

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

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

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

Примени это

Примени этот урок в реальном проекте.

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

Trademarks belong to their respective owners. Editorial reference only.