Пошук уроків, статей та іншого контенту
Порівняєте TCP, Redis, NATS, RabbitMQ та Kafka і виберете транспорт відповідно до вимог системи.
У мікросервісній архітектурі транспортний рівень визначає, як сервіси обмінюються повідомленнями:
як встановлюється з’єднання;
чи зберігаються повідомлення;
що відбувається під час недоступності споживача;
чи підтримуються підтвердження доставки;
як масштабується обробка;
чи зберігається порядок повідомлень;
як реалізуються повторні спроби;
чи можна відтворити старі повідомлення.
У NestJS транспорт обирається через Transport і налаштовується під час створення мікросервісу або клієнта:
import { Transport } from '@nestjs/microservices';
{
transport: Transport.TCP,
options: {
host: '127.0.0.1',
port: 8877,
},
}Вбудовані транспортні стратегії NestJS:
TCP;
REDIS;
NATS;
RMQ для RabbitMQ;
KAFKA.
Важливо розрізняти транспорт NestJS і брокер повідомлень. NestJS надає адаптер, який перетворює виклики @MessagePattern() та ClientProxy на протокол конкретної системи.
Перед вибором технології потрібно визначити тип взаємодії.
Клієнт надсилає повідомлення і очікує результат:
Orders Service → Payments Service
← результат оплатиТакий підхід зручний для:
отримання даних;
перевірки доступності;
коротких операцій;
команд, без результату яких клієнт не може продовжити роботу.
Недолік — сильніший зв’язок між сервісами. Якщо сервіс-отримувач недоступний або працює повільно, клієнт може чекати або завершитися помилкою.
Сервіс публікує факт, а інші сервіси реагують на нього:
Orders Service → OrderCreated
├─ Billing Service
├─ Notification Service
└─ Analytics ServiceТакий підхід зменшує зв’язність і дає змогу обробляти повідомлення пізніше. Натомість потрібно продумати:
повторну доставку;
дублікати;
ідемпотентність;
порядок подій;
зберігання повідомлень;
моніторинг невдалих обробок.
Один і той самий брокер може підтримувати обидва сценарії, але з різними гарантіями та обмеженнями.
TCP-транспорт NestJS використовує пряме мережеве з’єднання між клієнтом і мікросервісом. За замовчуванням це найпростіший транспорт для внутрішнього запит-відповідь-обміну.
// src/main.ts
import { NestFactory } from '@nestjs/core';
import { Transport } from '@nestjs/microservices';
import { AppModule } from './app.module';
async function bootstrap() {
const app = await NestFactory.createMicroservice(AppModule, {
transport: Transport.TCP,
options: {
host: process.env.TCP_HOST ?? '0.0.0.0',
port: Number(process.env.TCP_PORT ?? 8877),
},
});
await app.listen();
}
bootstrap();// src/app.controller.ts
import { Controller } from '@nestjs/common';
import { MessagePattern } from '@nestjs/microservices';
@Controller()
export class AppController {
@MessagePattern({ cmd: 'health' })
health() {
return {
status: 'ok',
service: 'catalog',
timestamp: new Date().toISOString(),
};
}
}// src/app.module.ts
import { Module } from '@nestjs/common';
import { AppController } from './app.controller';
@Module({
controllers: [AppController],
})
export class AppModule {}Клієнт такого мікросервісу може бути налаштований так:
import { ClientProxyFactory, Transport } from '@nestjs/microservices';
const client = ClientProxyFactory.create({
transport: Transport.TCP,
options: {
host: 'catalog-service',
port: 8877,
},
});
const result = await client
.send(
{ cmd: 'health' },
{},
)
.toPromise();
console.log(result);У сучасному коді замість застарілого toPromise() зазвичай використовують firstValueFrom з RxJS:
import { firstValueFrom } from 'rxjs';
const result = await firstValueFrom(
client.send({ cmd: 'health' }, {}),
);низька затримка;
проста конфігурація;
немає окремого брокера;
природний запит-відповідь;
зручно використовувати для внутрішніх RPC-викликів;
легко почати локальну розробку.
TCP-транспорт NestJS не є чергою повідомлень:
повідомлення не зберігаються в брокері;
недоступний сервіс не отримує повідомлення пізніше;
немає вбудованого журналу подій;
повторна доставка потребує логіки клієнта;
масштабування потрібно організовувати на рівні мережі або сервісної інфраструктури.
TCP добре підходить для:
синхронних команд;
внутрішніх запитів між сервісами;
операцій, де відповідь потрібна негайно;
простих систем із невеликою кількістю сервісів.
TCP не є хорошим вибором для подій, які не можна втратити.
Транспорт Redis у NestJS зазвичай використовує Redis Pub/Sub. Повідомлення публікується в канал, а підписані екземпляри сервісу його отримують.
import { Transport } from '@nestjs/microservices';
const redisOptions = {
transport: Transport.REDIS,
options: {
host: process.env.REDIS_HOST ?? 'localhost',
port: Number(process.env.REDIS_PORT ?? 6379),
},
};Redis Pub/Sub працює за принципом повідомлення без довготривалого журналу:
якщо споживач підписаний на канал, він отримує повідомлення;
якщо споживач був недоступний у момент публікації, повідомлення втрачається;
повідомлення не очікують споживача в черзі;
Redis не підтверджує обробку повідомлення так, як це робить RabbitMQ.
Тому Redis Pub/Sub часто використовують для:
швидких подій без вимоги гарантованого зберігання;
інвалідації кешу;
повідомлень про оновлення стану;
синхронізації кількох екземплярів застосунку;
простих внутрішніх broadcast-подій.
дуже проста конфігурація;
низька затримка;
Redis часто вже використовується як кеш;
зручно розсилати повідомлення багатьом підписникам;
не потрібно керувати окремою складною брокерною інфраструктурою.
Pub/Sub не зберігає повідомлення;
немає надійної повторної доставки;
складніше реалізовувати контроль невдалих обробок;
Redis Pub/Sub не замінює повноцінну чергу;
порядок і гарантії не повинні трактуватися як довготривала історія подій.
Redis Streams — інша можливість Redis, яка підтримує зберігання записів і групи споживачів. Однак це не те саме, що Redis Pub/Sub, і автоматичний вибір Streams замість Pub/Sub транспортом NestJS не слід припускати.
NATS — легкий брокер повідомлень, орієнтований на швидкий обмін через subjects.
Налаштування NATS-транспорту в NestJS:
import { Transport } from '@nestjs/microservices';
const natsOptions = {
transport: Transport.NATS,
options: {
servers: [
process.env.NATS_URL ?? 'nats://localhost:4222',
],
},
};Повідомлення публікуються в subject, наприклад:
orders.created
billing.payment_requested
users.profile_updatedУ NestJS шаблон повідомлення відображається на механізм subject-ів, а обробник оголошується через @MessagePattern().
Core NATS підходить для швидких повідомлень, коли втрата повідомлення під час недоступності споживача допустима або обробка може бути повторно ініційована джерелом.
Його властивості:
дуже низька затримка;
проста модель publish/subscribe;
підтримка queue groups для розподілу роботи між споживачами;
повідомлення не утворюють довготривалий журнал;
базовий транспорт не слід розглядати як сховище подій.
Queue group означає, що з кількох екземплярів одного логічного споживача повідомлення отримає лише один екземпляр. Це зручно для горизонтального масштабування обробників.
JetStream додає до NATS функції, потрібні для надійнішої доставки:
зберігання повідомлень;
підтвердження обробки;
повторну доставку;
відтворення повідомлень;
споживачів із певним станом.
Вибираючи NATS, потрібно окремо визначити, чи достатньо Core NATS, чи потрібен JetStream. Не можна переносити гарантії JetStream на звичайний Core NATS.
NATS доречний, коли потрібні:
низька затримка;
проста модель subjects;
масштабовані групи споживачів;
легкий брокер для внутрішньої комунікації;
за потреби — перехід до персистентної моделі JetStream.
RabbitMQ — брокер повідомлень із розвиненою моделлю маршрутизації та черг.
У NestJS він налаштовується через транспорт RMQ:
import { Transport } from '@nestjs/microservices';
const rabbitOptions = {
transport: Transport.RMQ,
options: {
urls: [
process.env.RABBITMQ_URL ?? 'amqp://localhost:5672',
],
queue: 'orders_queue',
queueOptions: {
durable: true,
},
},
};У цьому прикладі NestJS підключається до RabbitMQ і використовує чергу orders_queue.
Producer — публікує повідомлення.
Exchange — приймає повідомлення і маршрутизує їх.
Binding — правило зв’язку exchange із чергою.
Queue — зберігає повідомлення до обробки споживачем.
Consumer — отримує та обробляє повідомлення.
Acknowledgement — підтвердження успішної обробки.
RabbitMQ особливо сильний у маршрутизації. Exchange може направити подію:
в одну конкретну чергу;
у кілька черг;
за routing key;
за шаблоном topic;
відповідно до типу exchange.
явні черги та підтвердження;
добре розвинена маршрутизація;
контроль кількості одночасних повідомлень через prefetch;
зручна модель робочих команд;
підтримка тимчасової недоступності споживачів;
природна інтеграція з retry-чергами та dead-letter-маршрутизацією.
складніша експлуатація, ніж у TCP або Redis;
повідомлення зазвичай видаляються після підтвердження, а не зберігаються як довічний журнал;
повторна обробка потребує окремої політики;
глобальний порядок повідомлень стає складним при кількох споживачах;
масштабування вимагає розуміння черг, exchange і розподілу навантаження.
RabbitMQ зазвичай вибирають для команд і робочих черг:
GenerateInvoice → invoice_queue
SendEmail → email_queue
ResizeImage → image_queueЯкщо обробник тимчасово недоступний, повідомлення може залишатися в черзі. Це принципово відрізняє RabbitMQ від Redis Pub/Sub.
Kafka — розподілений журнал подій, орієнтований на високий потік повідомлень, масштабування та повторне читання історії.
Налаштування Kafka-транспорту в NestJS:
import { Transport } from '@nestjs/microservices';
const kafkaOptions = {
transport: Transport.KAFKA,
options: {
client: {
clientId: 'orders-service',
brokers: [
process.env.KAFKA_BROKER ?? 'localhost:9092',
],
},
consumer: {
groupId: 'orders-service-consumer',
},
subscribe: {
fromBeginning: false,
},
},
};Topic — логічний потік подій.
Partition — упорядкована частина topic.
Offset — позиція повідомлення в partition.
Producer — записує повідомлення.
Consumer group — група споживачів, яка спільно обробляє partition.
Retention — політика зберігання повідомлень.
Kafka не видаляє повідомлення одразу після обробки конкретним споживачем. Повідомлення зберігаються відповідно до політики retention, а споживач зберігає свій offset.
У межах однієї consumer group:
кожна partition обробляється одним активним споживачем;
кількість паралельних споживачів ефективно обмежена кількістю partition;
різні consumer group можуть незалежно читати той самий topic;
нова group може прочитати історичні повідомлення, якщо вони ще зберігаються.
Наприклад:
orders.events
├─ partition 0
├─ partition 1
└─ partition 2Якщо topic має три partition, одна consumer group може обробляти його максимум трьома паралельними активними споживачами.
Kafka гарантує порядок у межах однієї partition, але не гарантує глобальний порядок усього topic.
Для подій одного агрегату зазвичай використовують стабільний ключ:
key = orderIdТоді події конкретного замовлення потрапляють в одну partition і зберігають порядок:
OrderCreated
PaymentConfirmed
OrderShippedНе слід покладатися на глобальний порядок подій, якщо вони розподілені між різними partition.
високий throughput;
довготривале зберігання подій;
можливість повторного читання;
незалежні consumer groups;
масштабування через partitions;
зручна основа для event-driven та аналітичних систем.
складніша експлуатація;
потрібно проєктувати topics, partitions і keys;
неправильний ключ може створити нерівномірне навантаження;
споживачі повинні правильно працювати з offset;
Kafka не є простою заміною RabbitMQ для кожної команди;
малі системи можуть отримати надмірну операційну складність.
Kafka доречна для:
потоку доменних подій;
інтеграції багатьох незалежних споживачів;
високого обсягу повідомлень;
аудиту та повторного відтворення;
потокової аналітики;
побудови read-моделей з історії подій.
Обирайте TCP, якщо:
потрібен простий RPC;
відповідь потрібна одразу;
повідомлення не потрібно зберігати;
брокер додав би зайву складність;
сервіси доступні одночасно.
Не обирайте TCP, якщо команда повинна чекати в черзі або подія не може бути втрачена.
Обирайте Redis Pub/Sub, якщо:
повідомлення короткоживучі;
важлива низька затримка;
втрата події допустима;
Redis уже є частиною інфраструктури;
потрібна проста розсилка оновлень.
Не обирайте Redis Pub/Sub для платежів, замовлень або інших критичних команд без додаткового механізму надійної доставки.
Обирайте NATS, якщо:
важливі низька затримка і простота subjects;
потрібні queue groups;
потрібна легка внутрішня шина повідомлень;
ви чітко розрізняєте Core NATS і JetStream;
персистентність потрібна лише для частини сценаріїв.
Обирайте RabbitMQ, якщо:
потрібні робочі черги;
потрібні підтвердження обробки;
важлива гнучка маршрутизація;
необхідно регулювати швидкість споживачів;
потрібні retry- та dead-letter-сценарії.
Обирайте Kafka, якщо:
події потрібно зберігати та перечитувати;
є багато незалежних споживачів;
потрібен високий throughput;
потрібна масштабованість через partition;
команда готова керувати topics, offsets і consumer groups.
Жоден транспорт не робить бізнес-операцію автоматично надійною. Потрібно розрізняти транспортну доставку та успішне виконання бізнес-логіки.
Повідомлення обробляється не більше одного разу, але може бути втрачено.
Це типовий сценарій для ефемерних Pub/Sub-повідомлень. Він прийнятний для:
оновлення кешу;
некритичних метрик;
сигналів про зміну стану, якщо стан можна отримати повторно.
Повідомлення буде доставлено щонайменше один раз, але можуть виникнути дублікати.
Це звичайна практична модель для черг і брокерів із підтвердженнями. Обробник повинен бути ідемпотентним.
Наприклад, перед обробкою платежу можна перевірити унікальний eventId:
type PaymentRequested = {
eventId: string;
orderId: string;
amount: number;
};
async function handlePayment(event: PaymentRequested) {
const alreadyProcessed = await paymentsRepository.hasEvent(
event.eventId,
);
if (alreadyProcessed) {
return;
}
await paymentsRepository.process({
orderId: event.orderId,
amount: event.amount,
});
await paymentsRepository.markEventAsProcessed(event.eventId);
}У реальній системі перевірку та запис потрібно захистити транзакцією або унікальним обмеженням у базі даних. Інакше два паралельні обробники можуть одночасно пройти перевірку.
Гарантія «рівно один раз» у розподіленій системі складна. Навіть якщо брокер гарантує певні властивості запису або читання, бізнес-операція може бути виконана повторно через збій між:
записом у базу даних;
підтвердженням повідомлення;
фіксацією offset.
Тому замість сліпої віри в exactly-once зазвичай застосовують:
ідемпотентні обробники;
унікальні ідентифікатори подій;
транзакції;
outbox-патерн;
контроль повторних спроб;
спостережуваність і ручне відновлення.
Запитайте:
клієнту потрібна відповідь зараз?
це команда чи подія?
чи може джерело продовжити роботу без відповіді?
чи повинні кілька сервісів незалежно отримати подію?
Для простого внутрішнього RPC почніть із TCP. Для асинхронної обробки розглядайте брокер.
Якщо втрата допустима:
TCP може бути достатнім для запиту;
Redis Pub/Sub або Core NATS можуть бути достатніми для події.
Якщо втрата неприпустима:
розглядайте RabbitMQ;
Kafka;
NATS JetStream;
інший персистентний механізм із підтвердженнями.
Одна робоча черга, де завдання обробляє один із доступних воркерів:
RabbitMQ;
NATS queue group;
Kafka consumer group.
Багато незалежних читачів однієї історії подій:
Kafka;
NATS JetStream за відповідної моделі споживання.
Broadcast без зберігання:
Redis Pub/Sub;
Core NATS.
Якщо сервіс, який з’явився пізніше, повинен прочитати старі події, потрібен журнал або персистентний stream:
Kafka;
NATS JetStream;
окремо спроєктований Redis Streams-сценарій.
TCP і звичайний Redis Pub/Sub для цього не підходять.
Для TCP порядок визначається конкретним з’єднанням і сценарієм взаємодії.
У RabbitMQ порядок може порушуватися через кілька споживачів, повторну доставку та паралельну обробку.
У Kafka порядок гарантується в межах partition.
У NATS не слід будувати глобальну бізнес-логіку на припущенні про загальний порядок усіх повідомлень.
У Redis Pub/Sub не слід використовувати транспорт як журнал упорядкованих подій.
Оцініть:
чи вже працює потрібна інфраструктура;
чи потрібна кластеризація;
хто буде моніторити брокер;
як оброблятимуться відмови;
як видалятимуться або зберігатимуться повідомлення;
як відбуватиметься повторна обробка;
як тестуватиметься відновлення після збою.
Технологічно потужніший транспорт не завжди є кращим. Він має відповідати вимогам системи, а не збільшувати кількість компонентів без потреби.
API Gateway → Catalog Service
← список товарівРекомендований початковий вибір — TCP. Якщо в організації вже є NATS або інший брокер і потрібні відповідні гарантії, можна використати його request-відповідь-модель.
API → команда GenerateReport → workerРекомендований вибір — RabbitMQ або персистентний NATS JetStream. Команда повинна залишатися доступною, якщо воркер тимчасово перезапущений.
Orders → OrderCreated
├─ Billing
├─ Notifications
└─ AnalyticsЯкщо потрібні зберігання, replay і багато незалежних consumer groups, Kafka зазвичай підходить краще. Для меншої системи з простішою експлуатацією може бути достатньо RabbitMQ або NATS JetStream.
Catalog → ProductUpdated → екземпляри APIRedis Pub/Sub або Core NATS можуть бути достатніми, якщо отримувач у разі пропуску здатний відновити актуальний стан із бази даних.
Не змішуйте налаштування різних брокерів в один неявний конфігураційний об’єкт. Для кожного транспорту є власні параметри.
import { Transport } from '@nestjs/microservices';
export const tcpTransport = {
transport: Transport.TCP,
options: {
host: process.env.CATALOG_HOST ?? 'catalog-service',
port: Number(process.env.CATALOG_PORT ?? 8877),
},
};
export const redisTransport = {
transport: Transport.REDIS,
options: {
host: process.env.REDIS_HOST ?? 'localhost',
port: Number(process.env.REDIS_PORT ?? 6379),
},
};
export const natsTransport = {
transport: Transport.NATS,
options: {
servers: [process.env.NATS_URL ?? 'nats://localhost:4222'],
},
};
export const rabbitTransport = {
transport: Transport.RMQ,
options: {
urls: [process.env.RABBITMQ_URL ?? 'amqp://localhost:5672'],
queue: 'catalog_queue',
queueOptions: {
durable: true,
},
},
};
export const kafkaTransport = {
transport: Transport.KAFKA,
options: {
client: {
clientId: 'catalog-service',
brokers: [
process.env.KAFKA_BROKER ?? 'localhost:9092',
],
},
consumer: {
groupId: 'catalog-service-consumer',
},
},
};У production-системі важливо також:
не зберігати облікові дані в коді;
перевіряти обов’язкові змінні середовища під час запуску;
налаштовувати тайм-аути для запит-відповідь-сценаріїв;
логувати ідентифікатор повідомлення та кореляції;
контролювати кількість повторних спроб;
вимірювати lag, розмір черг і час обробки;
перевіряти поведінку після перезапуску споживача.
Якщо сервіс не був підключений у момент публікації, він не отримає повідомлення. Для платежів, замовлень і задач, які не можна втратити, потрібна персистентна черга або журнал.
Kafka потребує проєктування partition, ключів, consumer groups, retention і обробки offset. Для невеликої черги фонових задач RabbitMQ може бути простішим і доречнішим.
RabbitMQ — передусім брокер черг. Після підтвердження повідомлення воно зазвичай більше не доступне конкретній черзі. Якщо потрібен replay для нових споживачів, потрібна інша модель зберігання.
Порядок гарантується в межах partition. Якщо події одного агрегату розподілені між різними partition, глобального порядку між ними немає.
Повторна доставка — нормальний сценарій. Обробник, який двічі списує кошти або створює два однакові записи, є небезпечним незалежно від вибраного брокера.
У групі споживачів повідомлення розподіляється між екземплярами. У broadcast-сценарії повідомлення мають отримати всі підписники. Потрібно чітко визначити, яка поведінка потрібна бізнес-операції.
Потрібно заздалегідь відповісти:
що відбудеться, якщо брокер недоступний;
де залишиться повідомлення;
скільки разів виконуватиметься retry;
як уникнути дубліката;
як знайти повідомлення, яке постійно завершується помилкою;
як відновити обробку після простою.
TCP — простий і швидкий транспорт для внутрішнього запит-відповідь-обміну без брокера.
Redis Pub/Sub — швидкий, але ефемерний канал; пропущені повідомлення не очікують споживача.
NATS — легка високопродуктивна шина subjects; Core NATS і JetStream мають різні гарантії.
RabbitMQ — брокер черг із підтвердженнями, маршрутизацією та зручними worker-сценаріями.
Kafka — розподілений журнал подій із partitions, consumer groups, retention і можливістю replay.
Вибір транспорту визначається не популярністю технології, а вимогами до доставки, зберігання, порядку, масштабування та відновлення.
Незалежно від транспорту, критичні обробники повинні бути ідемпотентними, а політики retry та відмов мають бути явними.