Do we need to connect everytime we producer Kafka message?

Viewed 1033

I have created this class thingy to make the code more reusable and clean for sending/producing Kafka messages.

I am using Node and Kafka using KafkaJS.

I am new to this and I could not find the perfect/good way anywhere on the internet to use it in production applications.

QUESTION (in short) : Do we need to connect as producer everytime we produce a new message. Can't we just stay connected as Redis or NATS.

This is what I have tried so far:

Use case

Lets assume I need to send a message everytime a new user is created.

1. Created afka client

Created kafka client so that we can don't have to re-configure it everytime

import { Kafka } from 'kafkajs';

class KafkaClient {
  private _client: Kafka;

  get client() {
    if (!this._client) {
      throw new Error('Cannot access client before initializing it');
    } else {
      return this._client;
    }
  }
  connect(clientId: string, brokers: string[]) {
    this._client = new Kafka({
      clientId,
      brokers,
    });
  }
}

export const producerClient = new KafkaClient();

2. Created abstract publisher class

Created kafka producer abstract class for all types of producers

import { Kafka, Message } from 'kafkajs';
import { Topics } from './topics';

interface Event {
  topic: Topics;
  data: Message;
}

export abstract class Publisher<T extends Event> {
  private client: Kafka;
  abstract topic: Topics;
  constructor(client: Kafka) {
    this.client = client;
  }

  async publish(data: T['data']): Promise<void> {
    const producer = this.client.producer();
    await producer.connect();
    await producer.send({
      topic: this.topic,
      messages: [data],
    });
  }
}

3. Finally a user-created producer (inherited from base class)

All producers will have these classes of their own

import { Publisher } from './publisher';
import { Topics } from './topics';

interface Event {
  topic: Topics.USER_CREATED;
  data: {
    value: string;
  };
}

export class UserCreatedPublisher extends Publisher<Event> {
  topic: Topics = Topics.USER_CREATED;
}

PRODUCING user-created event/message

Using it in the nodejs route

import { Router } from 'express';
import { producerClient } from './kafka-client';
import { UserCreatedPublisher } from './user-created-publisher';

const router = Router();

router.post('/create-user', async (req, res) => {
  const { email, password } = req.body;

  // send this to the other service using kafka
    await new UserCreatedPublisher(producerClient.client).publish({
      value: JSON.stringify({ email, password }),
    console.log('Message published');

  res.status(201).send();
});
1 Answers

Q: Do we need to connect as producer everytime we produce a new message? My Answer: No

What I would do in this situation is creating a singleton class to share the same instance of kafka producer and disconnect it when I am gracefully shutting down my app.

For example create a class for kafka producer

import { Kafka } from 'kafkajs';

export class KafkaProducer {
  private static instance: KafkaProducer;
  private _producer = null;
  private _isConnected = false;

  private constructor() {
    const kafka = new Kafka({
      clientId: process.env.KAFKA_CLIENTID,
      brokers: [process.env.KAFKA_SERVER],
    });
    this._producer = kafka.producer();
  }

  public static getInstance(): KafkaProducer {
    if (!KafkaProducer.instance) {
      KafkaProducer.instance = new KafkaProducer();
    }
    return KafkaProducer.instance;
  }

  public get isConnected() {
    return this._isConnected;
  }

  async connect(): Promise<void> {
    try {
      await this._producer.connect();
      this._isConnected = true;
    } catch (err) {
      console.error(err);
    }
  }

  get producer() {
    return this._producer;
  }
}

And use this producer to send messages in other places of code ( I created isConnected() method because I couldn't find build-in function for that to check if client is connected)

     let kafka = KafkaProducer.getInstance();
     if (!kafka.isConnected) {
       await kafka.connect();
     }
     await kafka.producer.send({
       topic: 'topicName',
       messages: [
         {
           value: JSON.stringify(valueObj),
         },
       ],
     });

And disconnecting before app shuts down

try {
      let kafka = KafkaProducer.getInstance();
      if (kafka.isConnected) {
        let producer = kafka.producer;
        await producer.disconnect();
      }
    } catch (err) {
      console.error(err);
    }
Related