Windows Errors? Fix Them Before They Spread
Repair common Windows errors and clear accumulated junk for a smoother, more stable PC - no reinstall needed.Free scan · no reinstallOutdated Drivers Are Slowing You Down
One free scan finds every outdated or missing driver and matches the right update for your exact hardware.Free scan · exact hardware matchSome 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:
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.
#1 Best Overall
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.
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.
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.
Do these 3 things before closing this tab:
1Fix the driver behind crashes, sound loss and screen glitches2Clear out junk files and repair common Windows errors3Scan for outdated or missing drivers - takes under a minuteBefore 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.
Rank #3
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.
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.
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:
- Transport retries: temporary broker or network failures handled by the client.
- Business retries: database deadlocks, rate limits, or temporary downstream outages.
- Permanent failures: malformed JSON, invalid schema, or impossible business state.
- 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.
Rank #4
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.
The Tool Desk
Outbyte Driver Updater FREEScan for outdated or missing drivers - takes under a minuteDriver Scan →Outbyte PC Repair FREERepair Windows errors before they cause bigger problemsFix Now →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.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.
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.
Best Value
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.
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.
Recommended Free Tools
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.
SIGTERMandSIGINTshutdown 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.
Quick Recap
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.
Quick wins for a faster PC:
Scan for outdated or missing drivers - takes under a minuteDriver Scan →Clear out junk files and repair common Windows errorsFree Scan →

