【发布时间】:2021-07-18 11:32:26
【问题描述】:
我创建了这个类,以使代码在发送/生成 Kafka 消息时更加可重用和清洁。
我正在使用 Node 和使用 KafkaJS 的 Kafka。
我对此很陌生,我在互联网上的任何地方都找不到完美/好的方法来在生产应用程序中使用它。
QUESTION(简而言之):我们是否需要在每次生成新消息时都以生产者身份连接。我们不能以 Redis 或 NATS 的身份保持连接。
这是我迄今为止尝试过的:
用例
假设每次创建新用户时我都需要发送一条消息。
1。创建 afka 客户端
创建了 kafka 客户端,这样我们就不必每次都重新配置它了
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。创建抽象发布者类
为所有类型的生产者创建了kafka生产者抽象类
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。最后是用户创建的生产者(继承自基类)
所有生产者都有自己的这些类
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 用户创建的事件/消息
在nodejs路由中使用
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();
});
【问题讨论】:
-
您不应该按记录连接/断开连接。您是否从某个地方复制了此模式?如果在导出之前连接生产者会发生什么?
标签: javascript node.js typescript apache-kafka microservices