umami/lib/kafka.ts

119 lines
2.5 KiB
TypeScript
Raw Normal View History

2022-08-12 18:21:43 +02:00
import dateFormat from 'dateformat';
2022-08-28 06:38:35 +02:00
import debug from 'debug';
import { Kafka, Mechanism, Producer, RecordMetadata, SASLOptions, logLevel } from 'kafkajs';
2022-08-28 06:38:35 +02:00
import { KAFKA, KAFKA_PRODUCER } from 'lib/db';
import * as tls from 'tls';
2022-08-12 18:21:43 +02:00
2022-08-29 05:20:54 +02:00
const log = debug('umami:kafka');
2022-08-28 06:38:35 +02:00
let kafka: Kafka;
let producer: Producer;
2022-10-07 00:00:16 +02:00
const enabled = Boolean(process.env.KAFKA_URL && process.env.KAFKA_BROKER);
2022-08-28 06:38:35 +02:00
function getClient() {
const { username, password } = new URL(process.env.KAFKA_URL);
2022-08-12 18:21:43 +02:00
const brokers = process.env.KAFKA_BROKER.split(',');
const ssl: { ssl?: tls.ConnectionOptions | boolean; sasl?: SASLOptions | Mechanism } =
2022-08-28 06:38:35 +02:00
username && password
? {
ssl: {
checkServerIdentity: () => undefined,
2022-09-22 19:36:23 +02:00
ca: [process.env.CA_CERT],
key: process.env.CLIENT_KEY,
cert: process.env.CLIENT_CERT,
},
2022-08-28 06:38:35 +02:00
sasl: {
mechanism: 'plain',
username,
password,
},
}
: {};
2022-08-12 18:21:43 +02:00
const client: Kafka = new Kafka({
2022-08-28 06:38:35 +02:00
clientId: 'umami',
brokers: brokers,
connectionTimeout: 3000,
logLevel: logLevel.ERROR,
...ssl,
});
2022-08-25 20:07:47 +02:00
if (process.env.NODE_ENV !== 'production') {
2022-08-28 06:38:35 +02:00
global[KAFKA] = client;
2022-08-25 20:07:47 +02:00
}
2022-08-12 18:21:43 +02:00
log('Kafka initialized');
2022-08-28 06:38:35 +02:00
return client;
}
2022-08-12 18:21:43 +02:00
async function getProducer(): Promise<Producer> {
2022-08-12 18:21:43 +02:00
const producer = kafka.producer();
await producer.connect();
2022-08-28 06:38:35 +02:00
if (process.env.NODE_ENV !== 'production') {
global[KAFKA_PRODUCER] = producer;
}
log('Kafka producer initialized');
2022-08-25 20:07:47 +02:00
return producer;
}
2023-07-11 00:30:34 +02:00
function getDateFormat(date, format?): string {
return dateFormat(date, format ? format : 'UTC:yyyy-mm-dd HH:MM:ss');
2022-08-26 07:04:32 +02:00
}
async function sendMessage(
message: { [key: string]: string | number },
topic: string,
): Promise<RecordMetadata[]> {
2022-10-07 00:00:16 +02:00
await connect();
2022-09-22 19:36:23 +02:00
return producer.send({
2022-08-12 18:21:43 +02:00
topic,
messages: [
{
value: JSON.stringify(message),
2022-08-12 18:21:43 +02:00
},
],
acks: -1,
});
}
async function sendMessages(messages: { [key: string]: string | number }[], topic: string) {
await connect();
await producer.send({
topic,
messages: messages.map(a => {
return { value: JSON.stringify(a) };
}),
acks: 1,
2022-08-12 18:21:43 +02:00
});
}
2022-08-26 07:43:22 +02:00
async function connect(): Promise<Kafka> {
2022-09-22 19:36:23 +02:00
if (!kafka) {
kafka = process.env.KAFKA_URL && process.env.KAFKA_BROKER && (global[KAFKA] || getClient());
if (kafka) {
producer = global[KAFKA_PRODUCER] || (await getProducer());
}
}
return kafka;
}
2022-08-26 07:43:22 +02:00
export default {
2022-10-07 00:00:16 +02:00
enabled,
2022-08-28 06:38:35 +02:00
client: kafka,
2022-09-22 19:36:23 +02:00
producer,
2022-08-28 06:38:35 +02:00
log,
2022-10-07 00:00:16 +02:00
connect,
2022-08-26 07:43:22 +02:00
getDateFormat,
sendMessage,
sendMessages,
2022-08-26 07:43:22 +02:00
};