Files
Edu/services/msg/src/shared/kafka/kafka.consumer.ts
SpecialX 765f7da4c0 feat(msg): v2 ARB-013 topic 命名统一 + 考试事件消费
ARB-013 P0 修复:PRODUCER_TOPIC_MAP 从 edu.notification.* 改为 edu.notify.notification.*

kafka.consumer 新增 3 考试实时事件消费(exam.extended/force_submitted/question_reordered)

嵌套 payload 解包支持 + topic-map 扩展

新增 6 测试数据文件(docker-notify + 5 kafka 事件 json)

101 单元测试通过 + Docker 真实环境验证
2026-07-14 23:00:33 +08:00

636 lines
21 KiB
TypeScript
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
import {
Injectable,
type OnModuleDestroy,
type OnModuleInit,
} from "@nestjs/common";
import type { KafkaMessage } from "kafkajs";
import { connectKafka, getConsumer } from "./kafka.client.js";
import { CONSUMER_TOPICS } from "./topic-map.js";
import { checkAndMark } from "../redis/idempotency.guard.js";
import { logger } from "../observability/logger.js";
import { NotificationsService } from "../../notifications/notifications.service.js";
/**
* KafkaConsumer —— 消费 iam/core-edu/data-ana 事件,触发通知。
*
* 仲裁依据:
* - 02-architecture-design.md §5.112 类消费事件)
* - JSON payload 降级events.proto 未定义全部事件类型,直接解析 JSON
* - 幂等event_id 去重Redis SETNX + DB 降级)
* - at-least-once消费失败不 commit offsetKafka 重投
*
* 事件路由:
* - iamuser.created/updated/deleted/role_changed, role.created/updated
* - core-eduexam.published, homework.assigned, assignment.submitted/graded, grade.recorded, attendance.recorded
* - data-anamastery.updated
*/
@Injectable()
export class KafkaConsumerService implements OnModuleInit, OnModuleDestroy {
private running = false;
constructor(private readonly notificationsService: NotificationsService) {}
async onModuleInit(): Promise<void> {
await this.start();
}
async onModuleDestroy(): Promise<void> {
await this.stop();
}
async start(): Promise<void> {
if (this.running) return;
// 确保 Kafka producer+consumer 已连接幂等connectKafka 内部处理重复调用)
await connectKafka();
const consumer = getConsumer();
try {
await consumer.subscribe({
topics: [...CONSUMER_TOPICS],
fromBeginning: false,
});
this.running = true;
await consumer.run({
eachMessage: async ({ topic, partition, message }) => {
await this.handleMessage(topic, partition, message);
},
});
logger.info({ topics: CONSUMER_TOPICS }, "KafkaConsumer started");
} catch (err) {
logger.warn(
{ err },
"KafkaConsumer start failed (running without consumer)",
);
}
}
async stop(): Promise<void> {
this.running = false;
try {
const consumer = getConsumer();
await consumer.stop();
await consumer.disconnect();
logger.info("KafkaConsumer stopped");
} catch (err) {
logger.warn({ err }, "KafkaConsumer stop error");
}
}
private async handleMessage(
topic: string,
partition: number,
message: KafkaMessage,
): Promise<void> {
const eventId = this.extractEventId(message, topic, partition);
const payload = this.parsePayload(message);
if (!payload) {
logger.warn({ topic, eventId }, "Failed to parse Kafka message payload");
return;
}
// 幂等检查
const { isFirst } = await checkAndMark(eventId, topic);
if (!isFirst) {
logger.debug(
{ topic, eventId },
"Event already processed (idempotent skip)",
);
return;
}
// 路由到处理器
try {
await this.routeEvent(topic, eventId, payload);
logger.info({ topic, eventId }, "Kafka event processed");
} catch (err) {
logger.error({ topic, eventId, err }, "Failed to process Kafka event");
// 不抛出at-least-once 语义下,失败的 event 已标记 processed
// 后续靠人工或监控重处理(避免无限重试阻塞消费)
}
}
private extractEventId(
message: KafkaMessage,
topic: string,
partition: number,
): string {
// 优先从 headers 取 eventId
const headerValue = message.headers?.eventId;
if (headerValue) {
if (typeof headerValue === "string") {
return headerValue;
}
// KafkaJS header 可能是 Buffer 或 (string|Buffer)[],取首元素
const buf = Array.isArray(headerValue) ? headerValue[0] : headerValue;
if (buf) {
return Buffer.from(buf).toString("utf-8");
}
}
// 降级topic + partition + offset 组合
return `${topic}:${partition}:${message.offset}`;
}
/**
* JSON payload 降级events.proto 未定义全部事件类型,
* 直接解析 JSON。若解析失败返回 null。
*/
private parsePayload(message: KafkaMessage): Record<string, unknown> | null {
try {
const value = message.value;
if (!value) return null;
const str =
typeof value === "string"
? value
: Buffer.from(value).toString("utf-8");
return JSON.parse(str) as Record<string, unknown>;
} catch {
return null;
}
}
/**
* 根据 topic 路由到具体的事件处理器。
*
* payload 兼容两种格式:
* - 扁平格式(测试用):消息 value 直接是业务字段 { examId, studentIds, ... }
* - 嵌套格式core-edu Outbox 标准):{ event_id, event_type, payload: { examId, ... } }
* routeEvent 内部统一解包到 businessPayloadhandler 只处理业务字段。
*/
private async routeEvent(
topic: string,
eventId: string,
payload: Record<string, unknown>,
): Promise<void> {
// 解包嵌套 payloadcore-edu Outbox 发送 { event_id, event_type, payload: {...} }
const inner = payload.payload;
const businessPayload =
inner && typeof inner === "object" && !Array.isArray(inner)
? (inner as Record<string, unknown>)
: payload;
switch (topic) {
// iam 事件
case "edu.identity.user.created":
await this.handleUserCreated(eventId, businessPayload);
break;
case "edu.identity.user.role_changed":
await this.handleRoleChanged(eventId, businessPayload);
break;
case "edu.identity.role.updated":
await this.handleRoleUpdated(eventId, businessPayload);
break;
case "edu.identity.user.updated":
case "edu.identity.user.deleted":
case "edu.identity.role.created":
// 无需通知,仅幂等标记
logger.debug(
{ topic, eventId },
"Event acknowledged (no notification)",
);
break;
// core-edu 事件6 基础)
case "edu.teaching.exam.published":
await this.handleExamPublished(eventId, businessPayload);
break;
case "edu.teaching.homework.assigned":
await this.handleHomeworkAssigned(eventId, businessPayload);
break;
case "edu.teaching.assignment.submitted":
await this.handleAssignmentSubmitted(eventId, businessPayload);
break;
case "edu.teaching.assignment.graded":
await this.handleAssignmentGraded(eventId, businessPayload);
break;
case "edu.teaching.grade.recorded":
await this.handleGradeRecorded(eventId, businessPayload);
break;
case "edu.teaching.attendance.recorded":
await this.handleAttendanceRecorded(eventId, businessPayload);
break;
// core-edu 考试实时事件P3.14 新增3 个)— 转发到 push-gateway
case "edu.teaching.exam.extended":
await this.handleExamExtended(eventId, businessPayload, payload);
break;
case "edu.teaching.exam.force_submitted":
await this.handleExamForceSubmitted(eventId, businessPayload, payload);
break;
case "edu.teaching.exam.question_reordered":
await this.handleExamQuestionReordered(
eventId,
businessPayload,
payload,
);
break;
// data-ana 事件
case "edu.insight.mastery.updated":
await this.handleMasteryUpdated(eventId, businessPayload);
break;
default:
logger.warn({ topic, eventId }, "Unknown topic, skipping");
}
}
// ============================================================
// 考试实时事件处理器P3.14:转发 core-edu → push-gateway → WebSocket
// ============================================================
/**
* 考试时间延长通知。
* payload: { examId, classId, subjectId, extensionSeconds, newDuration, studentIds? }
*/
private async handleExamExtended(
eventId: string,
payload: Record<string, unknown>,
rawPayload: Record<string, unknown>,
): Promise<void> {
const examId = String(payload.examId ?? payload.exam_id ?? "");
const className = String(payload.className ?? payload.class_name ?? "");
const extensionSeconds = Number(
payload.extensionSeconds ?? payload.extension_seconds ?? 0,
);
const newDuration = Number(
payload.newDuration ?? payload.new_duration ?? 0,
);
const studentIds = payload.studentIds ?? payload.student_ids;
const extensionMin = Math.round(extensionSeconds / 60);
const newDurationMin = Math.round(newDuration / 60);
const content = `你的考试时间已延长 ${extensionMin} 分钟,新时长:${newDurationMin} 分钟${className ? `${className}` : ""}`;
await this.broadcastExamRealtimeEvent(
eventId,
"exam",
"考试时间延长通知",
content,
"exam_extended",
examId,
studentIds,
rawPayload,
);
}
/**
* 考试强制提交通知。
* payload: { examId, classId, studentIds? }
*/
private async handleExamForceSubmitted(
eventId: string,
payload: Record<string, unknown>,
rawPayload: Record<string, unknown>,
): Promise<void> {
const examId = String(payload.examId ?? payload.exam_id ?? "");
const examTitle = String(payload.examTitle ?? payload.exam_title ?? "考试");
const studentIds = payload.studentIds ?? payload.student_ids;
const content = `你的考试「${examTitle}」已被教师强制提交`;
await this.broadcastExamRealtimeEvent(
eventId,
"exam",
"考试强制提交通知",
content,
"exam_force_submitted",
examId,
studentIds,
rawPayload,
);
}
/**
* 考试题目顺序调整通知。
* payload: { examId, classId, studentIds? }
*/
private async handleExamQuestionReordered(
eventId: string,
payload: Record<string, unknown>,
rawPayload: Record<string, unknown>,
): Promise<void> {
const examId = String(payload.examId ?? payload.exam_id ?? "");
const examTitle = String(payload.examTitle ?? payload.exam_title ?? "考试");
const studentIds = payload.studentIds ?? payload.student_ids;
const content = `你的考试「${examTitle}」题目顺序已调整,请刷新查看`;
await this.broadcastExamRealtimeEvent(
eventId,
"exam",
"题目顺序调整通知",
content,
"exam_question_reordered",
examId,
studentIds,
rawPayload,
);
}
/**
* 考试实时事件广播:为每个学生创建 in_app 通知(触发 ChannelDispatcher → push-gateway HTTP + Outbox → Kafka
* 若 payload 无 studentIds记录 warn 并跳过(等待 core-edu 补全 payload
*/
private async broadcastExamRealtimeEvent(
eventId: string,
type: string,
title: string,
content: string,
eventTag: string,
examId: string,
studentIds: unknown,
rawPayload: Record<string, unknown>,
): Promise<void> {
if (!Array.isArray(studentIds) || studentIds.length === 0) {
logger.warn(
{ eventId, eventTag, examId },
"Exam realtime event missing studentIds, skipping notification (core-edu should include studentIds in payload)",
);
return;
}
const businessEventId = String(
rawPayload.eventId ?? rawPayload.event_id ?? eventId,
);
const groupId =
businessEventId.length > 32
? businessEventId.slice(0, 32)
: businessEventId;
for (const userId of studentIds) {
await this.notificationsService.send({
userId: String(userId),
type,
title,
content,
channel: "in_app",
groupId,
eventId: `${businessEventId}:${userId}`,
relatedEntityType: "exam",
relatedEntityId: examId,
metadata: { source: "core-edu", event: eventTag },
});
}
}
// ============================================================
// 事件处理器(每个创建对应通知)
// ============================================================
private async handleUserCreated(
eventId: string,
payload: Record<string, unknown>,
): Promise<void> {
const userId = String(payload.userId ?? payload.user_id ?? "");
if (!userId) return;
const name = String(payload.name ?? payload.username ?? "新用户");
await this.notificationsService.send({
userId,
type: "system",
title: "欢迎加入 Edu 云课堂",
content: `你好 ${name},欢迎加入 Edu 云课堂!开始你的学习之旅吧。`,
channel: "in_app",
eventId,
metadata: { source: "iam", event: "user.created" },
});
}
private async handleRoleChanged(
eventId: string,
payload: Record<string, unknown>,
): Promise<void> {
const userId = String(payload.userId ?? payload.user_id ?? "");
if (!userId) return;
const oldRole = String(payload.oldRole ?? payload.old_role ?? "未知");
const newRole = String(payload.newRole ?? payload.new_role ?? "未知");
await this.notificationsService.send({
userId,
type: "system",
title: "角色变更通知",
content: `你的角色已从「${oldRole}」变更为「${newRole}`,
channel: "in_app",
eventId,
metadata: { source: "iam", event: "user.role_changed" },
});
}
private async handleRoleUpdated(
eventId: string,
payload: Record<string, unknown>,
): Promise<void> {
const affectedUserIds =
payload.affectedUserIds ?? payload.affected_user_ids;
if (!Array.isArray(affectedUserIds)) return;
for (const userId of affectedUserIds) {
await this.notificationsService.send({
userId: String(userId),
type: "system",
title: "权限变更通知",
content: "你的角色权限已更新,请查看最新权限。",
channel: "in_app",
eventId: `${eventId}:${userId}`,
metadata: { source: "iam", event: "role.updated" },
});
}
}
private async handleExamPublished(
eventId: string,
payload: Record<string, unknown>,
): Promise<void> {
const studentIds = payload.studentIds ?? payload.student_ids;
if (!Array.isArray(studentIds)) return;
const examTitle = String(payload.examTitle ?? payload.exam_title ?? "考试");
const className = String(payload.className ?? payload.class_name ?? "");
// groupId 优先用 payload 业务 eventId降级路径的 Kafka eventId 可能超长truncate 到 32
const businessEventId = String(
payload.eventId ?? payload.event_id ?? eventId,
);
const groupId =
businessEventId.length > 32
? businessEventId.slice(0, 32)
: businessEventId;
for (const userId of studentIds) {
await this.notificationsService.send({
userId: String(userId),
type: "exam",
title: "新考试通知",
content: `${className}」班级发布了新考试:${examTitle}`,
channel: "in_app",
groupId,
eventId: `${businessEventId}:${userId}`,
relatedEntityType: "exam",
relatedEntityId: String(payload.examId ?? payload.exam_id ?? ""),
metadata: { source: "core-edu", event: "exam.published" },
});
}
}
private async handleHomeworkAssigned(
eventId: string,
payload: Record<string, unknown>,
): Promise<void> {
const studentIds = payload.studentIds ?? payload.student_ids;
if (!Array.isArray(studentIds)) return;
const homeworkTitle = String(
payload.homeworkTitle ?? payload.homework_title ?? "作业",
);
const className = String(payload.className ?? payload.class_name ?? "");
const subject = String(payload.subject ?? "");
// groupId 优先用 payload 业务 eventId降级路径的 Kafka eventId 可能超长truncate 到 32
const businessEventId = String(
payload.eventId ?? payload.event_id ?? eventId,
);
const groupId =
businessEventId.length > 32
? businessEventId.slice(0, 32)
: businessEventId;
for (const userId of studentIds) {
await this.notificationsService.send({
userId: String(userId),
type: "homework",
title: "新作业通知",
content: `${className ? `${className}` : ""}${subject ? `${subject}` : ""}布置了新作业:${homeworkTitle}`,
channel: "in_app",
groupId,
eventId: `${businessEventId}:${userId}`,
relatedEntityType: "homework",
relatedEntityId: String(
payload.homeworkId ?? payload.homework_id ?? "",
),
metadata: { source: "core-edu", event: "homework.assigned" },
});
}
}
private async handleAssignmentSubmitted(
eventId: string,
payload: Record<string, unknown>,
): Promise<void> {
const teacherId = String(payload.teacherId ?? payload.teacher_id ?? "");
if (!teacherId) return;
const studentName = String(
payload.studentName ?? payload.student_name ?? "学生",
);
const homeworkTitle = String(
payload.homeworkTitle ?? payload.homework_title ?? "作业",
);
await this.notificationsService.send({
userId: teacherId,
type: "homework",
title: "作业提交通知",
content: `${studentName} 提交了作业:${homeworkTitle}`,
channel: "in_app",
eventId,
relatedEntityType: "homework",
relatedEntityId: String(payload.homeworkId ?? payload.homework_id ?? ""),
metadata: { source: "core-edu", event: "assignment.submitted" },
});
}
private async handleAssignmentGraded(
eventId: string,
payload: Record<string, unknown>,
): Promise<void> {
const studentId = String(payload.studentId ?? payload.student_id ?? "");
if (!studentId) return;
const homeworkTitle = String(
payload.homeworkTitle ?? payload.homework_title ?? "作业",
);
const score = payload.score ?? payload.grade;
await this.notificationsService.send({
userId: studentId,
type: "grade",
title: "作业批改通知",
content: `你的作业「${homeworkTitle}」已批改${score ? `,得分:${score}` : ""}`,
channel: "in_app",
eventId,
relatedEntityType: "homework",
relatedEntityId: String(payload.homeworkId ?? payload.homework_id ?? ""),
metadata: { source: "core-edu", event: "assignment.graded" },
});
}
private async handleGradeRecorded(
eventId: string,
payload: Record<string, unknown>,
): Promise<void> {
const studentId = String(payload.studentId ?? payload.student_id ?? "");
if (!studentId) return;
const subject = String(payload.subject ?? "科目");
const score = payload.score ?? payload.grade;
await this.notificationsService.send({
userId: studentId,
type: "grade",
title: "成绩录入通知",
content: `你的${subject}成绩已录入${score ? `${score}` : ""}`,
channel: "in_app",
eventId,
relatedEntityType: "grade",
relatedEntityId: String(payload.gradeId ?? payload.grade_id ?? ""),
metadata: { source: "core-edu", event: "grade.recorded" },
});
}
private async handleAttendanceRecorded(
eventId: string,
payload: Record<string, unknown>,
): Promise<void> {
const parentId = String(payload.parentId ?? payload.parent_id ?? "");
if (!parentId) return;
const studentName = String(
payload.studentName ?? payload.student_name ?? "学生",
);
const status = String(payload.status ?? "缺勤");
const date = String(payload.date ?? "");
await this.notificationsService.send({
userId: parentId,
type: "attendance",
title: "出勤异常通知",
content: `${studentName}${date} 的出勤状态为:${status}`,
channel: "in_app",
eventId,
relatedEntityType: "attendance",
relatedEntityId: String(
payload.attendanceId ?? payload.attendance_id ?? "",
),
metadata: { source: "core-edu", event: "attendance.recorded" },
});
}
private async handleMasteryUpdated(
eventId: string,
payload: Record<string, unknown>,
): Promise<void> {
const studentId = String(payload.studentId ?? payload.student_id ?? "");
if (!studentId) return;
const subject = String(payload.subject ?? "科目");
const mastery = payload.mastery ?? payload.masteryLevel;
const trend = String(payload.trend ?? "下降");
await this.notificationsService.send({
userId: studentId,
type: "mastery",
title: "学情预警通知",
content: `你的${subject}掌握度${trend}(当前:${mastery ?? "未知"}),建议加强复习`,
channel: "in_app",
eventId,
relatedEntityType: "mastery",
relatedEntityId: String(payload.masteryId ?? payload.mastery_id ?? ""),
metadata: { source: "data-ana", event: "mastery.updated" },
});
}
}