<?xml version="1.0" encoding="UTF-8"?>
<rss version="2.0" xmlns:atom="http://www.w3.org/2005/Atom" xmlns:dc="http://purl.org/dc/elements/1.1/">
  <channel>
    <title>DEV Community: fadi romdhan</title>
    <description>The latest articles on DEV Community by fadi romdhan (@fadiroot).</description>
    <link>https://dev.to/fadiroot</link>
    <image>
      <url>https://media2.dev.to/dynamic/image/width=90,height=90,fit=cover,gravity=auto,format=auto/https:%2F%2Fdev-to-uploads.s3.us-east-2.amazonaws.com%2Fuploads%2Fuser%2Fprofile_image%2F1547282%2F74e2be02-4c28-4f29-aba1-6539ff54f412.jpeg</url>
      <title>DEV Community: fadi romdhan</title>
      <link>https://dev.to/fadiroot</link>
    </image>
    <atom:link rel="self" type="application/rss+xml" href="https://dev.to/feed/fadiroot"/>
    <language>en</language>
    <item>
      <title>NestJS on Kafka without kafkajs: building a wire-compatible transport</title>
      <dc:creator>fadi romdhan</dc:creator>
      <pubDate>Mon, 28 Sep 2026 12:50:28 +0000</pubDate>
      <link>https://dev.to/fadiroot/nestjs-on-kafka-without-kafkajs-building-a-wire-compatible-transport-366l</link>
      <guid>https://dev.to/fadiroot/nestjs-on-kafka-without-kafkajs-building-a-wire-compatible-transport-366l</guid>
      <description>&lt;h2&gt;
  
  
  The problem nobody shipped a fix for
&lt;/h2&gt;

&lt;p&gt;If you run NestJS microservices on Kafka, your service talks to the broker through &lt;code&gt;kafkajs&lt;/code&gt;. It still works. It also has had no maintainer activity for years, and the NestJS issue asking for an alternative (&lt;a href="https://github.com/nestjs/nest/issues/13223" rel="noopener noreferrer"&gt;nestjs/nest#13223&lt;/a&gt;) has been open since February 2024 with 60+ comments.&lt;/p&gt;

&lt;p&gt;The thread is worth reading because it shows why this stayed unsolved:&lt;/p&gt;

&lt;ul&gt;
&lt;li&gt;The NestJS core keeps kafkajs for backwards compatibility. The maintainer's position was: "ideally a new transport strategy, so folks can migrate over time".&lt;/li&gt;
&lt;li&gt;Two maintained clients appeared: Confluent's &lt;code&gt;@confluentinc/kafka-javascript&lt;/code&gt; (librdkafka bindings) and &lt;code&gt;@platformatic/kafka&lt;/code&gt; (pure JavaScript, written by Node.js core contributors).&lt;/li&gt;
&lt;li&gt;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 &lt;code&gt;any&lt;/code&gt;, no reconnection support". They wrote their own and could not share it.&lt;/li&gt;
&lt;/ul&gt;

&lt;p&gt;So I wrote the missing piece: &lt;a href="https://github.com/fadiroot/nestjs-kafka-transport" rel="noopener noreferrer"&gt;&lt;code&gt;nestjs-kafka-transport&lt;/code&gt;&lt;/a&gt;, a drop-in transport for &lt;code&gt;@nestjs/microservices&lt;/code&gt; built on &lt;code&gt;@platformatic/kafka&lt;/code&gt;.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight shell"&gt;&lt;code&gt;npm &lt;span class="nb"&gt;install &lt;/span&gt;nestjs-kafka-transport @platformatic/kafka
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;





&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight typescript"&gt;&lt;code&gt;&lt;span class="c1"&gt;// before&lt;/span&gt;
&lt;span class="kd"&gt;const&lt;/span&gt; &lt;span class="nx"&gt;app&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="k"&gt;await&lt;/span&gt; &lt;span class="nx"&gt;NestFactory&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;createMicroservice&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="nx"&gt;AppModule&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="p"&gt;{&lt;/span&gt;
  &lt;span class="na"&gt;transport&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt; &lt;span class="nx"&gt;Transport&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nx"&gt;KAFKA&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;
  &lt;span class="na"&gt;options&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt; &lt;span class="p"&gt;{&lt;/span&gt;
    &lt;span class="na"&gt;client&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt; &lt;span class="p"&gt;{&lt;/span&gt; &lt;span class="na"&gt;clientId&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt; &lt;span class="dl"&gt;'&lt;/span&gt;&lt;span class="s1"&gt;orders&lt;/span&gt;&lt;span class="dl"&gt;'&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="na"&gt;brokers&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt; &lt;span class="p"&gt;[&lt;/span&gt;&lt;span class="dl"&gt;'&lt;/span&gt;&lt;span class="s1"&gt;localhost:9092&lt;/span&gt;&lt;span class="dl"&gt;'&lt;/span&gt;&lt;span class="p"&gt;]&lt;/span&gt; &lt;span class="p"&gt;},&lt;/span&gt;
    &lt;span class="na"&gt;consumer&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt; &lt;span class="p"&gt;{&lt;/span&gt; &lt;span class="na"&gt;groupId&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt; &lt;span class="dl"&gt;'&lt;/span&gt;&lt;span class="s1"&gt;orders&lt;/span&gt;&lt;span class="dl"&gt;'&lt;/span&gt; &lt;span class="p"&gt;},&lt;/span&gt;
  &lt;span class="p"&gt;},&lt;/span&gt;
&lt;span class="p"&gt;});&lt;/span&gt;

