包含 classes/exams/homework/grades/attendance/scheduling 域、outbox、iam-consumer、redis 配置等完整实现
54 lines
1.3 KiB
TypeScript
54 lines
1.3 KiB
TypeScript
import { Kafka } from "kafkajs";
|
|
import { env } from "./env.js";
|
|
import { logger } from "../shared/observability/logger.js";
|
|
|
|
export const kafka = new Kafka({
|
|
brokers: env.KAFKA_BROKERS.split(","),
|
|
clientId: "core-edu-service",
|
|
});
|
|
|
|
export const producer = kafka.producer({
|
|
idempotent: true,
|
|
transactionalId: "core-edu-tx",
|
|
});
|
|
|
|
export const consumer = kafka.consumer({ groupId: "core-edu-group" });
|
|
|
|
let producerConnected = false;
|
|
let consumerConnected = false;
|
|
|
|
producer.on("producer.connect", () => {
|
|
producerConnected = true;
|
|
});
|
|
producer.on("producer.disconnect", () => {
|
|
producerConnected = false;
|
|
});
|
|
consumer.on("consumer.connect", () => {
|
|
consumerConnected = true;
|
|
});
|
|
consumer.on("consumer.disconnect", () => {
|
|
consumerConnected = false;
|
|
});
|
|
|
|
export function isKafkaConnected(): boolean {
|
|
return producerConnected && consumerConnected;
|
|
}
|
|
|
|
export async function connectKafka(): Promise<void> {
|
|
try {
|
|
await producer.connect();
|
|
await consumer.connect();
|
|
logger.info("Kafka connected");
|
|
} catch (err) {
|
|
logger.warn(
|
|
{ err: err instanceof Error ? err.message : String(err) },
|
|
"Kafka connect failed, running without Kafka",
|
|
);
|
|
}
|
|
}
|
|
|
|
export async function disconnectKafka(): Promise<void> {
|
|
await producer.disconnect();
|
|
await consumer.disconnect();
|
|
}
|