Kafka
Apache Kafka is a distributed event streaming platform capable of handling trillions of events a day. Unlike traditional message queues, Kafka acts as an append-only log, making it the standard choice for Event Sourcing architectures and high-throughput data pipelines.
Overview
While RabbitMQ is like a post office (messages are delivered and then thrown away), Kafka is like a ledger (messages are written permanently, and clients read through the ledger).
NestJS provides a robust Kafka transporter using the kafkajs library. Because Kafka operates differently than standard RPC transporters, configuring it in NestJS requires understanding Consumer Groups, Partitions, and Offsets.
Key Concepts
- Topics: The equivalent of “queues” or “channels”. You publish to a Topic.
- Consumer Groups: If 3 instances of a microservice share the same
groupId, Kafka will distribute the load evenly among them. If they have differentgroupIds, they will all receive the exact same messages (Broadcast). - Offsets: Kafka remembers exactly which message a Consumer Group read last. If a microservice is offline for 3 days, it will reconnect and resume reading exactly where it left off.
Code Examples
1. Installation
Install the required Kafka driver.
npm i kafkajs
2. The Kafka Server (Microservice)
Configure the microservice to connect to a Kafka cluster and join a Consumer Group.
// server/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.createMicroservice<MicroserviceOptions>(
AppModule,
{
transport: Transport.KAFKA,
options: {
client: {
brokers: ['localhost:9092'], // Array of Kafka brokers
},
consumer: {
// IMPORTANT: Services that do the same job must share a groupId!
groupId: 'billing-consumer-group',
},
},
},
);
await app.listen();
}
// server/app.controller.ts
import { Controller } from '@nestjs/common';
import { EventPattern, Payload, Ctx, KafkaContext } from '@nestjs/microservices';
@Controller()
export class AppController {
// Listen to the 'order.created' topic
@EventPattern('order.created')
async handleOrderCreated(@Payload() message: any, @Ctx() context: KafkaContext) {
const originalMessage = context.getMessage();
const partition = context.getPartition();
const topic = context.getTopic();
console.log(`Received message on partition ${partition}:`, message);
// Process billing...
}
}
3. The Kafka Client
When configuring a Kafka client, you only need to specify the brokers.
// client/app.module.ts
import { Module } from '@nestjs/common';
import { ClientsModule, Transport } from '@nestjs/microservices';
@Module({
imports: [
ClientsModule.register([
{
name: 'KAFKA_SERVICE',
transport: Transport.KAFKA,
options: {
client: {
brokers: ['localhost:9092'],
},
producerOnlyMode: true, // Optimizes the client if it ONLY sends messages
},
},
]),
],
})
export class AppModule {}
// client/app.controller.ts
import { Controller, Post, Inject, OnModuleInit } from '@nestjs/common';
import { ClientKafka } from '@nestjs/microservices';
@Controller('orders')
export class AppController implements OnModuleInit {
constructor(
@Inject('KAFKA_SERVICE') private readonly client: ClientKafka,
) {}
// Kafka clients must explicitly connect before sending messages
async onModuleInit() {
this.client.subscribeToResponseOf('order.created'); // Only if expecting a response
await this.client.connect();
}
@Post()
createOrder() {
const orderData = { id: 123, total: 50.0 };
// Emit the event to the 'order.created' topic
// We provide a 'key'. Kafka guarantees that all messages with the same key
// go to the exact same partition, ensuring strict chronological ordering!
this.client.emit('order.created', {
key: 'user_456',
value: orderData,
});
return 'Order queued';
}
}
Best Practices
- Use Keys for Ordering: Kafka guarantees message ordering only within a single partition. If you are processing bank transactions, you MUST ensure deposits and withdrawals for the same account arrive in order. Always use the
accountIdas the Kafkakeywhen emitting the message to guarantee ordering. - Manual Commit: By default, NestJS auto-commits the Kafka offset before the controller method even finishes executing. If your method throws an error, the message is permanently lost. For critical systems, disable
autoCommitand manually commit offsets usingcontext.getConsumer().commitOffsets()only after your database transaction succeeds.