&lt;span class="c1"&gt;// after&lt;/span&gt;
&lt;span class="kd"&gt;const&lt;/span&gt; &lt;span class="nx"&gt;app&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="k"&gt;await&lt;/span&gt; &lt;span class="nx"&gt;NestFactory&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;createMicroservice&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="nx"&gt;AppModule&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="p"&gt;{&lt;/span&gt;
  &lt;span class="na"&gt;strategy&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt; &lt;span class="k"&gt;new&lt;/span&gt; &lt;span class="nc"&gt;KafkaTransportServer&lt;/span&gt;&lt;span class="p"&gt;({&lt;/span&gt;
    &lt;span class="na"&gt;client&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt; &lt;span class="p"&gt;{&lt;/span&gt; &lt;span class="na"&gt;clientId&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt; &lt;span class="dl"&gt;'&lt;/span&gt;&lt;span class="s1"&gt;orders&lt;/span&gt;&lt;span class="dl"&gt;'&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="na"&gt;bootstrapBrokers&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt; &lt;span class="p"&gt;[&lt;/span&gt;&lt;span class="dl"&gt;'&lt;/span&gt;&lt;span class="s1"&gt;localhost:9092&lt;/span&gt;&lt;span class="dl"&gt;'&lt;/span&gt;&lt;span class="p"&gt;]&lt;/span&gt; &lt;span class="p"&gt;},&lt;/span&gt;
    &lt;span class="na"&gt;consumer&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt; &lt;span class="p"&gt;{&lt;/span&gt; &lt;span class="na"&gt;groupId&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt; &lt;span class="dl"&gt;'&lt;/span&gt;&lt;span class="s1"&gt;orders&lt;/span&gt;&lt;span class="dl"&gt;'&lt;/span&gt; &lt;span class="p"&gt;},&lt;/span&gt;
  &lt;span class="p"&gt;}),&lt;/span&gt;
&lt;span class="p"&gt;});&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;Your controllers do not change. This article is about the part that made "drop-in" true: the wire format.&lt;/p&gt;

&lt;h2&gt;
  
  
  What "drop-in" has to mean
&lt;/h2&gt;

&lt;p&gt;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.&lt;/p&gt;

&lt;p&gt;NestJS's Kafka request-reply is a small protocol on top of plain records. A request looks like this:&lt;/p&gt;

