NestJS (6 services): implement @RequirePermission decorator with SetMetadata+Reflector, register APP_GUARD globally, fix as assertions to type guards, add explicit return types, fix import type for express, fix /metrics implicit any, replace native Error with ApplicationError, remove typeorm remnants, register LifecycleService. teacher-bff: add logger, ApplicationError, GlobalErrorFilter, forward real userId to downstream, log downstream failures, migrate health controller to shared/health. Go (2 services): interface to any, doc comments, CORS dev whitelist, JWT secret fail-fast, push-gateway internal API auth, metrics and readyz endpoints, remove dead code. Python (2 services): lifespan return type, dev_mode to bool, data-ana APIRouter, ai POST body model, ClickHouse async wrapping.
300 lines
9.2 KiB
Python
300 lines
9.2 KiB
Python
"""数据分析服务入口.
|
||
|
||
支持 ClickHouse 降级模式:当 CLICKHOUSE_HOST 未配置或不可达时,
|
||
查询端点返回骨架数据,服务仍可启动与响应。
|
||
|
||
支持 CDC 消费者:当 KAFKA_BROKERS 配置时,
|
||
后台启动 aiokafka 消费者,监听 Debezium CDC 事件写入 ClickHouse。
|
||
"""
|
||
|
||
import asyncio
|
||
import contextlib
|
||
from collections.abc import AsyncGenerator
|
||
from contextlib import asynccontextmanager
|
||
from datetime import UTC, datetime
|
||
|
||
import structlog
|
||
from fastapi import APIRouter, FastAPI
|
||
from opentelemetry import trace
|
||
from opentelemetry.exporter.otlp.proto.http.trace_exporter import OTLPSpanExporter
|
||
from opentelemetry.instrumentation.fastapi import FastAPIInstrumentor
|
||
from opentelemetry.sdk.trace import TracerProvider
|
||
from opentelemetry.sdk.trace.export import BatchSpanProcessor
|
||
from prometheus_client import make_asgi_app
|
||
|
||
from .cdc_consumer import run_consumer as run_cdc_consumer
|
||
from .clickhouse_client import (
|
||
close_client,
|
||
query_class_performance,
|
||
query_dashboard,
|
||
query_student_errors,
|
||
)
|
||
from .clickhouse_client import ping as ch_ping
|
||
from .config import settings
|
||
|
||
_logger: structlog.stdlib.BoundLogger | None = None
|
||
tracer = trace.get_tracer(__name__)
|
||
|
||
# CDC 消费者后台任务句柄
|
||
_cdc_task: asyncio.Task | None = None
|
||
|
||
# 日志级别映射
|
||
_LOG_LEVELS: dict[str, int] = {
|
||
"DEBUG": 10,
|
||
"INFO": 20,
|
||
"WARNING": 30,
|
||
"ERROR": 40,
|
||
"CRITICAL": 50,
|
||
}
|
||
|
||
|
||
def init_logger() -> structlog.stdlib.BoundLogger:
|
||
"""初始化 structlog logger.
|
||
|
||
根据配置的 log_level 设置日志级别。
|
||
"""
|
||
global _logger
|
||
level = _LOG_LEVELS.get(settings.log_level.upper(), 20)
|
||
structlog.configure(
|
||
wrapper_class=structlog.make_filtering_bound_logger(level),
|
||
processors=[
|
||
structlog.contextvars.merge_contextvars,
|
||
structlog.processors.add_log_level,
|
||
structlog.processors.TimeStamper(fmt="iso"),
|
||
structlog.dev.ConsoleRenderer(),
|
||
],
|
||
cache_logger_on_first_use=True,
|
||
)
|
||
_logger = structlog.get_logger(__name__)
|
||
return _logger
|
||
|
||
|
||
def get_logger() -> structlog.stdlib.BoundLogger:
|
||
"""获取已初始化的 logger(未初始化时自动初始化)."""
|
||
global _logger
|
||
if _logger is None:
|
||
return init_logger()
|
||
return _logger
|
||
|
||
|
||
def init_tracer() -> None:
|
||
"""初始化 OpenTelemetry.
|
||
|
||
endpoint 从 settings.otel_endpoint 读取(不硬编码)。
|
||
"""
|
||
provider = TracerProvider()
|
||
endpoint = settings.otel_endpoint.rstrip("/")
|
||
exporter = OTLPSpanExporter(endpoint=f"{endpoint}/v1/traces")
|
||
provider.add_span_processor(BatchSpanProcessor(exporter))
|
||
trace.set_tracer_provider(provider)
|
||
|
||
|
||
@asynccontextmanager
|
||
async def lifespan(app: FastAPI) -> AsyncGenerator[None, None]:
|
||
"""应用生命周期.
|
||
|
||
1. 初始化 logger(structlog)
|
||
2. 初始化 OTel tracer(endpoint 从 config 读)
|
||
3. 触发 ClickHouse 客户端惰性初始化(不阻塞启动,失败进入降级模式)
|
||
4. 若配置了 kafka_brokers,后台启动 CDC 消费者任务
|
||
5. 关闭时停止 CDC 任务并释放 ClickHouse 客户端
|
||
"""
|
||
global _cdc_task
|
||
logger = init_logger()
|
||
init_tracer()
|
||
logger.info(
|
||
"data_ana_service_starting",
|
||
port=settings.port,
|
||
dev_mode=settings.dev_mode,
|
||
clickhouse_configured=bool(settings.clickhouse_host),
|
||
kafka_brokers=settings.kafka_brokers,
|
||
kafka_cdc_topics=settings.kafka_cdc_topics,
|
||
)
|
||
# 启动 CDC 消费者后台任务(若未配置 kafka_brokers,run_consumer 内部直接返回)
|
||
_cdc_task = asyncio.create_task(run_cdc_consumer())
|
||
yield
|
||
logger.info("data_ana_service_stopping")
|
||
# 取消 CDC 任务
|
||
if _cdc_task is not None and not _cdc_task.done():
|
||
_cdc_task.cancel()
|
||
with contextlib.suppress(asyncio.CancelledError):
|
||
await _cdc_task
|
||
await close_client()
|
||
|
||
|
||
app = FastAPI(
|
||
title="Data Analytics Service",
|
||
version="0.1.0",
|
||
lifespan=lifespan,
|
||
)
|
||
|
||
# OpenTelemetry FastAPI 自动埋点(HTTP 请求/响应 span)
|
||
FastAPIInstrumentor.instrument_app(app)
|
||
|
||
# Prometheus 指标
|
||
app.mount("/metrics", make_asgi_app())
|
||
|
||
# 业务路由
|
||
router = APIRouter()
|
||
|
||
|
||
@app.get("/healthz")
|
||
async def healthz() -> dict:
|
||
"""健康检查(liveness).
|
||
|
||
只要进程存活即返回 ok,不依赖 ClickHouse。
|
||
"""
|
||
return {"status": "ok", "service": "data-ana"}
|
||
|
||
|
||
@app.get("/readyz")
|
||
async def readyz() -> dict:
|
||
"""就绪检查(readiness).
|
||
|
||
ClickHouse 为可选依赖:
|
||
- 已配置且可达:ready=true
|
||
- 未配置:ready=true,degraded=true(降级模式仍可服务)
|
||
- 已配置但不可达:ready=false
|
||
|
||
CDC 消费者状态附加在响应中:
|
||
- cdc_consumer: running / disabled / failed
|
||
"""
|
||
cdc_status = "disabled"
|
||
if _cdc_task is not None:
|
||
if _cdc_task.done():
|
||
cdc_status = "failed"
|
||
elif not settings.kafka_brokers:
|
||
cdc_status = "disabled"
|
||
else:
|
||
cdc_status = "running"
|
||
|
||
if not settings.clickhouse_host:
|
||
return {
|
||
"status": "ok",
|
||
"service": "data-ana",
|
||
"ready": True,
|
||
"degraded": True,
|
||
"clickhouse": "not_configured",
|
||
"cdc_consumer": cdc_status,
|
||
"kafka_brokers": settings.kafka_brokers or None,
|
||
"timestamp": datetime.now(UTC).isoformat(),
|
||
}
|
||
|
||
ch_ok = await ch_ping()
|
||
return {
|
||
"status": "ok" if ch_ok else "degraded",
|
||
"service": "data-ana",
|
||
"ready": ch_ok,
|
||
"degraded": not ch_ok,
|
||
"clickhouse": "ok" if ch_ok else "unreachable",
|
||
"cdc_consumer": cdc_status,
|
||
"kafka_brokers": settings.kafka_brokers or None,
|
||
"timestamp": datetime.now(UTC).isoformat(),
|
||
}
|
||
|
||
|
||
@router.get("/analytics/class/{class_id}/performance")
|
||
async def class_performance(class_id: str) -> dict:
|
||
"""班级成绩分析.
|
||
|
||
优先查 ClickHouse;降级时返回骨架数据。
|
||
"""
|
||
logger = get_logger()
|
||
with tracer.start_as_current_span("class_performance") as span:
|
||
span.set_attribute("class_id", class_id)
|
||
result = await query_class_performance(class_id)
|
||
if result is None:
|
||
logger.info("class_performance_degraded", class_id=class_id)
|
||
return {
|
||
"success": True,
|
||
"data": {
|
||
"classId": class_id,
|
||
"averageScore": 0,
|
||
"passRate": 0,
|
||
"totalStudents": 0,
|
||
"message": "ClickHouse unavailable - skeleton data",
|
||
"degraded": True,
|
||
},
|
||
}
|
||
return {"success": True, "data": {**result, "degraded": False}}
|
||
|
||
|
||
@router.get("/analytics/student/{student_id}/weakness")
|
||
async def student_weakness(student_id: str) -> dict:
|
||
"""学生薄弱知识点分析.
|
||
|
||
优先查 ClickHouse;降级时返回骨架数据。
|
||
"""
|
||
logger = get_logger()
|
||
with tracer.start_as_current_span("student_weakness") as span:
|
||
span.set_attribute("student_id", student_id)
|
||
result = await query_dashboard(student_id)
|
||
if result is None:
|
||
logger.info("student_weakness_degraded", student_id=student_id)
|
||
return {
|
||
"success": True,
|
||
"data": {
|
||
"studentId": student_id,
|
||
"weakPoints": [],
|
||
"message": "ClickHouse unavailable - skeleton data",
|
||
"degraded": True,
|
||
},
|
||
}
|
||
|
||
# 从宽表提取薄弱知识点:mastery_level < 0.6 视为薄弱
|
||
weak_points = [
|
||
{
|
||
"knowledgePointId": r["knowledge_point_id"],
|
||
"masteryLevel": r["mastery_level"],
|
||
"errorCount": r["error_count"],
|
||
}
|
||
for r in result["records"]
|
||
if r.get("mastery_level") is not None and r["mastery_level"] < 0.6
|
||
]
|
||
return {
|
||
"success": True,
|
||
"data": {
|
||
"studentId": student_id,
|
||
"weakPoints": weak_points,
|
||
"records": result["records"],
|
||
"total": result["total"],
|
||
"degraded": False,
|
||
},
|
||
}
|
||
|
||
|
||
@router.get("/analytics/student/{student_id}/errorbook")
|
||
async def student_errorbook(student_id: str) -> dict:
|
||
"""学生错题本.
|
||
|
||
优先查 ClickHouse;降级时返回空列表。
|
||
"""
|
||
logger = get_logger()
|
||
with tracer.start_as_current_span("student_errorbook") as span:
|
||
span.set_attribute("student_id", student_id)
|
||
result = await query_student_errors(student_id)
|
||
if result is None:
|
||
logger.info("student_errorbook_degraded", student_id=student_id)
|
||
return {
|
||
"success": True,
|
||
"data": {
|
||
"studentId": student_id,
|
||
"errors": [],
|
||
"total": 0,
|
||
"message": "ClickHouse unavailable - empty errorbook",
|
||
"degraded": True,
|
||
},
|
||
}
|
||
return {
|
||
"success": True,
|
||
"data": {
|
||
"studentId": student_id,
|
||
"errors": result,
|
||
"total": len(result),
|
||
"degraded": False,
|
||
},
|
||
}
|
||
|
||
|
||
app.include_router(router)
|