DEV Community

fadi romdhan
fadi romdhan

Posted on

NestJS on Kafka without kafkajs: building a wire-compatible transport

The problem nobody shipped a fix for

If you run NestJS microservices on Kafka, your service talks to the broker through kafkajs. It still works. It also has had no maintainer activity for years, and the NestJS issue asking for an alternative (nestjs/nest#13223) has been open since February 2024 with 60+ comments.

The thread is worth reading because it shows why this stayed unsolved:

  • The NestJS core keeps kafkajs for backwards compatibility. The maintainer's position was: "ideally a new transport strategy, so folks can migrate over time".
  • Two maintained clients appeared: Confluent's @confluentinc/kafka-javascript (librdkafka bindings) and @platformatic/kafka (pure JavaScript, written by Node.js core contributors).
  • A copy-paste strategy for Platformatic was posted in the thread, event-only, and a team going to production described it as "not production ready: full of any, no reconnection support". They wrote their own and could not share it.

So I wrote the missing piece: nestjs-kafka-transport, a drop-in transport for @nestjs/microservices built on @platformatic/kafka.

npm install nestjs-kafka-transport @platformatic/kafka
Enter fullscreen mode Exit fullscreen mode
// before
const app = await NestFactory.createMicroservice(AppModule, {
  transport: Transport.KAFKA,
  options: {
    client: { clientId: 'orders', brokers: ['localhost:9092'] },
    consumer: { groupId: 'orders' },
  },
});

// after
const app = await NestFactory.createMicroservice(AppModule, {
  strategy: new KafkaTransportServer({
    client: { clientId: 'orders', bootstrapBrokers: ['localhost:9092'] },
    consumer: { groupId: 'orders' },
  }),
});
Enter fullscreen mode Exit fullscreen mode

Your controllers do not change. This article is about the part that made "drop-in" true: the wire format.

What "drop-in" has to mean

A transport you can swap one service at a time must produce and consume exactly the same Kafka records as the old one. Otherwise you get a big-bang migration, which is the thing everyone in that thread wanted to avoid.

NestJS's Kafka request-reply is a small protocol on top of plain records. A request looks like this:

Record field Value
topic the pattern, e.g. order.total
value the payload, JSON-encoded if it is an object or array
header kafka_correlationId a unique id per request
header kafka_replyTopic order.total.reply
header kafka_replyPartition the partition of the reply topic the caller owns

And the reply:

Record field Value
topic / partition order.total.reply, the partition from the request
value the handler's return value
header kafka_correlationId copied from the request
header kafka_nest-err present when the handler threw (serialized error)
header kafka_nest-is-disposed present on the last reply of an observable result

The header names come from Spring Kafka, which is why they look the way they do. The parser rules matter too: a value starting with { or [ is JSON-parsed, a value whose first byte is 0 is a Confluent Schema Registry payload and is passed through untouched, everything else stays a string.

nestjs-kafka-transport reproduces all of it. The e2e suite has a test where the built-in ClientKafka (kafkajs) sends requests to the new server, and another where the new client sends requests to the built-in ServerKafka, so both directions are covered on every CI run.

The hard part: who owns the reply partition

Events are easy: subscribe, consume, dispatch. Request-reply has a trap.

The server does not "reply to the caller". It produces a record to the reply topic, to the partition the caller named in kafka_replyPartition. For that to work, every client instance must own at least one partition of every reply topic it subscribed to, and it must know which one before sending the request.

Kafka does not guarantee that. With the default assignment strategy and more clients than partitions, a client may own none, and its replies land somewhere nobody reads. The built-in transport solves this with a custom partition assigner registered on the kafkajs consumer. Until June 2025 @platformatic/kafka had no way to plug one in, which is exactly the feature the maintainer said was missing from his snippet. platformatic/kafka#62 added it, and the assigner in the new transport is about thirty lines:

export function replyPartitionAssigner(
  _current: string,
  members: Map<string, ExtendedGroupProtocolSubscription>,
  topics: Set<string>,
  metadata: ClusterMetadata,
): GroupPartitionsAssignments[] {
  const memberIds = [...members.keys()].sort(); // every member computes the same result
  const assignments = new Map(memberIds.map((id) => [id, new Map<string, GroupAssignment>()]));

  for (const topic of [...topics].sort()) {
    const count = metadata.topics.get(topic)?.partitionsCount ?? 0;
    for (let partition = 0; partition < count; partition++) {
      const memberId = memberIds[partition % memberIds.length]!;
      const mine = assignments.get(memberId)!;
      const existing = mine.get(topic);
      if (existing) existing.partitions.push(partition);
      else mine.set(topic, { topic, partitions: [partition] });
    }
  }
  return memberIds.map((memberId) => ({ memberId, assignments: assignments.get(memberId)! }));
}
Enter fullscreen mode Exit fullscreen mode

Round-robin over sorted member ids. After each group join the client reads its assignments and stamps the first partition it owns on every outgoing request. If the group rebalances while a request is in flight, the reply may land on a partition the client no longer owns; that request times out and the caller retries. The built-in transport has the same window; it is inherent to the design.

Retries: an exception that returns itself

NestJS has KafkaRetriableException: throw it from a handler and the record is redelivered instead of being answered with an error. Implementing it taught me a Nest internal I did not know.

Every handler is wrapped by Nest's RPC exceptions handler, which turns exceptions into an erroring observable. For a normal RpcException the error value is exception.getError(), so the class is lost. KafkaRetriableException overrides getError() to return this:

class KafkaRetriableException extends RpcException {
  getError() {
    return this;
  }
}
Enter fullscreen mode Exit fullscreen mode

That is how it survives the filter and can be instanceof-checked by the transport. The new transport does exactly what the built-in one does: it waits for the handler result (including observables), and if the failure is a KafkaRetriableException it re-runs the handler with exponential backoff (retriableAttempts, retriableDelay) and only commits the offset afterwards. Any other error is answered to the caller (requests) or logged (events), and the record is committed so the partition is not blocked.

Commits are per record, after the handler finished: at-least-once. Records of one partition stay in order; consumer.concurrency only overlaps different partitions.

One deliberate deviation

Send the number 6 through the built-in transport and the other side receives the string "6". Values are toString()-ed on the way out and only {/[ payloads are parsed on the way in. Every team I know has a Number(...) somewhere because of this.

The new transport JSON-encodes numbers, booleans and null as well, and stamps a kafka_nest-content-type: application/json header on such records. A receiver on the new transport decodes them back to their type. A receiver on the built-in transport ignores the header and sees the same strings it always saw. Compatibility preserved, wart removed.

Proving it, locally

git clone https://github.com/fadiroot/nestjs-kafka-transport
cd nestjs-kafka-transport && pnpm install
docker compose up -d kafka      # single-node Kafka 3.9 (KRaft)
pnpm test:e2e                   # request-reply, events, RegExp patterns, retries, kafkajs interop
Enter fullscreen mode Exit fullscreen mode

CI runs the same suite on Node 22 and 24 against Kafka 3.9.1 and 4.0.0.

Migrating

  1. Migrate consumers first, producers last. A new-transport consumer understands everything an old producer sends.
  2. Replace Transport.KAFKA + options with strategy: new KafkaTransportServer(options); replace ClientKafka with KafkaTransportClient. subscribeToResponseOf, send, emit, connect, close keep their names.
  3. Rename a handful of options (brokers → bootstrapBrokers, ssl → tls, subscribe.fromBeginning → consumer.mode: 'earliest'). The full option-by-option table is in docs/migration.md.
  4. Watch for two behaviour differences: KafkaContext.getConsumer() now returns a Platformatic consumer, and RegExp patterns are resolved against the topics that exist at startup.

Requirements: Node ≥ 22.22 (Platformatic's floor), @nestjs/microservices 10 or 11.

What is next

This is v0.1: the built-in transport's feature set, done properly, typed all the way down (no any in the public API). The roadmap is driven by what production users ask for:

  • retry topics and dead-letter topics with the kafka_dlt-* headers Nest already defines,
  • manual commits (ctx.commit()) and a @nestjs/terminus health indicator,
  • batch handlers, Prometheus metrics and OpenTelemetry spans,
  • a Confluent (librdkafka) adapter behind the same interface.

If you run Kafka with NestJS, try it on one consumer and tell me what breaks: issues. And if you are one of the people who commented on nestjs/nest#13223 over the last two years, this one is for you.


Fadi Romdhan is a full-stack engineer (NestJS, Node.js, Python) and open-source contributor to NestJS, Mongoose and the MCP SDKs. Package: npm · GitHub.

Top comments (1)

Collapse
 
raju_dandigam profile image
Raju Dandigam •

@fadiroot The bidirectional interop tests are what make the incremental-migration claim convincing. The in-flight rebalance window is an important boundary: the handler may have committed a side effect even though the caller loses its reply partition. Do you expose that as an unknown outcome rather than a straightforward retryable failure? Keeping an application operation key stable across retries, separately from each request's correlation ID, would help avoid duplicate writes during that window.