&lt;div class="table-wrapper-paragraph"&gt;&lt;table&gt;
&lt;thead&gt;
&lt;tr&gt;
&lt;th&gt;Record field&lt;/th&gt;
&lt;th&gt;Value&lt;/th&gt;
&lt;/tr&gt;
&lt;/thead&gt;
&lt;tbody&gt;
&lt;tr&gt;
&lt;td&gt;topic&lt;/td&gt;
&lt;td&gt;the pattern, e.g. &lt;code&gt;order.total&lt;/code&gt;
&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;value&lt;/td&gt;
&lt;td&gt;the payload, JSON-encoded if it is an object or array&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;header &lt;code&gt;kafka_correlationId&lt;/code&gt;
&lt;/td&gt;
&lt;td&gt;a unique id per request&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;header &lt;code&gt;kafka_replyTopic&lt;/code&gt;
&lt;/td&gt;
&lt;td&gt;&lt;code&gt;order.total.reply&lt;/code&gt;&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;header &lt;code&gt;kafka_replyPartition&lt;/code&gt;
&lt;/td&gt;
&lt;td&gt;the partition of the reply topic the caller owns&lt;/td&gt;
&lt;/tr&gt;
&lt;/tbody&gt;
&lt;/table&gt;&lt;/div&gt;

&lt;p&gt;And the reply:&lt;/p&gt;

&lt;div class="table-wrapper-paragraph"&gt;&lt;table&gt;
&lt;thead&gt;
&lt;tr&gt;
&lt;th&gt;Record field&lt;/th&gt;
&lt;th&gt;Value&lt;/th&gt;
&lt;/tr&gt;
&lt;/thead&gt;
&lt;tbody&gt;
&lt;tr&gt;
&lt;td&gt;topic / partition&lt;/td&gt;
&lt;td&gt;
&lt;code&gt;order.total.reply&lt;/code&gt;, the partition from the request&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;value&lt;/td&gt;
&lt;td&gt;the handler's return value&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;header &lt;code&gt;kafka_correlationId&lt;/code&gt;
&lt;/td&gt;
&lt;td&gt;copied from the request&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;header &lt;code&gt;kafka_nest-err&lt;/code&gt;
&lt;/td&gt;
&lt;td&gt;present when the handler threw (serialized error)&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;header &lt;code&gt;kafka_nest-is-disposed&lt;/code&gt;
&lt;/td&gt;
&lt;td&gt;present on the last reply of an observable result&lt;/td&gt;
&lt;/tr&gt;
&lt;/tbody&gt;
&lt;/table&gt;&lt;/div&gt;

