Untitled

 avatar
unknown
javascript
a year ago
3.9 kB
23
Indexable
const { Kafka, logLevel } = require('kafkajs');

class KafkaConsumer {
    isConnected = false;
    constructor(name, topic, consumerConfig = {}) {
        this.name = name;
        this.topic = topic;
        this.consumerConfig = consumerConfig;
        this.setConsumner(consumerConfig);
    }

    async setConsumner(consumerConfig = {}) {
        this.consumer = KafkaClient.instance.broker.consumer(consumerConfig);
    }

    async connect() {
        try {
            await this.consumer.connect();
            this.isConnected = true;
            console.log(`Consumer ${this.name} connected to Kafka successfully, start cosuming...`);
            this.run();
        } catch (error) {
            console.error(`Error connecting consumer ${this.name} to Kafka:`, error);
        }
    }

    start() {
        KafkaClient.addConsumer(this.name, this);
    }

    async run() {
        throw 'Not Implemented';
    }
}

class KafkaClient {
    static topics = {
        SYNC_TCHC_CAN_BO: {
            topic: 'SYNC_TCHC_CAN_BO',
            numPartitions: 3,
            replicationFactor: 1,
        }
    }
    static producer = null;
    static _instance;
    static get instance() {
        if (!this._instance) {
            this._instance = new KafkaClient();
        }
        return this._instance;
    }


    static get isConnected() {
        return this.instance.isConnected;
    }

    static addConsumer(name, consumer = new KafkaConsumer()) {
        const instance = this.instance;
        instance._addConsumer(name, consumer);
    }

    static async send(topic, messages = []) {
        if (!this.isConnected) throw 'Kafka is not connected';
        if (!Array.isArray(messages)) messages = [messages];
        return this.instance.producer.send({ topic, messages });
    }


    isConnected = false;

    constructor() {
        this.broker = new Kafka({
            clientId: process.env.kafkaClientId || 'hcmussh',
            brokers: process.env.kafkaBrokers ? JSON.parse(process.env.kafkaBrokers) : ['localhost:9092'],
            logLevel: logLevel.ERROR
        });
        this.producer = this.broker.producer();
        this.consumner = {};
    }

    async initTopics() {
        const admin = this.broker.admin();
        await admin.connect();
        const topics = await admin.listTopics();
        await admin.createTopics({
            topics: Object.values(KafkaClient.topics).filter(t => !topics.includes(t.topic)),
            waitForLeaders: true
        });
        await admin.disconnect();
    }


    async connectConsumer(consumer) {
        try {
            if (!this.isConnected) {
                return setTimeout(() => this.connectConsumer(consumer), 1000);
            }
            await consumer.connect();
            consumer.run();
        } catch (error) {
            console.error(` - Error connecting consumer ${consumer.name} to Kafka:`, error, 'retrying in 5s');
            setTimeout(() => this.connectConsumer(consumer), 5000);
        }
    }

    _addConsumer(name, consumer) {
        if (!this.consumner[name]) {
            this.consumner[name] = consumer;
            this.connectConsumer(consumer);
        }
        else
            throw new Error(` - Consumer with name ${name} already exists`);
    }

    async connect() {
        try {
            //TODO: reconsider this
            await Promise.all([
                this.producer.connect(),
            ]);
            this.isConnected = true;
            console.info(' - Connected to Kafka successfully');
        } catch (error) {
            console.error(' - Error connecting to Kafka:', error, 'retrying in 5s');
            setTimeout(() => this.connect(), 5000);
        }
    }
}

const instance = KafkaClient.instance;
instance.connect();

module.exports = { KafkaClient, instance, KafkaConsumer };
Editor is loading...
Leave a Comment