"""数据分析服务入口. 支持 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)