&lt;p&gt;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 &lt;code&gt;{&lt;/code&gt; or &lt;code&gt;[&lt;/code&gt; is JSON-parsed, a value whose first byte is &lt;code&gt;0&lt;/code&gt; is a Confluent Schema Registry payload and is passed through untouched, everything else stays a string.&lt;/p&gt;

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

&lt;h2&gt;
  
  
  The hard part: who owns the reply partition
&lt;/h2&gt;

&lt;p&gt;Events are easy: subscribe, consume, dispatch. Request-reply has a trap.&lt;/p&gt;

&lt;p&gt;The server does not "reply to the caller". It produces a record to the reply topic, to the partition the caller named in &lt;code&gt;kafka_replyPartition&lt;/code&gt;. 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 &lt;em&gt;before&lt;/em&gt; sending the request.&lt;/p&gt;

&lt;p&gt;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 &lt;code&gt;@platformatic/kafka&lt;/code&gt; had no way to plug one in, which is exactly the feature the maintainer said was missing from his snippet. &lt;a href="https://github.com/platformatic/kafka/pull/62" rel="noopener noreferrer"&gt;platformatic/kafka#62&lt;/a&gt; added it, and the assigner in the new transport is about thirty lines:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight typescript"&gt;&lt;code&gt;&lt;span class="k"&gt;export&lt;/span&gt; &lt;span class="kd"&gt;function&lt;/span&gt; &lt;span class="nf"&gt;replyPartitionAssigner&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;
  &lt;span class="nx"&gt;_current&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt; &lt;span class="kr"&gt;string&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;
  &lt;span class="nx"&gt;members&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt; &lt;span class="nb"&gt;Map&lt;/span&gt;&lt;span class="o"&gt;&amp;lt;&lt;/span&gt;&lt;span class="kr"&gt;string&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="nx"&gt;ExtendedGroupProtocolSubscription&lt;/span&gt;&lt;span class="o"&gt;&amp;gt;&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;
  &lt;span class="nx"&gt;topics&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt; &lt;span class="nb"&gt;Set&lt;/span&gt;&lt;span class="o"&gt;&amp;lt;&lt;/span&gt;&lt;span class="kr"&gt;string&lt;/span&gt;&lt;span class="o"&gt;&amp;gt;&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;
  &lt;span class="nx"&gt;metadata&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt; &lt;span class="nx"&gt;ClusterMetadata&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;
&lt;span class="p"&gt;):&lt;/span&gt; &lt;span class="nx"&gt;GroupPartitionsAssignments&lt;/span&gt;&lt;span class="p"&gt;[]&lt;/span&gt; &lt;span class="p"&gt;{&lt;/span&gt;
  &lt;span class="kd"&gt;const&lt;/span&gt; &lt;span class="nx"&gt;memberIds&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="p"&gt;[...&lt;/span&gt;&lt;span class="nx"&gt;members&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;keys&lt;/span&gt;&lt;span class="p"&gt;()].&lt;/span&gt;&lt;span class="nf"&gt;sort&lt;/span&gt;&lt;span class="p"&gt;();&lt;/span&gt; &lt;span class="c1"&gt;// every member computes the same result&lt;/span&gt;
  &lt;span class="kd"&gt;const&lt;/span&gt; &lt;span class="nx"&gt;assignments&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="k"&gt;new&lt;/span&gt; &lt;span class="nc"&gt;Map&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="nx"&gt;memberIds&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;map&lt;/span&gt;&lt;span class="p"&gt;((&lt;/span&gt;&lt;span class="nx"&gt;id&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt; &lt;span class="o"&gt;=&amp;gt;&lt;/span&gt; &lt;span class="p"&gt;[&lt;/span&gt;&lt;span class="nx"&gt;id&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="k"&gt;new&lt;/span&gt; &lt;span class="nb"&gt;Map&lt;/span&gt;&lt;span class="o"&gt;&amp;lt;&lt;/span&gt;&lt;span class="kr"&gt;string&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="nx"&gt;GroupAssignment&lt;/span&gt;&lt;span class="o"&gt;&amp;gt;&lt;/span&gt;&lt;span class="p"&gt;()]));&lt;/span&gt;

  &lt;span class="k"&gt;for &lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="kd"&gt;const&lt;/span&gt; &lt;span class="nx"&gt;topic&lt;/span&gt; &lt;span class="k"&gt;of&lt;/span&gt; &lt;span class="p"&gt;[...&lt;/span&gt;&lt;span class="nx"&gt;topics&lt;/span&gt;&lt;span class="p"&gt;].&lt;/span&gt;&lt;span class="nf"&gt;sort&lt;/span&gt;&lt;span class="p"&gt;())&lt;/span&gt; &lt;span class="p"&gt;{&lt;/span&gt;
    &lt;span class="kd"&gt;const&lt;/span&gt; &lt;span class="nx"&gt;count&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="nx"&gt;metadata&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nx"&gt;topics&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;get&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="nx"&gt;topic&lt;/span&gt;&lt;span class="p"&gt;)?.&lt;/span&gt;&lt;span class="nx"&gt;partitionsCount&lt;/span&gt; &lt;span class="o"&gt;??&lt;/span&gt; &lt;span class="mi"&gt;0&lt;/span&gt;&lt;span class="p"&gt;;&lt;/span&gt;
    &lt;span class="k"&gt;for &lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="kd"&gt;let&lt;/span&gt; &lt;span class="nx"&gt;partition&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="mi"&gt;0&lt;/span&gt;&lt;span class="p"&gt;;&lt;/span&gt; &lt;span class="nx"&gt;partition&lt;/span&gt; &lt;span class="o"&gt;&amp;lt;&lt;/span&gt; &lt;span class="nx"&gt;count&lt;/span&gt;&lt;span class="p"&gt;;&lt;/span&gt; &lt;span class="nx"&gt;partition&lt;/span&gt;&lt;span class="o"&gt;++&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt; &lt;span class="p"&gt;{&lt;/span&gt;
      &lt;span class="kd"&gt;const&lt;/span&gt; &lt;span class="nx"&gt;memberId&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="nx"&gt;memberIds&lt;/span&gt;&lt;span class="p"&gt;[&lt;/span&gt;&lt;span class="nx"&gt;partition&lt;/span&gt; &lt;span class="o"&gt;%&lt;/span&gt; &lt;span class="nx"&gt;memberIds&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nx"&gt;length&lt;/span&gt;&lt;span class="p"&gt;]&lt;/span&gt;&lt;span class="o"&gt;!&lt;/span&gt;&lt;span class="p"&gt;;&lt;/span&gt;
      &lt;span class="kd"&gt;const&lt;/span&gt; &lt;span class="nx"&gt;mine&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="nx"&gt;assignments&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;get&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="nx"&gt;memberId&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;&lt;span class="o"&gt;!&lt;/span&gt;&lt;span class="p"&gt;;&lt;/span&gt;
      &lt;span class="kd"&gt;const&lt;/span&gt; &lt;span class="nx"&gt;existing&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="nx"&gt;mine&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;get&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="nx"&gt;topic&lt;/span&gt;&lt;span class="p"&gt;);&lt;/span&gt;
      &lt;span class="k"&gt;if &lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="nx"&gt;existing&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt; &lt;span class="nx"&gt;existing&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nx"&gt;partitions&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;push&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="nx"&gt;partition&lt;/span&gt;&lt;span class="p"&gt;);&lt;/span&gt;
      &lt;span class="k"&gt;else&lt;/span&gt; &lt;span class="nx"&gt;mine&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;set&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="nx"&gt;topic&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="p"&gt;{&lt;/span&gt; &lt;span class="nx"&gt;topic&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="na"&gt;partitions&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt; &lt;span class="p"&gt;[&lt;/span&gt;&lt;span class="nx"&gt;partition&lt;/span&gt;&lt;span class="p"&gt;]&lt;/span&gt; &lt;span class="p"&gt;});&lt;/span&gt;
    &lt;span class="p"&gt;}&lt;/span&gt;
  &lt;span class="p"&gt;}&lt;/span&gt;
  &lt;span class="k"&gt;return&lt;/span&gt; &lt;span class="nx"&gt;memberIds&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;map&lt;/span&gt;&lt;span class="p"&gt;((&lt;/span&gt;&lt;span class="nx"&gt;memberId&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt; &lt;span class="o"&gt;=&amp;gt;&lt;/span&gt; &lt;span class="p"&gt;({&lt;/span&gt; &lt;span class="nx"&gt;memberId&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="na"&gt;assignments&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt; &lt;span class="nx"&gt;assignments&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;get&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="nx"&gt;memberId&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;&lt;span class="o"&gt;!&lt;/span&gt; &lt;span class="p"&gt;}));&lt;/span&gt;
&lt;span class="p"&gt;}&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;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.&lt;/p&gt;

