npm package discovery and stats viewer.

Discover Tips

  • General search

    [free text search, go nuts!]

  • Package details

    pkg:[package-name]

  • User packages

    @[username]

Sponsor

Optimize Toolset

I’ve always been into building performant and accessible sites, but lately I’ve been taking it extremely seriously. So much so that I’ve been building a tool to help me optimize and monitor the sites that I build to make sure that I’m making an attempt to offer the best experience to those who visit them. If you’re into performant, accessible and SEO friendly sites, you might like it too! You can check it out at Optimize Toolset.

About

Hi, 👋, I’m Ryan Hefner  and I built this site for me, and you! The goal of this site was to provide an easy way for me to check the stats on my npm packages, both for prioritizing issues and updates, and to give me a little kick in the pants to keep up on stuff.

As I was building it, I realized that I was actually using the tool to build the tool, and figured I might as well put this out there and hopefully others will find it to be a fast and useful way to search and browse npm packages as I have.

If you’re interested in other things I’m working on, follow me on Twitter or check out the open source projects I’ve been publishing on GitHub.

I am also working on a Twitter bot for this site to tweet the most popular, newest, random packages from npm. Please follow that account now and it will start sending out packages soon–ish.

Open Software & Tools

This site wouldn’t be possible without the immense generosity and tireless efforts from the people who make contributions to the world and share their work via open source initiatives. Thank you 🙏

© 2026 – Pkg Stats / Ryan Hefner

sentinel-kafka-manager

v1.0.5

Published

Reusable Kafka manager for Node.js microservices

Readme

sentinel-kafka-manager

Reusable Kafka manager for Node.js microservices built on top of kafkajs.

This package gives you one shared way to:

  • connect to Kafka
  • reuse a single producer safely
  • reuse consumers by groupId
  • provision topics if they do not exist
  • publish JSON messages
  • build broker lists from environment variables
  • attach logging and metrics hooks
  • publish standard event envelopes
  • run consumers with DLQ support
  • inspect connection health

What This Package Solves

In most services, Kafka setup gets repeated:

  • build broker arrays
  • create producers
  • create consumers
  • manage connect and disconnect
  • handle Docker-internal vs host-external broker addresses

sentinel-kafka-manager wraps those common tasks in one small class: KafkaManager.

Enterprise Features

This version now includes a stronger production baseline:

  • config validation for client id, brokers, and broker discovery inputs
  • shared producer, admin, and consumer lifecycle
  • logger and metrics hooks
  • standard message envelope publishing
  • managed consumer runner with dead-letter topic support
  • health snapshot reporting
  • safer startup behavior for concurrent connections
  • stable error codes and structured exception details

Recommended Usage Level

This package is now suitable as a shared internal Kafka module for SaaS microservices.

Recommended use:

  • internal platform package for Node.js services
  • event publishing with standard envelopes
  • managed consumers with DLQ handling
  • consistent service-layer success and error responses

Still recommended outside the package:

  • automated unit and integration tests in the consuming service or platform repo
  • schema validation for business payloads when required by your organization
  • organization-specific tracing, alerting, and compliance rules

Installation

Install from npm:

npm install sentinel-kafka-manager

For local package development:

npm install
npm run build

Requirements

  • Node.js service
  • Kafka cluster reachable from the service
  • kafkajs is included automatically as a dependency of this package

Full Service Example

import { KafkaManager } from 'sentinel-kafka-manager'

const kafkaManager = KafkaManager.fromEnv({
  clientId: 'claims-api',
  mode: 'external',
})

async function bootstrap(): Promise<void> {
  await kafkaManager.provisionTopics([
    {
      topic: 'claims.created',
      numPartitions: 6,
      replicationFactor: 3,
    },
  ])

  await kafkaManager.publish({
    topic: 'claims.created',
    key: 'claim-1001',
    payload: {
      claimId: 'claim-1001',
      source: 'claims-api',
    },
  })

  await kafkaManager.runConsumer({
    groupId: 'claims-api-group',
    topic: 'claims.created',
    dlqTopic: 'claims.created.dlq',
    onMessage: async ({ message, headers }) => {
      console.log('Received event:', {
        traceId: headers['x-trace-id'],
        payload: message,
      })
    },
  })
}

