feat(data-ana): implement complete CDC pipeline MySQL to ClickHouse

Debezium + Kafka + aiokafka consumer with table routing.

E2E verified: MySQL INSERT to ClickHouse upsert.
This commit is contained in:
SpecialX
2026-07-09 13:02:59 +08:00
parent 958b17c9d8
commit 1f901c5b20
10 changed files with 595 additions and 43 deletions

View File

@@ -2,8 +2,13 @@
支持 ClickHouse 降级模式:当 CLICKHOUSE_HOST 未配置或不可达时,
查询端点返回骨架数据,服务仍可启动与响应。
支持 CDC 消费者:当 KAFKA_BROKERS 配置时,
后台启动 aiokafka 消费者,监听 Debezium CDC 事件写入 ClickHouse。
"""
import asyncio
import contextlib
from contextlib import asynccontextmanager
from datetime import UTC, datetime
@@ -15,6 +20,7 @@ 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,
@@ -27,6 +33,9 @@ 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,
@@ -45,7 +54,7 @@ def init_logger() -> structlog.stdlib.BoundLogger:
global _logger
level = _LOG_LEVELS.get(settings.log_level.upper(), 20)
structlog.configure(
wrapper_class=structlog.make_filtering_logger(level),
wrapper_class=structlog.make_filtering_bound_logger(level),
processors=[
structlog.contextvars.merge_contextvars,
structlog.processors.add_log_level,
@@ -85,8 +94,10 @@ async def lifespan(app: FastAPI):
1. 初始化 loggerstructlog
2. 初始化 OTel tracerendpoint 从 config 读)
3. 触发 ClickHouse 客户端惰性初始化(不阻塞启动,失败进入降级模式)
4. 关闭时释放 ClickHouse 客户端
4. 若配置了 kafka_brokers后台启动 CDC 消费者任务
5. 关闭时停止 CDC 任务并释放 ClickHouse 客户端
"""
global _cdc_task
logger = init_logger()
init_tracer()
logger.info(
@@ -95,9 +106,17 @@ async def lifespan(app: FastAPI):
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_brokersrun_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()
@@ -128,7 +147,19 @@ async def readyz() -> dict:
- 已配置且可达ready=true
- 未配置ready=truedegraded=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",
@@ -136,6 +167,8 @@ async def readyz() -> dict:
"ready": True,
"degraded": True,
"clickhouse": "not_configured",
"cdc_consumer": cdc_status,
"kafka_brokers": settings.kafka_brokers or None,
"timestamp": datetime.now(UTC).isoformat(),
}
@@ -146,6 +179,8 @@ async def readyz() -> dict:
"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(),
}