Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.

Some links on this page are affiliate links: if you buy through them we may earn a commission, at no extra cost to you.

The practical way to use Kafka with Node.js is to create a producer, publish records to a topic, and run a consumer in a named consumer group. For production, that basic flow is only the beginning: you also need TLS/SASL authentication, deliberate partitioning, controlled offset commits, idempotent handlers, bounded retries, dead-letter handling, observability, and graceful shutdown.

This guide uses Confluent’s JavaScript client, @confluentinc/kafka-javascript, with a KafkaJS-compatible configuration style. The client is based on librdkafka and provides promisified and callback APIs. Check the current support matrix before choosing a Node.js version, operating system, architecture, or container image.

What you will build

The example publishes an order.created JSON event to an orders topic and consumes it with an orders-service consumer group:

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.
Node.js API
   │
   ├── producer ──> orders topic ──> orders-service group
   │                                      ├── worker 1
   │                                      └── worker 2
   │
   └── analytics group ──> receives the same events independently

Kafka is a distributed event-streaming platform, not simply an in-process queue. Producers write records to topics. Topics contain partitions, and consumers read records by offset. Kafka normally retains records according to topic policy; reading a record does not delete it.

A record can contain a key, value, headers, timestamp, topic, partition, and offset. Ordering is guaranteed within a partition, not across an entire topic. A key commonly determines the partition, so using the same key for all events belonging to one order helps preserve that order.

Consumer groups allow several application instances to share a topic’s partitions. Each partition is assigned to at most one active consumer in a group at a time. A different group receives its own view of the topic and can process the same records independently.

When Kafka is—and is not—the right choice

Kafka is a strong fit for event-driven services, durable asynchronous workflows, fan-out to independent consumers, audit and activity streams, high-volume telemetry, data pipelines, and decoupling an HTTP request from slow downstream work.

Free tools Windows power users keep installed

One-click scans. No signup required.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.

It may be excessive for a small application that needs only a few delayed jobs, a simple in-process queue, or strictly synchronous request/response behavior. Redis, a database-backed queue, or a cloud task service may be simpler when replay, multiple consumer groups, high throughput, and durable retention are not requirements. Kafka adds infrastructure, serialization, partitioning, offset, consumer-group, security, and operational concerns.

Choose a Node.js Kafka client

Confluent JavaScript client

This guide uses @confluentinc/kafka-javascript. Confluent describes it as a JavaScript client built on librdkafka, with a promisified API, callback support, and compatibility with KafkaJS-style patterns. Its native foundation can provide mature protocol behavior, but it also means you should verify prebuilt-binary availability for your Node.js version, operating system, CPU architecture, container base image, and CI runner.

Install it with:

npm install @confluentinc/kafka-javascript

KafkaJS

KafkaJS remains a reasonable choice when a project already uses it, has existing wrappers and examples, or prefers a JavaScript-native implementation. Do not assume that KafkaJS and Confluent’s client have identical defaults, retry behavior, transaction options, subscription behavior, or configuration nesting. Follow the documentation for the package actually installed. Confluent provides migration guidance.

Set up a broker

For the shortest path through this tutorial, use Confluent Cloud or another reachable Kafka service. In the Confluent Cloud console, select an environment and cluster, choose Clients, select JavaScript, create or use API keys, and copy the generated configuration. Confluent Cloud documents TLS 1.2 and SASL/PLAIN or SASL/OAUTHBEARER authentication requirements; its Kafka connections also require correct SNI behavior.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.

A local broker is useful for learning and integration tests. A typical local configuration uses an address such as localhost:9092 and omits TLS and SASL, but local commands and advertised-listener settings vary by Kafka distribution and version. A single local broker also does not represent production availability, replication, networking, or security.

Create the topic through deployment or infrastructure automation when possible. Define its partition count, replication, retention, access policy, and naming deliberately. Do not rely blindly on automatic topic creation in production: its behavior is affected by client, broker, and provider configuration.

Create the Node.js project

mkdir node-kafka-example
cd node-kafka-example
npm init -y
npm install @confluentinc/kafka-javascript

Set configuration outside the source tree:

