Untitled
unknown
javascript
10 months ago
3.9 kB
22
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