bootstrap().catch(async (error) => {
  console.error('Kafka bootstrap failed:', error)
  await kafkaManager.disconnect()
  process.exit(1)
})

process.on('SIGINT', async () => {
  await kafkaManager.disconnect()
  process.exit(0)
})

process.on('SIGTERM', async () => {
  await kafkaManager.disconnect()
  process.exit(0)
})

Basic Usage

1. Import the package

import { KafkaManager } from 'sentinel-kafka-manager'

If you want to catch package-specific errors explicitly:

import {
  KafkaErrorCode,
  KafkaManager,
  KafkaManagerError,
  OperationResponse,
} from 'sentinel-kafka-manager'

2. Create a manager

You can create it in two ways:

  1. pass brokers directly
  2. build brokers from environment variables

Direct broker example

const kafkaManager = new KafkaManager({
  clientId: 'claims-service',
  brokers: ['localhost:29092', 'localhost:39092', 'localhost:49092'],
})

Environment-based example

const kafkaManager = KafkaManager.fromEnv({
  clientId: 'claims-service',
  mode: 'external',
})

Recommended Project Pattern

In a real service, create one shared manager instance and reuse it.

Example:

import { KafkaManager } from 'sentinel-kafka-manager'

export const kafkaManager = KafkaManager.fromEnv({
  clientId: 'claims-service',
  mode: 'external',
})

Then import that shared instance anywhere you need Kafka access.

You can also attach enterprise hooks:

import { KafkaManager } from 'sentinel-kafka-manager'

export const kafkaManager = KafkaManager.fromEnv({
  clientId: 'claims-service',
  mode: 'external',
  logger: {
    debug: (message, meta) => console.debug(message, meta),
    info: (message, meta) => console.info(message, meta),
    warn: (message, meta) => console.warn(message, meta),
    error: (message, meta) => console.error(message, meta),
  },
  metrics: {
    emit: (event) => {
      console.log('METRIC', event)
    },
  },
})

How To Use In Your Service

Publish a message

await kafkaManager.publish({
  topic: 'claims.created',
  key: 'claim-1001',
  payload: {
    claimId: 'claim-1001',
    status: 'created',
    source: 'claims-service',
  },
})

publish() automatically:

  • gets a shared producer
  • connects it once
  • serializes payload with JSON.stringify()

Publish an enterprise event envelope

Use this when you want traceable event metadata such as eventType, version, traceId, or tenantId.

await kafkaManager.publishEnvelope({
  topic: 'claims.created',
  key: 'claim-1001',
  envelope: {
    eventId: 'evt-claim-1001',
    eventType: 'claims.created',
    version: '1.0.0',
    timestamp: new Date().toISOString(),
    source: 'claims-service',
    traceId: 'trace-123',
    tenantId: 'tenant-abc',
    correlationId: 'corr-001',
    payload: {
      claimId: 'claim-1001',
      status: 'created',
    },
  },
})

Create or reuse a producer manually

If you need direct producer access:

const producer = await kafkaManager.getProducer()

await producer.send({
  topic: 'claims.created',
  messages: [
    {
      key: 'claim-1001',
      value: JSON.stringify({ claimId: 'claim-1001' }),
    },
  ],
})

Create or reuse a consumer

const consumer = await kafkaManager.getConsumer('claims-service-group')

await consumer.subscribe({
  topic: 'claims.created',
  fromBeginning: false,
})

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

    console.log({
      topic,
      partition,
      key: message.key?.toString(),
      payload,
    })
  },
})

getConsumer(groupId) returns one shared connected consumer per group id.

Run a managed consumer with DLQ support

For SaaS-style services, this is the recommended consumer pattern.

await kafkaManager.runConsumer({
  groupId: 'claims-service-group',
  topic: 'claims.created',
  dlqTopic: 'claims.created.dlq',
  onMessage: async ({ message, headers, key }) => {
    console.log('Processing:', {
      key,
      traceId: headers['x-trace-id'],
      payload: message,
    })
  },
  onError: async (error, context) => {
    console.error('Consumer error:', {
      error: error.message,
      topic: context.topic,
      offset: context.offset,
    })
  },
})

