"""数据分析服务入口(FastAPI HTTP :3006). 端点清单(3 基础 + 11 业务 = 14 个): 基础: GET / 根信息 GET /healthz 活性检查(liveness) GET /readyz 就绪检查(readiness,检查 4 依赖) 业务(全部返回 ActionState[T] 信封): GET /analytics/class/{class_id}/performance 班级成绩分析 GET /analytics/student/{student_id}/weakness 学生薄弱知识点 GET /analytics/student/{student_id}/trend 学习趋势 GET /analytics/student/{student_id}/errorbook 错题本(额外) GET /analytics/teacher/dashboard 教师仪表盘 GET /analytics/student/dashboard 学生仪表盘 GET /analytics/parent/dashboard 家长仪表盘 GET /analytics/admin/dashboard 管理员仪表盘 GET /analytics/warnings 预警列表 POST /analytics/warnings/trigger 手动触发预警 GET /analytics/class/{class_id}/mastery-distribution 班级掌握度分布 GET /analytics/student/{student_id}/mastery 学生掌握度明细 设计要点: - 所有业务端点返回 ActionState[T](coord-cross-review §5.3 P0 整改) - 降级标记在顶层 details.degraded(不放 error.details) - /readyz 检查 4 依赖:clickhouse / cdc_consumer / redis / iam_grpc - gRPC server :50055 在 lifespan 启动 - CDC 消费者在 lifespan 启动 """ from collections.abc import AsyncGenerator from contextlib import asynccontextmanager from datetime import UTC, datetime from typing import Any import structlog from fastapi import APIRouter, Depends, FastAPI, Query 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 . import analytics_service, cdc_consumer, grpc_server, warning_service from .config import settings from .repository import ( clickhouse_repository, iam_client, kafka_producer, redis_client, ) from .shared.action_state import ActionState from .shared.permissions import UserContext, get_user_context _logger: structlog.stdlib.BoundLogger | None = None tracer = trace.get_tracer(__name__) _LOG_LEVELS: dict[str, int] = { "DEBUG": 10, "INFO": 20, "WARNING": 30, "ERROR": 40, "CRITICAL": 50, } def init_logger() -> structlog.stdlib.BoundLogger: """初始化 structlog logger.""" 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.""" 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 + tracer 2. 启动 gRPC server :50055(AnalyticsService 12 RPC) 3. 启动 CDC 消费者后台任务(手动 commit) 4. Kafka producer 惰性初始化(首次发布时触发) 关闭顺序: 1. 停止 CDC 消费者 2. 停止 gRPC server 3. 关闭 Kafka producer 4. 关闭 Redis / iam gRPC / ClickHouse 客户端 """ logger = init_logger() init_tracer() logger.info( "data_ana_service_starting", http_port=settings.http_port, grpc_port=settings.grpc_port, dev_mode=settings.dev_mode, clickhouse_configured=bool(settings.clickhouse_host), kafka_brokers=settings.kafka_brokers, iam_grpc_endpoint=settings.iam_grpc_endpoint, redis_url=settings.redis_url or "not_configured", ) # 1. 启动 gRPC server(grpcio 未安装则跳过,降级为仅 HTTP) grpc_server_obj = await grpc_server.start_grpc_server() if grpc_server_obj is None: logger.warning("grpc_server_not_started_http_only") # 2. 启动 CDC 消费者后台任务 await cdc_consumer.start_consumer() yield logger.info("data_ana_service_stopping") # 3. 停止 CDC 消费者 await cdc_consumer.stop_consumer() # 4. 停止 gRPC server await grpc_server.stop_grpc_server() # 5. 关闭 Kafka producer await kafka_producer.close_producer() # 6. 关闭 Redis / iam gRPC / ClickHouse 客户端 await redis_client.close_client() await iam_client.close_grpc() await clickhouse_repository.close_client() app = FastAPI( title="Data Analytics Service", version="1.0.0", description="D6 智能洞察领域服务(ClickHouse 宽表 + CDC + 掌握度算法 + 预警)", lifespan=lifespan, ) FastAPIInstrumentor.instrument_app(app) app.mount("/metrics", make_asgi_app()) router = APIRouter() # ===== 基础端点(3 个) ===== @app.get("/") async def root() -> dict[str, Any]: """根信息.""" return { "service": "data-ana", "version": "1.0.0", "http_port": settings.http_port, "grpc_port": settings.grpc_port, "docs": "/docs", "endpoints": { "healthz": "/healthz", "readyz": "/readyz", "business": "/analytics/*", }, } @app.get("/healthz") async def healthz() -> dict[str, str]: """健康检查(liveness,只要进程存活即返回 ok).""" return {"status": "ok", "service": "data-ana"} @app.get("/readyz") async def readyz() -> dict[str, Any]: """就绪检查(readiness,检查 4 依赖). 依赖检查: 1. clickhouse:已配置且可达(未配置算降级就绪) 2. cdc_consumer:running / disabled 3. redis:已配置且可达(未配置算降级就绪) 4. iam_grpc:已配置且可达(未配置算降级就绪) 返回 ready=true 的条件: - ClickHouse 已配置且可达,或未配置(降级就绪) - 不要求所有依赖都健康(降级模式下仍可服务骨架数据) """ # 1. ClickHouse ch_ok = await clickhouse_repository.ping() ch_status = "ok" if ch_ok else ("unreachable" if settings.clickhouse_host else "not_configured") # 2. CDC 消费者 cdc_status = ( "running" if cdc_consumer.is_running() else ("disabled" if not settings.kafka_brokers else "failed") ) # 3. Redis redis_ok = await redis_client.ping() if settings.redis_url else None redis_status = "ok" if redis_ok else ("unreachable" if settings.redis_url else "not_configured") # 4. iam gRPC iam_ok = await iam_client.ping() if settings.iam_grpc_endpoint else None iam_status = ( "ok" if iam_ok else ("unreachable" if settings.iam_grpc_endpoint else "not_configured") ) # 5. gRPC server grpc_status = "running" if grpc_server.is_running() else "stopped" # 就绪判定:ClickHouse 可达或未配置(降级就绪) ready = ch_ok or not settings.clickhouse_host degraded = not ch_ok or not redis_ok or not iam_ok return { "status": "ok" if ready else "not_ready", "service": "data-ana", "ready": ready, "degraded": degraded, "dependencies": { "clickhouse": ch_status, "cdc_consumer": cdc_status, "redis": redis_status, "iam_grpc": iam_status, "grpc_server": grpc_status, "kafka_producer": "ok" if settings.kafka_brokers else "not_configured", }, "timestamp": datetime.now(UTC).isoformat(), } # ===== 业务端点(11 个,全部返回 ActionState[T]) ===== @router.get("/analytics/class/{class_id}/performance") async def get_class_performance( class_id: str, subject_id: str = Query(""), start_date: int = Query(0), end_date: int = Query(0), user: UserContext = Depends(get_user_context), ) -> ActionState[dict[str, Any]]: """班级成绩分析.""" with tracer.start_as_current_span("get_class_performance") as span: span.set_attribute("class_id", class_id) result = await analytics_service.get_class_performance( user=user, class_id=class_id, subject_id=subject_id, start_date=start_date, end_date=end_date, ) degraded = result.get("degraded", False) return ActionState.ok( result, degraded=degraded, degraded_reason=result.get("degraded_reason", "") if degraded else "", ) @router.get("/analytics/student/{student_id}/weakness") async def get_student_weakness( student_id: str, subject_id: str = Query(""), user: UserContext = Depends(get_user_context), ) -> ActionState[dict[str, Any]]: """学生薄弱知识点.""" with tracer.start_as_current_span("get_student_weakness") as span: span.set_attribute("student_id", student_id) result = await analytics_service.get_student_weakness( user=user, student_id=student_id, subject_id=subject_id, ) degraded = result.get("degraded", False) return ActionState.ok( result, degraded=degraded, degraded_reason=result.get("degraded_reason", "") if degraded else "", ) @router.get("/analytics/student/{student_id}/trend") async def get_learning_trend( student_id: str, subject_id: str = Query(""), start_date: int = Query(0), end_date: int = Query(0), user: UserContext = Depends(get_user_context), ) -> ActionState[dict[str, Any]]: """学习趋势.""" with tracer.start_as_current_span("get_learning_trend") as span: span.set_attribute("student_id", student_id) result = await analytics_service.get_learning_trend( user=user, student_id=student_id, subject_id=subject_id, start_date=start_date, end_date=end_date, ) degraded = result.get("degraded", False) return ActionState.ok( result, degraded=degraded, degraded_reason=result.get("degraded_reason", "") if degraded else "", ) @router.get("/analytics/student/{student_id}/errorbook") async def get_student_errorbook( student_id: str, user: UserContext = Depends(get_user_context), ) -> ActionState[dict[str, Any]]: """学生错题本.""" with tracer.start_as_current_span("get_student_errorbook") as span: span.set_attribute("student_id", student_id) errors = await clickhouse_repository.query_student_errors(student_id) if errors is None: return ActionState.ok( {"studentId": student_id, "errors": [], "total": 0}, degraded=True, degraded_reason="clickhouse_unavailable", ) return ActionState.ok( { "studentId": student_id, "errors": errors, "total": len(errors), } ) @router.get("/analytics/teacher/dashboard") async def get_teacher_dashboard( class_id: str = Query(""), subject_id: str = Query(""), user: UserContext = Depends(get_user_context), ) -> ActionState[dict[str, Any]]: """教师仪表盘.""" with tracer.start_as_current_span("get_teacher_dashboard"): result = await analytics_service.get_teacher_dashboard( user=user, class_id=class_id, subject_id=subject_id, ) degraded = result.get("degraded", False) return ActionState.ok( result, degraded=degraded, degraded_reason=result.get("degraded_reason", "") if degraded else "", ) @router.get("/analytics/student/dashboard") async def get_student_dashboard( subject_id: str = Query(""), user: UserContext = Depends(get_user_context), ) -> ActionState[dict[str, Any]]: """学生仪表盘.""" with tracer.start_as_current_span("get_student_dashboard"): result = await analytics_service.get_student_dashboard( user=user, subject_id=subject_id, ) degraded = result.get("degraded", False) return ActionState.ok( result, degraded=degraded, degraded_reason=result.get("degraded_reason", "") if degraded else "", ) @router.get("/analytics/parent/dashboard") async def get_parent_dashboard( student_id: str = Query(""), user: UserContext = Depends(get_user_context), ) -> ActionState[dict[str, Any]]: """家长仪表盘.""" with tracer.start_as_current_span("get_parent_dashboard"): result = await analytics_service.get_parent_dashboard( user=user, child_id=student_id, ) degraded = result.get("degraded", False) return ActionState.ok( result, degraded=degraded, degraded_reason=result.get("degraded_reason", "") if degraded else "", ) @router.get("/analytics/admin/dashboard") async def get_admin_dashboard( school_id: str = Query(""), user: UserContext = Depends(get_user_context), ) -> ActionState[dict[str, Any]]: """管理员仪表盘.""" with tracer.start_as_current_span("get_admin_dashboard"): result = await analytics_service.get_admin_dashboard( user=user, school_id=school_id, ) degraded = result.get("degraded", False) return ActionState.ok( result, degraded=degraded, degraded_reason=result.get("degraded_reason", "") if degraded else "", ) @router.get("/analytics/warnings") async def get_warnings( student_id: str = Query(""), warning_type: str = Query(""), user: UserContext = Depends(get_user_context), ) -> ActionState[dict[str, Any]]: """预警列表查询.""" with tracer.start_as_current_span("get_warnings"): warnings = await warning_service.get_warnings( student_id=student_id, warning_type=warning_type, ) return ActionState.ok( { "warnings": warnings, "total": len(warnings), } ) @router.post("/analytics/warnings/trigger") async def trigger_warning( target_id: str = Query(...), warning_type: str = Query(...), severity: str = Query("WARN"), threshold: float = Query(0.0), current_value: float = Query(0.0), user: UserContext = Depends(get_user_context), ) -> ActionState[dict[str, Any]]: """手动触发预警.""" with tracer.start_as_current_span("trigger_warning"): result = await warning_service.trigger_warning_manual( target_id=target_id, warning_type=warning_type, threshold=threshold, current_value=current_value, severity=severity, ) return ActionState.ok(result) @router.get("/analytics/class/{class_id}/mastery-distribution") async def get_mastery_distribution( class_id: str, subject_id: str = Query(""), knowledge_point_id: str = Query(""), user: UserContext = Depends(get_user_context), ) -> ActionState[dict[str, Any]]: """班级掌握度分布.""" with tracer.start_as_current_span("get_mastery_distribution"): result = await analytics_service.get_mastery_distribution( user=user, class_id=class_id, subject_id=subject_id, knowledge_point_id=knowledge_point_id, ) degraded = result.get("degraded", False) return ActionState.ok( result, degraded=degraded, degraded_reason=result.get("degraded_reason", "") if degraded else "", ) @router.get("/analytics/student/{student_id}/mastery") async def get_student_mastery( student_id: str, subject_id: str = Query(""), user: UserContext = Depends(get_user_context), ) -> ActionState[dict[str, Any]]: """学生知识点掌握度明细.""" with tracer.start_as_current_span("get_student_mastery"): result = await analytics_service.get_student_mastery( user=user, student_id=student_id, subject_id=subject_id, ) degraded = result.get("degraded", False) return ActionState.ok( result, degraded=degraded, degraded_reason=result.get("degraded_reason", "") if degraded else "", ) app.include_router(router)