export KAFKA_BROKERS="your-bootstrap-server"
export KAFKA_USERNAME="your-api-key"
export KAFKA_PASSWORD="your-api-secret"
export KAFKA_TOPIC="orders"
export KAFKA_GROUP_ID="orders-service"

Use a secret manager in production. Never commit API secrets, certificates, .env files, or generated cloud configuration to source control.

Configure the Kafka connection

// kafka.js
const { Kafka } =
  require("@confluentinc/kafka-javascript").KafkaJS;

const brokers = process.env.KAFKA_BROKERS
  .split(",")
  .map((value) => value.trim());

const kafka = new Kafka({
  kafkaJS: {
    brokers,
    ssl: true,
    sasl: {
      mechanism: "plain",
      username: process.env.KAFKA_USERNAME,
      password: process.env.KAFKA_PASSWORD,
    },
    clientId: "node-kafka-example",
  },
});

module.exports = { kafka };

For a local unauthenticated broker, use its local bootstrap address and normally omit ssl and sasl. For a managed service, copy the provider’s generated settings rather than guessing them.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.

Before troubleshooting Kafka itself, verify that the variables exist without exposing the secret:

console.log({
  brokers,
  hasUsername: Boolean(process.env.KAFKA_USERNAME),
  hasPassword: Boolean(process.env.KAFKA_PASSWORD),
});

Publish an event

// producer.js
const { kafka } = require("./kafka");

async function main() {
  const producer = kafka.producer();
  await producer.connect();

  try {
    const order = {
      orderId: "order-123",
      customerId: "customer-456",
      total: 49.99,
      createdAt: new Date().toISOString(),
    };

    const result = await producer.send({
      topic: process.env.KAFKA_TOPIC,
      messages: [{
        key: order.orderId,
        value: JSON.stringify(order),
        headers: {
          "content-type": "application/json",
          "event-type": "order.created",
        },
      }],
    });

    console.log("Published:", result);
  } finally {
    await producer.disconnect();
  }
}

main().catch((error) => {
  console.error(error);
  process.exitCode = 1;
});

Kafka stores bytes. JSON is an application convention, not a Kafka requirement. The key affects partition selection, while headers can carry an event type, schema version, correlation ID, or tracing metadata. In a long-running HTTP service, connect once and reuse the producer; do not open and close a Kafka connection for every request.

Consume events

// consumer.js
const { kafka } = require("./kafka");

async function main() {
  const consumer = kafka.consumer({
    kafkaJS: {
      groupId: process.env.KAFKA_GROUP_ID,
      fromBeginning: false,
    },
  });

  await consumer.connect();
  await consumer.subscribe({
    topics: [process.env.KAFKA_TOPIC],
  });

  await consumer.run({
    eachMessage: async ({ topic, partition, message }) => {
      const rawValue = message.value?.toString();

      if (!rawValue) {
        console.warn("Skipping empty message", {
          topic,
          partition,
          offset: message.offset,
        });
        return;
      }

      const order = JSON.parse(rawValue);

      console.log({
        topic,
        partition,
        offset: message.offset,
        key: message.key?.toString(),
        order,
      });

      // Perform the business operation here.
    },
  });
}

main().catch((error) => {
  console.error(error);
  process.exitCode = 1;
});

A consumer needs a stable groupId. fromBeginning: false means the consumer normally starts at the group’s current position, or at the end when no committed position exists. A new group ID has its own offsets and can read the same topic independently.

Run the consumer before the producer:

node consumer.js
node producer.js

The consumer should print the topic, partition, offset, key, and parsed event. Start another consumer with the same group ID to see partitions shared between instances. Start one with a different group ID to receive the event independently. Parallel assignment is visible only when the topic has enough partitions and there is sufficient work.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.

Shut down cleanly

Containers commonly send SIGTERM before stopping a process. Keep the consumer in scope and disconnect it during shutdown:

async function shutdown(signal) {
  console.log(`Received ${signal}; shutting down`);
  try {
    await consumer.disconnect();
    process.exit(0);
  } catch (error) {
    console.error("Shutdown failed", error);
    process.exit(1);
  }
}

process.once("SIGINT", () => shutdown("SIGINT"));
process.once("SIGTERM", () => shutdown("SIGTERM"));