&lt;h2&gt;
  
  
  Retries: an exception that returns itself
&lt;/h2&gt;

&lt;p&gt;NestJS has &lt;code&gt;KafkaRetriableException&lt;/code&gt;: 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.&lt;/p&gt;

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

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight typescript"&gt;&lt;code&gt;&lt;span class="kd"&gt;class&lt;/span&gt; &lt;span class="nc"&gt;KafkaRetriableException&lt;/span&gt; &lt;span class="kd"&gt;extends&lt;/span&gt; &lt;span class="nc"&gt;RpcException&lt;/span&gt; &lt;span class="p"&gt;{&lt;/span&gt;
  &lt;span class="nf"&gt;getError&lt;/span&gt;&lt;span class="p"&gt;()&lt;/span&gt; &lt;span class="p"&gt;{&lt;/span&gt;
    &lt;span class="k"&gt;return&lt;/span&gt; &lt;span class="k"&gt;this&lt;/span&gt;&lt;span class="p"&gt;;&lt;/span&gt;
  &lt;span class="p"&gt;}&lt;/span&gt;
&lt;span class="p"&gt;}&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;That is how it survives the filter and can be &lt;code&gt;instanceof&lt;/code&gt;-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 &lt;code&gt;KafkaRetriableException&lt;/code&gt; it re-runs the handler with exponential backoff (&lt;code&gt;retriableAttempts&lt;/code&gt;, &lt;code&gt;retriableDelay&lt;/code&gt;) 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.&lt;/p&gt;

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

&lt;h2&gt;
  
  
  One deliberate deviation
&lt;/h2&gt;

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

&lt;p&gt;The new transport JSON-encodes numbers, booleans and &lt;code&gt;null&lt;/code&gt; as well, and stamps a &lt;code&gt;kafka_nest-content-type: application/json&lt;/code&gt; 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.&lt;/p&gt;

&lt;h2&gt;
  
  
  Proving it, locally
