Транспорты микросервисов: TCP, Redis, NATS, Kafka
Микросервис Nest слушает на сменном transport, а не на HTTP. send/@MessagePattern — запрос–ответ; emit/@EventPattern — fire-and-forget. Транспорт выбирают по доставке: TCP/Redis теряют сообщения, RabbitMQ/Kafka передоставляют — что требует идемпотентных consumer'ов.
Поток «заказ оформлен» работал на демо: сервис заказов звал 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’а в эксплуатации |
| Redis | Pub/sub | Нет | At-most-once | Дешёвый fan-out, где потерять сообщение приемлемо |
| NATS | Pub/sub (subjects) | Core нет; JetStream да | At-most-once (core) | Лёгкий, низколатентный обмен; JetStream добавляет персистентность |
| RabbitMQ | Очереди + routing | Да (ack, requeue) | At-least-once | Work-очереди, 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. Биллинг теперь изредка списывает с карты дважды. Какова корневая причина и верный фикс?
- 01Объясни разницу между @MessagePattern + send() и @EventPattern + emit(), включая то, что получает вызывающий и когда использовать каждый.
- 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-уровень. Открой, попробуй, потом открой ответ.
Что-то непонятно?
Задай вопрос по этому уроку. Вопросы анонимны и попадают напрямую автору — урок станет лучше.
Примени это
Примени этот урок в реальном проекте.