If processing fails:

  • the error is logged
  • a metric hook can receive the failure event
  • the message can be published to a dead-letter topic when dlqTopic is set

Catch structured errors

All package-thrown operational errors now use KafkaManagerError.

import {
  KafkaErrorCode,
  KafkaManagerError,
} from 'sentinel-kafka-manager'

try {
  await kafkaManager.publish({
    topic: 'claims.created',
    payload: { claimId: 'claim-1001' },
  })
} catch (error) {
  if (error instanceof KafkaManagerError) {
    console.error('Kafka error', {
      code: error.code,
      message: error.message,
      details: error.details,
    })

    if (error.code === KafkaErrorCode.MESSAGE_PUBLISH_FAILED) {
      // retry, alert, or degrade gracefully
    }
  }

  throw error
}

Use common success and error response types

If your service wraps Kafka operations in API or service-layer responses, you can use the shared response contracts from this package.

import {
  ErrorResponse,
  OperationResponse,
  SuccessResponse,
} from 'sentinel-kafka-manager'

type PublishClaimCreatedResponse = OperationResponse<{
  topic: string
  key: string
}>

Success example:

const response: SuccessResponse<{ topic: string; key: string }> = {
  success: true,
  message: 'Kafka message published successfully.',
  data: {
    topic: 'claims.created',
    key: 'claim-1001',
  },
  meta: {
    timestamp: new Date().toISOString(),
    traceId: 'trace-123',
  },
}

Error example:

const response: ErrorResponse = {
  success: false,
  message: 'Failed to publish Kafka message.',
  error: {
    code: KafkaErrorCode.MESSAGE_PUBLISH_FAILED,
    details: {
      topic: 'claims.created',
      key: 'claim-1001',
    },
  },
  meta: {
    timestamp: new Date().toISOString(),
    retryable: true,
  },
}

Provision topics

await kafkaManager.provisionTopics([
  {
    topic: 'claims.created',
    numPartitions: 6,
    replicationFactor: 3,
  },
  {
    topic: 'claims.updated',
    numPartitions: 6,
    replicationFactor: 3,
  },
])

This is useful during service startup when your topics should exist before producers or consumers begin work.

Disconnect on shutdown

Always close Kafka clients during application shutdown.

process.on('SIGINT', async () => {
  await kafkaManager.disconnect()
  process.exit(0)
})

process.on('SIGTERM', async () => {
  await kafkaManager.disconnect()
  process.exit(0)
})

Inspect health

const health = kafkaManager.getHealthSnapshot()

console.log(health)

Environment Setup

This package supports two environment styles:

  1. one explicit broker list with KAFKA_BROKERS
  2. derived brokers from your existing Kafka runtime variables

Option 1: Use KAFKA_BROKERS

This is the simplest option.

.env

KAFKA_BROKERS=localhost:29092,localhost:39092,localhost:49092

Code

const kafkaManager = KafkaManager.fromEnv({
  clientId: 'claims-service',
})

Option 2: Use your current Kafka cluster variables

This matches the Kafka server setup you shared.

For host machine apps

Use this when your Node.js service runs on the host machine.

.env

KAFKA_EXTERNAL_HOST=localhost
KAFKA_BROKER_1_EXTERNAL_PORT=29092
KAFKA_BROKER_2_EXTERNAL_PORT=39092
KAFKA_BROKER_3_EXTERNAL_PORT=49092
KAFKA_BROKER_INTERNAL_PORT=9092

Code

const kafkaManager = KafkaManager.fromEnv({
  clientId: 'claims-service',
  mode: 'external',
})

This resolves to:

localhost:29092
localhost:39092
localhost:49092

For Dockerized apps on the Kafka network

Use this when your Node.js service runs inside Docker on the same Kafka network.

.env

KAFKA_BROKER_INTERNAL_PORT=9092

Code

const kafkaManager = KafkaManager.fromEnv({
  clientId: 'claims-service',
  mode: 'internal',
})

This resolves to:

broker-1:9092
broker-2:9092
broker-3:9092

