Files
Edu/services/data-ana/src/data_ana/main.py
SpecialX a3f4fd013e feat(data-ana): python graphql federation subgraph with strawberry
- strawberry-graphql[asgi] dependency added

- 13 Federation 2 types: ClassPerformance/StudentWeakness/Dashboards/Mastery/ErrorBook

- 11 Query resolvers delegate to existing analytics/mastery services

- RouterAuthMiddleware validates Router-Authorization header on /graphql

- GraphQL endpoint mounted at /graphql alongside existing REST endpoints
2026-07-15 00:57:02 +08:00

607 lines
21 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 业务 + 1 GraphQL = 15 个):
基础:
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 学生掌握度明细
GraphQLv2.1 M1Apollo Federation 2 子图):
POST /graphql GraphQL 端点Router-Authorization 校验)
- Resolver 委托给 analytics_service / mastery_service
- 由 Apollo Router 组合不直接对外暴露dev_mode 除外)
设计要点:
- 所有业务端点返回 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 启动
- /graphql 由 strawberry-graphql 提供RouterAuthMiddleware 拦截
"""
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, Response
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 .graphql.router_auth import RouterAuthMiddleware
from .graphql.schema import graphql_app
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. 预热客户端(避免首次 readyz 探针因惰性初始化超时)
if settings.clickhouse_host:
ch_ok = await clickhouse_repository.ping()
logger.info("clickhouse_warmup", ok=ch_ok)
if settings.redis_url:
redis_ok = await redis_client.ping()
logger.info("redis_warmup", ok=redis_ok)
# 2.1 初始化监控指标初始值(确保 /metrics 端点暴露自定义指标)
from . import metrics as _metrics
_metrics.clickhouse_connection_status.set(1 if settings.clickhouse_host else 0)
_metrics.redis_connection_status.set(1 if settings.redis_url else 0)
_metrics.iam_grpc_connection_status.set(1 if settings.iam_grpc_endpoint else 0)
# 3. 启动 CDC 消费者后台任务
await cdc_consumer.start_consumer()
yield
logger.info("data_ana_service_stopping")
# 4. 停止 CDC 消费者
await cdc_consumer.stop_consumer()
# 5. 停止 gRPC server
await grpc_server.stop_grpc_server()
# 6. 关闭 Kafka producer
await kafka_producer.close_producer()
# 7. 关闭 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())
# GraphQL Federation 2 子图v2.1 M1
# RouterAuthMiddleware 仅对 /graphql 端点校验 Router-Authorization Header
# REST 路由放行(已有 PermissionGuard / DataScope 校验)
app.add_middleware(RouterAuthMiddleware)
# 挂载 strawberry ASGI handler 到 /graphqlApollo Router 访问入口)
app.mount(settings.graphql_path, graphql_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/*",
"graphql": settings.graphql_path,
},
}
@app.get("/healthz")
async def healthz() -> dict[str, str]:
"""健康检查liveness只要进程存活即返回 ok."""
return {"status": "ok", "service": "data-ana"}
@app.get("/readyz")
async def readyz(response: Response) -> dict[str, Any]:
"""就绪检查P6 深度硬化版,检查 4 依赖 + 超时控制 + lag 阈值).
依赖检查(带超时):
1. clickhouse1s 超时,已配置且可达(未配置算降级就绪)
2. cdc_consumerrunning / disabled + lag < 1000
3. redis200ms 超时,已配置且可达(未配置算降级就绪)
4. iam_grpc2s 超时,已配置且可达(未配置算降级就绪)
返回 ready=true 的条件:
- ClickHouse 已配置且可达1s 内),或未配置(降级就绪)
- CDC consumer lag < readyz_cdc_lag_threshold1000
不满足时返回 HTTP 503K8s 摘流量.
"""
import asyncio
# 1. ClickHouse1s 超时)
try:
ch_ok = await asyncio.wait_for(
clickhouse_repository.ping(),
timeout=settings.readyz_clickhouse_timeout_s,
)
except TimeoutError:
ch_ok = False
ch_status = "ok" if ch_ok else ("unreachable" if settings.clickhouse_host else "not_configured")
# 2. CDC 消费者 + lag 检查
cdc_running = cdc_consumer.is_running()
cdc_lag = 0
if cdc_running:
try:
cdc_lag = await asyncio.wait_for(
cdc_consumer.get_lag(),
timeout=2.0,
)
except TimeoutError:
cdc_lag = -1 # 查询超时标记
cdc_lag_ok = cdc_lag >= 0 and cdc_lag < settings.readyz_cdc_lag_threshold
if not settings.kafka_brokers:
cdc_status = "disabled"
cdc_lag_ok = True # 未配置 Kafka 时不阻塞 ready
elif cdc_running and cdc_lag_ok:
cdc_status = f"running(lag={cdc_lag})"
elif cdc_running:
cdc_status = f"running(lag={cdc_lag},exceeded)"
else:
cdc_status = "failed"
# 3. Redis200ms 超时)
if settings.redis_url:
try:
redis_ok = await asyncio.wait_for(
redis_client.ping(),
timeout=settings.readyz_redis_timeout_s,
)
except TimeoutError:
redis_ok = False
else:
redis_ok = None
redis_status = "ok" if redis_ok else ("unreachable" if settings.redis_url else "not_configured")
# 4. iam gRPC2s 超时iam_client.ping 已内置 1s 超时)
if settings.iam_grpc_endpoint:
try:
iam_ok = await asyncio.wait_for(
iam_client.ping(),
timeout=settings.readyz_iam_grpc_timeout_s,
)
except TimeoutError:
iam_ok = False
else:
iam_ok = 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"
# 就绪判定P6 硬化):
# - ClickHouse 可达或未配置(降级就绪)
# - CDC consumer lag 未超阈值
ready = (ch_ok or not settings.clickhouse_host) and cdc_lag_ok
degraded = not ch_ok or not redis_ok or not iam_ok
if not ready:
response.status_code = 503 # K8s 摘流量
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",
},
"thresholds": {
"clickhouse_timeout_s": settings.readyz_clickhouse_timeout_s,
"redis_timeout_s": settings.readyz_redis_timeout_s,
"iam_grpc_timeout_s": settings.readyz_iam_grpc_timeout_s,
"cdc_lag_threshold": settings.readyz_cdc_lag_threshold,
"cdc_lag_current": cdc_lag,
},
"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)