Пошук уроків, статей та іншого контенту
Підключите Kafka до NestJS і розберете топіки, партиції та потокову доставку подій.
Kafka — це розподілена платформа для передавання потоків подій. У Kafka продюсери публікують повідомлення, а консьюмери читають їх із топіків.
Типовий потік виглядає так:
NestJS producer → topic → NestJS consumerKafka не передає подію безпосередньо конкретному сервісу. Подія записується в топік і зберігається там певний час. Консьюмер може прочитати її одразу або пізніше, починаючи з потрібної позиції.
У NestJS Kafka використовується через транспорт мікросервісів:
ClientKafka — для публікації подій;
@EventPattern() — для обробки подій;
KafkaContext — для доступу до метаданих повідомлення;
Топік — це іменований потік подій. Наприклад:
orders.created
payments.completed
notifications.requestedПродюсер публікує повідомлення в топік, а консьюмер підписується на нього.
Топік не є чергою в класичному розумінні. Подія залишається доступною після її прочитання, доки Kafka не видалить її відповідно до налаштувань зберігання.
Кожен топік складається з однієї або кількох партицій.
orders.created
├── partition 0
├── partition 1
└── partition 2Партиція — це впорядкований журнал подій. Порядок гарантується лише всередині однієї партиції, а не для всього топіку.
Кілька партицій дозволяють обробляти події паралельно.
Кожне повідомлення в партиції має числову позицію — offset.
partition 0:
offset 0 → подія A
offset 1 → подія B
offset 2 → подія CКонсьюмер зберігає останній оброблений offset. Якщо сервіс перезапуститься, він може продовжити читання з цієї позиції.
Consumer group — це група консьюмерів, які спільно обробляють один топік.
Наприклад, якщо топік має три партиції та група має три екземпляри сервісу, Kafka може розподілити навантаження так:
consumer-1 → partition 0
consumer-2 → partition 1
consumer-3 → partition 2У межах однієї consumer group кожна партиція обробляється лише одним активним консьюмером.
Якщо дві різні групи читають один топік, кожна група отримує всі події незалежно від іншої:
orders-service-group → усі події
analytics-service-group → усі подіїДля локального запуску можна використати Docker Compose:
services:
kafka:
image: apache/kafka:3.8.1
container_name: kafka
ports:
- "9092:9092"
environment:
KAFKA_NODE_ID: 1
KAFKA_PROCESS_ROLES: broker,controller
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT
KAFKA_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka:9093
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1
KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0
KAFKA_NUM_PARTITIONS: 3Запустіть Kafka:
docker compose up -dСтворіть топік із трьома партиціями:
docker exec kafka /opt/kafka/bin/kafka-topics.sh \
--create \
--topic orders.created \
--partitions 3 \
--replication-factor 1 \
--bootstrap-server localhost:9092У production кількість партицій і фактор реплікації потрібно планувати окремо. Для локального прикладу достатньо одного брокера.
Встановіть залежності:
npm install @nestjs/microservices kafkajsУ NestJS Kafka-клієнт реєструється через ClientsModule.
// src/app.module.ts
import { Module } from '@nestjs/common';
import { ClientsModule, Transport } from '@nestjs/microservices';
import { OrdersController } from './orders.controller';
@Module({
imports: [
ClientsModule.register([
{
name: 'KAFKA_SERVICE',
transport: Transport.KAFKA,
options: {
client: {
clientId: 'orders-api-producer',
brokers: ['localhost:9092'],
},
consumer: {
groupId: 'orders-api-client',
},
producer: {
allowAutoTopicCreation: false,
},
},
},
]),
],
controllers: [OrdersController],
})
export class AppModule {}Основні параметри:
clientId — ідентифікатор Kafka-клієнта;
brokers — адреси Kafka-брокерів;
groupId — ідентифікатор consumer group;
allowAutoTopicCreation — чи дозволено автоматично створювати топіки.
У прикладі автоматичне створення вимкнене, тому топік потрібно створити заздалегідь.
Kafka-споживач підключається як мікросервіс у main.ts:
// src/main.ts
import { NestFactory } from '@nestjs/core';
import { MicroserviceOptions, Transport } from '@nestjs/microservices';
import { AppModule } from './app.module';
async function bootstrap() {
const app = await NestFactory.create(AppModule);
app.connectMicroservice<MicroserviceOptions>({
transport: Transport.KAFKA,
options: {
client: {
clientId: 'orders-api-consumer',
brokers: ['localhost:9092'],
},
consumer: {
groupId: 'orders-service-group',
},
},
});
await app.startAllMicroservices();
await app.listen(3000);
}
bootstrap();Тут використовуються дві Kafka-конфігурації:
ClientsModule — для публікації подій;
connectMicroservice() — для отримання та обробки подій.
У цього прикладу producer і consumer мають різні clientId та groupId. Це дозволяє використовувати їх в одному NestJS-застосунку без змішування ролей.
Створимо контролер, який:
приймає HTTP-запит;
формує подію;
публікує її в топік orders.created.
// src/orders.controller.ts
import {
Body,
Controller,
Inject,
OnModuleInit,
Post,
} from '@nestjs/common';
import {
ClientKafka,
Ctx,
EventPattern,
KafkaContext,
Payload,
} from '@nestjs/microservices';
import { firstValueFrom } from 'rxjs';
import { randomUUID } from 'node:crypto';
interface CreateOrderDto {
amount: number;
}
interface OrderCreatedEvent {
orderId: string;
amount: number;
createdAt: string;
}
@Controller('orders')
export class OrdersController implements OnModuleInit {
constructor(
@Inject('KAFKA_SERVICE')
private readonly kafkaClient: ClientKafka,
) {}
async onModuleInit() {
await this.kafkaClient.connect();
}
@Post()
async createOrder(@Body() body: CreateOrderDto) {
const event: OrderCreatedEvent = {
orderId: randomUUID(),
amount: body.amount,
createdAt: new Date().toISOString(),
};
await firstValueFrom(
this.kafkaClient.emit('orders.created', event),
);
return {
status: 'published',
event,
};
}
@EventPattern('orders.created')
handleOrderCreated(
@Payload() event: OrderCreatedEvent,
@Ctx() context: KafkaContext,
) {
const message = context.getMessage();
console.log('Отримано подію:', event);
console.log('Топік:', context.getTopic());
console.log('Партиція:', context.getPartition());
console.log('Offset:', message.offset);
}
}Запустіть застосунок:
npm run start:devОпублікуйте подію HTTP-запитом:
curl -X POST http://localhost:3000/orders \
-H "Content-Type: application/json" \
-d '{"amount": 1500}'Контролер публікує подію в топік, а метод handleOrderCreated() отримує цю подію через Kafka-мікросервіс.
firstValueFrom()Метод ClientKafka.emit() повертає RxJS Observable. Операція публікації виконується після підписки на цей Observable.
firstValueFrom() перетворює його на Promise, тому подію можна коректно дочекатися в async-методі:
await firstValueFrom(
this.kafkaClient.emit('orders.created', event),
);Без очікування результату HTTP-відповідь може бути сформована раніше, ніж NestJS завершить публікацію події.
@EventPattern()Декоратор @EventPattern() оголошує обробник події:
@EventPattern('orders.created')
handleOrderCreated(event: OrderCreatedEvent) {
// Обробка події
}@Payload() передає значення Kafka-повідомлення в метод обробника.
Якщо потрібні метадані Kafka, використовується KafkaContext:
@EventPattern('orders.created')
handleOrderCreated(
@Payload() event: OrderCreatedEvent,
@Ctx() context: KafkaContext,
) {
const topic = context.getTopic();
const partition = context.getPartition();
const message = context.getMessage();
console.log({
event,
topic,
partition,
offset: message.offset,
});
}Ці дані корисні для:
журналювання;
діагностики;
відстеження партицій;
аналізу offset;
повторної обробки проблемних повідомлень.
У прикладі топік має три партиції, але застосунок запускає один консьюмер. Тому цей консьюмер може отримати всі три партиції.
Якщо запустити кілька екземплярів застосунку з однаковим groupId, Kafka розподілить партиції між ними:
orders-service-group:
instance-1 → partition 0
instance-2 → partition 1
instance-3 → partition 2Кількість активних консьюмерів у групі, більша за кількість партицій, не дає додаткового паралелізму:
3 партиції + 5 консьюмерів = 2 консьюмери не отримують партиційKafka може змінити розподіл партицій, коли консьюмери підключаються або відключаються. Цей процес називається rebalance.
Kafka гарантує порядок лише в межах однієї партиції.
Наприклад:
partition 0:
offset 10 → order.created для замовлення A
offset 11 → order.paid для замовлення AТут порядок збережений.
Але якщо події потрапили в різні партиції, Kafka не гарантує, яка з них буде оброблена першою.
Тому події, для яких важливий порядок, потрібно спрямовувати в одну партицію за однаковим ключем. Наприклад, усі події одного замовлення мають використовувати orderId як ключ повідомлення.
У наведеному прикладі ключ явно не задається, тому Kafka розподіляє повідомлення між партиціями за стандартною стратегією продюсера.
Потокова обробка означає, що сервіс обробляє події в міру їх надходження, а не чекає завершення всієї вибірки.
Для кожної події NestJS викликає обробник:
подія 1 → handleOrderCreated()
подія 2 → handleOrderCreated()
подія 3 → handleOrderCreated()Обробник повинен бути коротким і передбачуваним. Якщо він виконує тривалу операцію, це впливає на швидкість обробки партиції.
Для потокових обробників важливо:
не блокувати event loop довгими синхронними операціями;
безпечно повторювати обробку;
обробляти помилки;
не покладатися на глобальний порядок подій;
зберігати необхідні дані до завершення обробки.
Kafka зазвичай використовується з моделлю доставки at-least-once. Це означає, що подія може бути доставлена повторно, наприклад після помилки або перезапуску консьюмера.
Тому обробник має бути ідемпотентним. Повторна обробка тієї самої події не повинна створювати некоректний результат.
Наприклад, перед створенням платежу можна перевірити, чи не був уже оброблений eventId:
interface OrderCreatedEvent {
eventId: string;
orderId: string;
amount: number;
}Сервіс може зберігати оброблені eventId у базі даних і пропускати дублікати.
Для потокової доставки подій використовується emit():
this.kafkaClient.emit('orders.created', event);Це асинхронна подія, яка не передбачає відповіді від консьюмера.
У NestJS також існує модель request-response із send(), але для потоків подій зазвичай використовують саме emit() та @EventPattern().
У цьому уроці події є односторонніми:
HTTP-запит → publish event → Kafka consumerПродюсеру не потрібно чекати результату бізнес-обробки консьюмером.
Назва в emit() повинна збігатися з назвою в @EventPattern():
this.kafkaClient.emit('orders.created', event);@EventPattern('orders.created')Помилка в одному символі призведе до того, що обробник не отримає подію.
groupIdKafka-консьюмер має бути частиною consumer group. Перевірте, що groupId заданий і стабільний:
consumer: {
groupId: 'orders-service-group',
}Зміна groupId створює нову групу, яка читатиме топік незалежно від попереднього консьюмера.
Порядок між різними партиціями не гарантується. Якщо порядок важливий, пов’язані події повинні потрапляти в одну партицію.
Якщо allowAutoTopicCreation вимкнено, топік потрібно створити заздалегідь. Інакше публікація події завершиться помилкою.
Обробник може отримати ту саму подію більше одного разу. Операції обробника мають бути ідемпотентними.
Додаткові екземпляри не збільшать паралелізм, якщо в топіку недостатньо партицій.
NestJS, запущений на хостовій машині, у локальному прикладі підключається до:
brokers: ['localhost:9092']Якщо NestJS також працює в Docker Compose, потрібно використовувати адресу сервісу всередині Docker-мережі, наприклад:
brokers: ['kafka:9092']Kafka зберігає події в топіках.
Топік складається з однієї або кількох партицій.
Порядок гарантується лише всередині окремої партиції.
Consumer group розподіляє партиції між екземплярами сервісу.
ClientKafka.emit() публікує події.
@EventPattern() обробляє події в NestJS.
KafkaContext дає доступ до топіку, партиції та offset.
Для потокової обробки важливі ідемпотентність і коректна робота з повторною доставкою.
Кількість паралельних консьюмерів обмежена кількістю партицій топіку.