Which Mode Should You Use?

  • Use mode: 'external' when your service uses localhost broker ports like 29092, 39092, and 49092
  • Use mode: 'internal' when your service is inside Docker and should reach brokers like broker-1:9092
  • Use KAFKA_BROKERS when you want the simplest and most explicit configuration

Your Current Running Kafka Cluster

Based on your running Docker containers, your Kafka cluster is exposed like this:

  • broker-1 -> host port 29092
  • broker-2 -> host port 39092
  • broker-3 -> host port 49092
  • controller-1, controller-2, and controller-3 are controller-only nodes and should not be used by application clients

If your Node.js service runs on the host machine, use:

KAFKA_EXTERNAL_HOST=localhost
KAFKA_BROKER_1_EXTERNAL_PORT=29092
KAFKA_BROKER_2_EXTERNAL_PORT=39092
KAFKA_BROKER_3_EXTERNAL_PORT=49092
KAFKA_BROKER_INTERNAL_PORT=9092
const kafkaManager = KafkaManager.fromEnv({
  clientId: 'claims-service',
  mode: 'external',
})

If your Node.js service runs inside Docker on the same Kafka network, use:

const kafkaManager = KafkaManager.fromEnv({
  clientId: 'claims-service',
  mode: 'internal',
})

This package is designed to work with that exact broker layout.

Local Infrastructure And Operations

Create A Kafka Server With This Repo

This repository already includes a full local Kafka server setup in docker-compose.yml and sample environment values in .env.local.example.

The stack includes:

  • 3 KRaft controllers
  • 3 Kafka brokers
  • 1 Kafka UI instance

1. Create your .env

Copy the example file:

cp .env.local.example .env

If you want a standard localhost setup, the example values are already suitable.

2. Start the Kafka cluster

docker compose --env-file .env up -d

This creates:

  • controller quorum on port 9093 inside Docker
  • internal broker traffic on port 9092 inside Docker
  • host-accessible broker ports 29092, 39092, and 49092
  • Kafka UI on 127.0.0.1:5001

3. Stop the Kafka cluster

docker compose --env-file .env down

To also remove persisted Kafka data volumes:

docker compose --env-file .env down -v

4. Check container status

docker compose --env-file .env ps

5. Current broker endpoints

For applications running on your host machine:

localhost:29092
localhost:39092
localhost:49092

For applications running inside Docker on the same Compose network:

broker-1:9092
broker-2:9092
broker-3:9092

Important:

  • application clients must connect to brokers, not controllers
  • controllers are internal cluster metadata nodes only
  • mode: 'external' is for host-based apps
  • mode: 'internal' is for Dockerized apps on the same network

Kafka UI Usage

The Docker stack also starts Kafka UI for cluster inspection.

Open the UI

Visit:

http://127.0.0.1:5001

Login credentials

By default, the UI uses the credentials from .env:

KAFKA_UI_USERNAME=admin
KAFKA_UI_PASSWORD=Admin@007

What the UI connects to

Kafka UI is preconfigured in docker-compose.yml to use all three brokers:

broker-1:9092,broker-2:9092,broker-3:9092

The cluster name shown in the UI comes from:

KAFKA_UI_CLUSTER_NAME=Project Omni Enterprise Kafka

What you can do in the UI

  • view brokers and cluster health
  • inspect topics and partitions
  • browse messages
  • inspect consumer groups
  • check offsets and lag

Read-only mode

The example .env sets:

KAFKA_UI_READONLY=true

That means the UI is intended for safe inspection only.

If you want to create topics or make changes from the UI, change it to:

KAFKA_UI_READONLY=false

Then restart the stack:

docker compose --env-file .env up -d

Typical UI flow

  1. Start the Docker stack.
  2. Open http://127.0.0.1:5001.
  3. Log in with KAFKA_UI_USERNAME and KAFKA_UI_PASSWORD.
  4. Open the configured cluster.
  5. Go to Topics, Consumer Groups, or Brokers as needed.

Docker Compose And .env Reference

This project does not use a separate Dockerfile for Kafka. The Kafka server is created from docker-compose.yml and environment values from .env.

Core Kafka image and cluster identity