A production shutdown should stop accepting new work, stop fetching records, finish or cancel in-flight work, advance offsets only for successful work, disconnect producers and consumers, and exit within the orchestrator’s termination grace period.

Understand offsets and delivery guarantees

An offset identifies a record’s position within one partition. Advancing an offset is not the same as completing a business operation.

  • At-most-once: advance or acknowledge before the operation. A crash can lose work.
  • At-least-once: complete the operation and then commit the offset. A crash between those actions can cause redelivery.
  • Exactly-once: requires a specific Kafka transaction design and compatible downstream behavior. An ordinary producer flag does not make an HTTP call or database update exactly once.

Automatic offset commits are convenient for demonstrations but can be unsafe for slow or non-idempotent handlers. Use the selected client’s documented manual or controlled commit API when offset advancement must follow successful processing. Be careful when copying examples between KafkaJS and the Confluent client because their exact APIs and semantics can differ.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.

Design most business handlers for at-least-once delivery. Include a stable event or operation ID, store it with a database uniqueness constraint, make updates conditional where appropriate, and use idempotency keys for external APIs that support them. A crash after a side effect but before the commit is normal failure behavior, not proof that Kafka lost the record.

Retries, poison messages, and dead-letter handling

Separate failure types:

  1. Transport retries: temporary broker or network failures handled by the client.
  2. Business retries: database deadlocks, rate limits, or temporary downstream outages.
  3. Permanent failures: malformed JSON, invalid schema, or impossible business state.
  4. Quarantine or dead-letter flow: preserve the original payload and diagnostic context for inspection or replay.

A useful flow is:

Kafka record
   ↓
Validate and deserialize
   ↓
Business handler
   ├── success → commit
   ├── transient failure → bounded retry
   └── permanent failure → dead-letter/quarantine, then commit original

Do not retry malformed messages forever. A poison message can repeatedly block progress through its partition. Client retries address transport conditions, not every business failure. The Confluent client’s documented KafkaJS-compatible defaults include a 300 ms initial retry backoff, 30-second maximum backoff, five producer retries, exponential multiplier 2, jitter 0.2, and consumer restart-on-failure enabled; treat these as client-specific defaults, not universal Kafka behavior.

Scale consumer groups correctly

Partitions limit parallelism. If a topic has three partitions, adding a fourth consumer to the same group cannot create a fourth active partition worker for that topic. More instances can also trigger rebalances when members join, leave, or assignments change.

Long handlers can cause heartbeat or processing-time problems, especially when they block the Node.js event loop. Do not perform CPU-heavy synchronous work in eachMessage. Limit concurrency deliberately, apply backpressure when a database or API slows down, and avoid unbounded promise creation or in-memory buffering.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.

Monitor consumer lag and processing duration rather than inferring health from application logs. The client documentation lists compatibility settings such as a 300,000 ms rebalance timeout and 3,000 ms heartbeat interval, but the right values depend on handler duration, partition count, deployment model, and broker configuration. Tune them only after measuring.

Choose serialization and evolve event contracts

JSON is an excellent starting point because it is easy to inspect, but it does not provide contract governance. For larger systems, Avro, Protobuf, or JSON Schema combined with a schema registry can enforce and document compatibility.

Useful event-envelope fields include:

  • Stable event ID.
  • Event type such as order.created.
  • Schema version.
  • Producer or service name.
  • Event timestamp.
  • Correlation or trace ID when needed.

Prefer additive changes that older consumers can tolerate. Consumers should generally ignore unknown fields when safe. Do not silently change the meaning or unit of an existing field. Confluent’s client and cloud configuration documentation distinguishes Kafka credentials from optional Schema Registry configuration.

Independent reader supportYour contribution helps us test, update, and keep practical guides available for everyone.Support on Ko-Fi

Secure Kafka connections

For Confluent Cloud, use TLS and the authentication method supplied by the service:

What’s actually slowing this PC down?

Pick the symptom - the matching free tool is one click away.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.
const kafka = new Kafka({
  kafkaJS: {
    brokers: [process.env.KAFKA_BROKER],
    ssl: true,
    sasl: {
      mechanism: "plain",
      username: process.env.KAFKA_API_KEY,
      password: process.env.KAFKA_API_SECRET,
    },
  },
});