&lt;/h2&gt;



&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight shell"&gt;&lt;code&gt;git clone https://github.com/fadiroot/nestjs-kafka-transport
&lt;span class="nb"&gt;cd &lt;/span&gt;nestjs-kafka-transport &lt;span class="o"&gt;&amp;amp;&amp;amp;&lt;/span&gt; pnpm &lt;span class="nb"&gt;install
&lt;/span&gt;docker compose up &lt;span class="nt"&gt;-d&lt;/span&gt; kafka      &lt;span class="c"&gt;# single-node Kafka 3.9 (KRaft)&lt;/span&gt;
pnpm &lt;span class="nb"&gt;test&lt;/span&gt;:e2e                   &lt;span class="c"&gt;# request-reply, events, RegExp patterns, retries, kafkajs interop&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;CI runs the same suite on Node 22 and 24 against Kafka 3.9.1 and 4.0.0.&lt;/p&gt;

&lt;h2&gt;
  
  
  Migrating
&lt;/h2&gt;

&lt;ol&gt;
&lt;li&gt;Migrate consumers first, producers last. A new-transport consumer understands everything an old producer sends.&lt;/li&gt;
&lt;li&gt;Replace &lt;code&gt;Transport.KAFKA&lt;/code&gt; + &lt;code&gt;options&lt;/code&gt; with &lt;code&gt;strategy: new KafkaTransportServer(options)&lt;/code&gt;; replace &lt;code&gt;ClientKafka&lt;/code&gt; with &lt;code&gt;KafkaTransportClient&lt;/code&gt;. &lt;code&gt;subscribeToResponseOf&lt;/code&gt;, &lt;code&gt;send&lt;/code&gt;, &lt;code&gt;emit&lt;/code&gt;, &lt;code&gt;connect&lt;/code&gt;, &lt;code&gt;close&lt;/code&gt; keep their names.&lt;/li&gt;
&lt;li&gt;Rename a handful of options (&lt;code&gt;brokers&lt;/code&gt; → &lt;code&gt;bootstrapBrokers&lt;/code&gt;, &lt;code&gt;ssl&lt;/code&gt; → &lt;code&gt;tls&lt;/code&gt;, &lt;code&gt;subscribe.fromBeginning&lt;/code&gt; → &lt;code&gt;consumer.mode: 'earliest'&lt;/code&gt;). The full option-by-option table is in &lt;a href="https://github.com/fadiroot/nestjs-kafka-transport/blob/main/docs/migration.md" rel="noopener noreferrer"&gt;docs/migration.md&lt;/a&gt;.&lt;/li&gt;
&lt;li&gt;Watch for two behaviour differences: &lt;code&gt;KafkaContext.getConsumer()&lt;/code&gt; now returns a Platformatic consumer, and RegExp patterns are resolved against the topics that exist at startup.&lt;/li&gt;
&lt;/ol&gt;

&lt;p&gt;Requirements: Node ≥ 22.22 (Platformatic's floor), &lt;code&gt;@nestjs/microservices&lt;/code&gt; 10 or 11.&lt;/p&gt;

&lt;h2&gt;
  
  
  What is next
&lt;/h2&gt;

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

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

&lt;p&gt;If you run Kafka with NestJS, try it on one consumer and tell me what breaks: &lt;a href="https://github.com/fadiroot/nestjs-kafka-transport/issues" rel="noopener noreferrer"&gt;issues&lt;/a&gt;. And if you are one of the people who commented on nestjs/nest#13223 over the last two years, this one is for you.&lt;/p&gt;




&lt;p&gt;&lt;em&gt;Fadi Romdhan is a full-stack engineer (NestJS, Node.js, Python) and open-source contributor to NestJS, Mongoose and the MCP SDKs. Package: &lt;a href="https://www.npmjs.com/package/nestjs-kafka-transport" rel="noopener noreferrer"&gt;npm&lt;/a&gt; · &lt;a href="https://github.com/fadiroot/nestjs-kafka-transport" rel="noopener noreferrer"&gt;GitHub&lt;/a&gt;.&lt;/em&gt;&lt;/p&gt;

</description>
      <category>nestjs</category>
      <category>kafka</category>
      <category>node</category>
      <category>typescript</category>
    </item>
  </channel>
</rss>