KAFKA_IMAGE=confluentinc/cp-kafka:7.6.6
KAFKA_CLUSTER_ID=MkU3OEVBNTcwNTJENDM2Qk
KAFKA_RESTART_POLICY=unless-stopped
KAFKA_STOP_GRACE_PERIOD=60s
KAFKA_NETWORK_NAME=kafka-network
  • KAFKA_IMAGE: Kafka container image used by controllers and brokers
  • KAFKA_CLUSTER_ID: shared KRaft cluster id for all nodes
  • KAFKA_NETWORK_NAME: Docker network name used by the full stack

Controller quorum settings

KAFKA_CONTROLLER_PORT=9093
KAFKA_CONTROLLER_LISTENER_NAMES=CONTROLLER
KAFKA_CONTROLLER_QUORUM_VOTERS=1@controller-1:9093,2@controller-2:9093,3@controller-3:9093
  • controllers run only inside Docker
  • apps should never use these controller addresses as Kafka client brokers

Broker listener settings

KAFKA_BROKER_BIND_ADDRESS=0.0.0.0
KAFKA_EXTERNAL_HOST=localhost
KAFKA_BROKER_INTERNAL_PORT=9092
KAFKA_BROKER_EXTERNAL_CONTAINER_PORT=19092
KAFKA_BROKER_1_EXTERNAL_PORT=29092
KAFKA_BROKER_2_EXTERNAL_PORT=39092
KAFKA_BROKER_3_EXTERNAL_PORT=49092
KAFKA_INTER_BROKER_LISTENER_NAME=INTERNAL
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP=CONTROLLER:PLAINTEXT,INTERNAL:PLAINTEXT,EXTERNAL:PLAINTEXT
  • KAFKA_EXTERNAL_HOST=localhost exposes brokers to your host machine
  • KAFKA_BROKER_1_EXTERNAL_PORT, KAFKA_BROKER_2_EXTERNAL_PORT, and KAFKA_BROKER_3_EXTERNAL_PORT are the ports your local apps should use
  • KAFKA_BROKER_INTERNAL_PORT=9092 is used by containers on the Docker network

Replication and durability defaults

KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR=3
KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR=3
KAFKA_TRANSACTION_STATE_LOG_MIN_ISR=2
KAFKA_DEFAULT_REPLICATION_FACTOR=3
KAFKA_MIN_INSYNC_REPLICAS=2
KAFKA_NUM_PARTITIONS=6

These values are designed for the included three-broker cluster.

Broker behavior

KAFKA_AUTO_CREATE_TOPICS_ENABLE=false
KAFKA_DELETE_TOPIC_ENABLE=true
KAFKA_LOG_RETENTION_HOURS=168
KAFKA_LOG_SEGMENT_BYTES=1073741824
KAFKA_MESSAGE_MAX_BYTES=10485880
KAFKA_REPLICA_FETCH_MAX_BYTES=10485880
KAFKA_SOCKET_REQUEST_MAX_BYTES=104857600
  • topics are not auto-created by default
  • deleting topics is allowed
  • retention is 7 days by default

Storage and safety limits

KAFKA_LOG_DIRS=/var/lib/kafka/data
KAFKA_ULIMIT_NOFILE_SOFT=65536
KAFKA_ULIMIT_NOFILE_HARD=65536
KAFKA_LOG_MAX_SIZE=100m
KAFKA_LOG_MAX_FILE=5

Kafka UI settings

KAFKA_UI_IMAGE=provectuslabs/kafka-ui:latest
KAFKA_UI_CONTAINER_NAME=kafka-ui
KAFKA_UI_HOSTNAME=kafka-ui
KAFKA_UI_RESTART_POLICY=unless-stopped
KAFKA_UI_BIND_ADDRESS=127.0.0.1
KAFKA_UI_PORT=5001
KAFKA_UI_DYNAMIC_CONFIG_ENABLED=false
KAFKA_UI_CLUSTER_NAME=Project Omni Enterprise Kafka
KAFKA_UI_READONLY=true
KAFKA_UI_AUTH_TYPE=LOGIN_FORM
KAFKA_UI_USERNAME=admin
KAFKA_UI_PASSWORD=Admin@007
  • KAFKA_UI_BIND_ADDRESS=127.0.0.1 keeps the UI local-only
  • KAFKA_UI_PORT=5001 publishes the UI on your machine
  • KAFKA_UI_READONLY=true prevents write actions from the UI

