Files
Edu/services/data-ana/src/data_ana/main.py
SpecialX ca3780aa24 feat(data-ana): 完整实现 data-ana 数据分析服务
包含 CDC consumer、analytics/mastery/warning service、grpc server、repository、ClickHouse DDL 等
2026-07-10 19:09:27 +08:00

522 lines
17 KiB
Python
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.
"""数据分析服务入口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 :50055AnalyticsService 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 servergrpcio 未安装则跳过,降级为仅 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_consumerrunning / 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)