Keep credentials in a secret manager, use least-privilege ACLs, separate producer and consumer credentials where practical, restrict topic access, and avoid logging sensitive payloads. Do not pin an intermediate certificate unless the provider explicitly requires it; certificate chains can change. Proxies must preserve the TLS SNI information required by some managed Kafka services.

Local Kafka versus managed Kafka

Option Advantages Trade-offs
Local Kafka Fast development, repeatable tests, no cloud account or usage bill Version-sensitive setup, advertised-listener problems, usually insecure defaults, and no production availability model
Managed Kafka Hosted brokers, upgrades, replication, TLS/SASL testing, team access Usage charges, IAM/API-key work, network and egress concerns, provider-specific behavior

Confluent Cloud is a convenient managed path for this tutorial, but the client and service are separate decisions: the same JavaScript client can connect to self-managed Kafka or another compatible provider. Confluent’s documentation has advertised a free-credit promotion, but promotions change. Cost depends on cloud, region, service tier, throughput, storage, retention, networking, connectors, and optional services. Check the official pricing page for a current estimate.

Other selection candidates include Amazon MSK, Azure Event Hubs with its Kafka endpoint, Google Cloud Managed Service for Apache Kafka, Aiven for Apache Kafka, and Redpanda. Compatibility, supported features, pricing, regions, and operational models differ; do not treat them as interchangeable without checking the specific APIs you need.

Troubleshoot common failures

ECONNREFUSED

Check that the broker is running, the port is correct, and the address is reachable from the same environment as Node.js. In Docker, a hostname that works inside the container may not work on the host. Verify advertised listeners, firewalls, security groups, cloud bootstrap addresses, and TLS settings separately.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.

Authentication failure

Check for reversed or unset credentials, the wrong SASL mechanism, credentials for another cluster, and missing topic permissions. Log only safe diagnostics such as broker names and whether username and password variables are present.

TLS or certificate failure

Verify that ssl: true is enabled where required, the trust store is current, no unsupported certificate pinning is in place, and a proxy preserves SNI. Confluent Cloud documents TLS 1.2 and its supported SASL mechanisms.

No messages received

Confirm that both applications use the same cluster and topic. Check the group ID, fromBeginning behavior, committed offset range, topic permissions, partition assignment, and whether the producer actually reported a successful send. A new group and an existing group can begin at different positions.

The consumer appears stuck

Inspect structured logs containing topic, partition, and offset. Measure handler duration, check for a poison message, observe rebalances and lag, and inspect downstream database or API latency. Bound retries and quarantine permanent failures before increasing timeouts or partitions.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.

Duplicate side effects occur

This usually means the handler completed but the process crashed before committing, a rebalance interrupted processing, or a retry followed an uncertain result. Use a durable event ID, database uniqueness constraints, conditional updates, and idempotent downstream requests. Kafka producer idempotence alone does not make consumer side effects exactly once.

Ordering is unexpected

Confirm that related records use the same key and therefore normally reach the same partition. Different partitions are processed concurrently, and asynchronous application work can complete out of order even when Kafka delivers records in partition order.

Production checklist

  • Topic partition count, retention, replication, and access policy are intentional.
  • Broker credentials are outside source control and TLS/SASL connectivity is tested.
  • The consumer group ID is stable and meaningful.
  • The business handler is idempotent.
  • Offset advancement follows successful processing.
  • Transport and business retries are separated and bounded.
  • A dead-letter or quarantine path preserves failed records and context.
  • Consumer lag, rebalances, errors, throughput, and handler duration are monitored.
  • CPU-heavy work and unbounded concurrency are avoided in the Node.js event loop.
  • Producers and consumers are reused rather than created per request.
  • SIGTERM and SIGINT shutdown behavior is tested.
  • Event names, IDs, schema versions, and compatibility rules are documented.
  • Node.js, client, OS, and container versions are supported by the chosen package.

The minimal producer and consumer are easy to write. The production-ready integration comes from treating partitions, offsets, retries, idempotency, schemas, security, and shutdown as part of the design—not as optional additions after the first message appears.

Product prices and availability are accurate as of the date/time indicated and are subject to change. Any price and availability information displayed on Amazon at the time of purchase will apply.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.