Healthcheck settings

KAFKA_HEALTHCHECK_INTERVAL=10s
KAFKA_HEALTHCHECK_TIMEOUT=5s
KAFKA_HEALTHCHECK_RETRIES=15
KAFKA_HEALTHCHECK_START_PERIOD=30s

Creating Topics

Because KAFKA_AUTO_CREATE_TOPICS_ENABLE=false, topics should be created deliberately.

You have two good options:

  • create topics from your service with kafkaManager.provisionTopics()
  • allow UI-based creation by setting KAFKA_UI_READONLY=false

Example with this package:

await kafkaManager.provisionTopics([
  {
    topic: 'claims.created',
    numPartitions: 6,
    replicationFactor: 3,
  },
])

API Summary

new KafkaManager(options)

Use this when you already know your brokers.

const kafkaManager = new KafkaManager({
  clientId: 'claims-service',
  brokers: ['localhost:29092'],
})

Options:

  • clientId: string
  • brokers: string[]
  • ssl?: boolean | { rejectUnauthorized?: boolean }
  • sasl?: SASLOptions
  • connectionTimeout?: number
  • requestTimeout?: number
  • retry?: { initialRetryTime?: number; retries?: number }
  • logger?: KafkaLogger
  • metrics?: KafkaMetrics
  • producerConfig?: ProducerConfig
  • consumerDefaults?: Omit<ConsumerConfig, 'groupId'>

KafkaManager.fromEnv(options)

Use this when you want brokers to come from environment variables.

const kafkaManager = KafkaManager.fromEnv({
  clientId: 'claims-service',
  mode: 'external',
})

Options:

  • clientId: string
  • envKey?: string
  • mode?: 'internal' | 'external'
  • brokerCount?: number
  • externalHostEnvKey?: string
  • internalPortEnvKey?: string
  • externalPortEnvKeyPrefix?: string
  • externalPortEnvKeySuffix?: string
  • internalBrokerHostPrefix?: string
  • ssl?: boolean | { rejectUnauthorized?: boolean }
  • sasl?: SASLOptions
  • connectionTimeout?: number
  • requestTimeout?: number
  • retry?: { initialRetryTime?: number; retries?: number }
  • logger?: KafkaLogger
  • metrics?: KafkaMetrics
  • producerConfig?: ProducerConfig
  • consumerDefaults?: Omit<ConsumerConfig, 'groupId'>

getProducer()

Returns a shared connected Kafka producer.

getConsumer(groupId)

Returns a shared connected Kafka consumer for that groupId.

publish(options)

Publishes one JSON message.

await kafkaManager.publish({
  topic: 'claims.created',
  key: 'claim-1001',
  payload: { claimId: 'claim-1001' },
})

publishEnvelope(options)

Publishes a standard metadata envelope for traceable event-driven systems.

provisionTopics(topics)

Creates missing topics.

runConsumer(options)

Runs a managed consumer with:

  • a typed message parser
  • application callback handling
  • optional DLQ publishing
  • optional error callback

getHealthSnapshot()

Returns a lightweight connection-health snapshot.

disconnect()

Disconnects the producer, admin client, and any connected consumers.

Error Handling

This package now throws structured errors with stable codes for production handling.

Error class

  • KafkaManagerError

Common response types

  • SuccessResponse<T>
  • ErrorResponse
  • OperationResponse<T>

Error fields

  • code: stable machine-readable error code
  • message: human-readable explanation
  • details: structured metadata for logs or monitoring
  • cause: original thrown error when available

Common error codes

  • INVALID_CLIENT_ID
  • INVALID_BROKERS
  • INVALID_GROUP_ID
  • INVALID_TOPIC
  • INVALID_ENVELOPE
  • INVALID_BROKER_COUNT
  • MISSING_ENVIRONMENT_VARIABLE
  • PRODUCER_CONNECTION_FAILED
  • CONSUMER_CONNECTION_FAILED
  • ADMIN_CONNECTION_FAILED
  • MESSAGE_PUBLISH_FAILED
  • TOPIC_PROVISION_FAILED
  • CONSUMER_RUN_FAILED
  • CONSUMER_HANDLER_FAILED
  • DLQ_PUBLISH_FAILED
  • MESSAGE_PARSE_FAILED
  • DISCONNECT_FAILED

Example

import { KafkaManagerError } from 'sentinel-kafka-manager'

try {
  await kafkaManager.runConsumer({
    groupId: 'claims-service-group',
    topic: 'claims.created',
    onMessage: async ({ message }) => {
      console.log(message)
    },
  })
} catch (error) {
  if (error instanceof KafkaManagerError) {
    console.error(error.code, error.message, error.details)
  }

  throw error
}

Advanced Examples

Logger and metrics example

const kafkaManager = KafkaManager.fromEnv({
  clientId: 'claims-service',
  mode: 'external',
  logger: {
    debug: (message, meta) => console.debug(message, meta),
    info: (message, meta) => console.info(message, meta),
    warn: (message, meta) => console.warn(message, meta),
    error: (message, meta) => console.error(message, meta),
  },
  metrics: {
    emit: (event) => {
      console.log('METRIC', event.name, event.tags, event.meta)
    },
  },
})

SSL example

const kafkaManager = new KafkaManager({
  clientId: 'claims-service',
  brokers: ['localhost:29092'],
  ssl: {
    rejectUnauthorized: true,
  },
})

Managed consumer with custom parser

await kafkaManager.runConsumer({
  groupId: 'claims-service-group',
  topic: 'claims.created',
  dlqTopic: 'claims.created.dlq',
  parser: (payload) => {
    const rawValue = payload.message.value?.toString() ?? '{}'
    return JSON.parse(rawValue) as {
      claimId: string
      status: string
    }
  },
  onMessage: async ({ message }) => {
    console.log(message.claimId, message.status)
  },
})

SASL example

const kafkaManager = new KafkaManager({
  clientId: 'claims-service',
  brokers: ['localhost:29092'],
  sasl: {
    mechanism: 'plain',
    username: 'kafka-user',
    password: 'secret',
  },
})

Retry and timeout example

const kafkaManager = new KafkaManager({
  clientId: 'claims-service',
  brokers: ['localhost:29092'],
  connectionTimeout: 5000,
  requestTimeout: 30000,
  retry: {
    initialRetryTime: 300,
    retries: 10,
  },
})

Operational Guidance

Common Mistakes

  • Do not connect application clients to KRaft controller nodes. Connect only to brokers.
  • Do not use mode: 'external' inside Docker unless you truly want host-published ports.
  • Do not create a new KafkaManager for every request. Reuse one shared instance.
  • Do not forget disconnect() during shutdown.
  • Do not assume topics will exist unless your platform creates them or your service provisions them.
  • Do not ignore failed consumer messages in production. Use runConsumer() with a dlqTopic.
  • Do not publish business-critical events without metadata. Prefer publishEnvelope() for traceability.
  • Do not depend on raw driver error strings in application logic. Use KafkaManagerError.code.

Recommended SaaS Pattern

For most SaaS services, the cleanest pattern is:

  1. Create one shared KafkaManager during service bootstrap.
  2. Use publishEnvelope() for business events.
  3. Use runConsumer() instead of raw consumer.run() for managed handling.
  4. Configure a dlqTopic for important consumers.
  5. Catch KafkaManagerError and branch on error.code.
  6. Wrap service outcomes in OperationResponse<T> where useful.

Notes

  • Your current Kafka cluster works fine with PLAINTEXT, so SSL and SASL are optional unless your environment changes
  • This package is designed for both host-based services and Dockerized services
  • If you want the least surprise, use KAFKA_BROKERS

Additional Documentation

About Me

I'm a full stack developer. Experienced developer with over 20 years of expertise in crafting scalable web applications. Proficient in frontend technologies such as Angular and React, alongside extensive experience in WordPress, Drupal, and backend frameworks like Node.js. With a proven track record of delivering high-quality solutions to meet client needs and drive business objectives, I bring a versatile skill set and a commitment to excellence to every project.

Feedback

If you have any feedback, please reach out to us at [email protected]