Files
atc-agents/api/cdc_consumer.py
T

247 lines
8.0 KiB
Python
Raw Normal View History

"""Live CDC consumer for the Command Center "Changes" dashboard.
A background aiokafka consumer subscribes to the Debezium CDC topics on the
Kafka broker, parses each change event (insert/update/delete with before/after
images), keeps a bounded in-memory ring buffer, and fans every change out over
the WebSocket bus as a ``cdc_change`` event so the UI can render new & changed
data in real time.
Endpoints:
GET /api/changes -> recent changes from the ring buffer (filterable)
GET /api/changes/stats -> volume per source/table/op + per-minute buckets
"""
from __future__ import annotations
import asyncio
import json
import os
import re
import uuid
from collections import deque
from datetime import datetime, timezone
from typing import Any
from fastapi import APIRouter, Query
from fastapi.responses import JSONResponse
router = APIRouter(prefix="/api/changes", tags=["changes"])
KAFKA_BOOTSTRAP = os.getenv("KAFKA_BOOTSTRAP", "10.0.21.36:9092")
RING_SIZE = int(os.getenv("CDC_RING_SIZE", "1000"))
# Active CDC topic prefixes -> logical source. Matches Debezium topic.prefix.
PREFIX_SOURCE = {
"postgres_sales": "postgres",
"mysql_hr": "mysql",
"mongodb_supplychain": "mongodb",
"cassandra_telemetry": "cassandra",
"neo4j_graph": "neo4j",
}
TOPIC_PATTERN = re.compile(
r"^(postgres_sales|mysql_hr|mongodb_supplychain|cassandra_telemetry|neo4j_graph)\..+"
)
OP_MAP = {"c": "insert", "u": "update", "d": "delete", "r": "snapshot"}
_ring: deque[dict[str, Any]] = deque(maxlen=RING_SIZE)
_state: dict[str, Any] = {
"started": False,
"connected": False,
"consumed": 0,
"last_error": None,
"last_ts": None,
"topics": [],
}
def _source_of(topic: str) -> tuple[str, str]:
prefix = topic.split(".", 1)[0]
source = PREFIX_SOURCE.get(prefix, prefix)
table = topic.split(".")[-1]
return source, table
def _coerce(v: Any) -> Any:
"""Debezium/Mongo nests the document as a JSON string sometimes."""
if isinstance(v, str):
try:
return json.loads(v)
except Exception:
return v
return v
def _key_fields(source: str, table: str, after: Any, before: Any) -> str:
row = after or before or {}
if not isinstance(row, dict):
return ""
prefer = ["order_id", "event_id", "_id", "id", "region", "order_status",
"event_type", "type", "amount", "salary_change"]
parts = []
for k in prefer:
if k in row and row[k] is not None:
val = row[k]
if isinstance(val, (dict, list)):
continue
sval = str(val)
if len(sval) > 40:
sval = sval[:40] + "…"
parts.append(f"{k}={sval}")
if len(parts) >= 4:
break
return " ".join(parts)
def _parse(topic: str, value: bytes | None) -> dict[str, Any] | None:
if value is None: # tombstone
return None
try:
payload = json.loads(value.decode("utf-8"))
except Exception:
return None
# Debezium envelope: {schema, payload:{before,after,op,source,ts_ms}}
if isinstance(payload, dict) and "payload" in payload and isinstance(payload["payload"], dict):
payload = payload["payload"]
if not isinstance(payload, dict):
return None
op = OP_MAP.get(payload.get("op"), payload.get("op") or "change")
before = _coerce(payload.get("before"))
after = _coerce(payload.get("after"))
src = payload.get("source") or {}
source, table = _source_of(topic)
if isinstance(src, dict) and src.get("table"):
table = src.get("table")
ts_ms = payload.get("ts_ms")
ts = datetime.fromtimestamp(ts_ms / 1000, timezone.utc).isoformat() if ts_ms else datetime.now(timezone.utc).isoformat()
return {
"id": uuid.uuid4().hex[:10],
"ts": ts,
"source": source,
"table": table,
"topic": topic,
"op": op,
"summary": _key_fields(source, table, after, before),
"before": before if isinstance(before, dict) else None,
"after": after if isinstance(after, dict) else None,
}
async def cdc_consumer_loop() -> None:
"""Background loop: consume CDC topics and publish each change live."""
_state["started"] = True
await asyncio.sleep(6)
try:
from aiokafka import AIOKafkaConsumer
except Exception as exc:
_state["last_error"] = f"aiokafka import failed: {exc}"
return
while True:
consumer = None
try:
consumer = AIOKafkaConsumer(
bootstrap_servers=KAFKA_BOOTSTRAP,
group_id=f"atc-cc-cdc-{uuid.uuid4().hex[:8]}",
auto_offset_reset="latest",
enable_auto_commit=False,
client_id="atc-command-center-cdc",
)
consumer.subscribe(pattern=TOPIC_PATTERN)
await consumer.start()
_state["connected"] = True
_state["last_error"] = None
from main import publish_event
async for msg in consumer:
entry = _parse(msg.topic, msg.value)
if not entry:
continue
_ring.append(entry)
_state["consumed"] += 1
_state["last_ts"] = entry["ts"]
try:
_state["topics"] = sorted(consumer.subscription() or [])
except Exception:
pass
try:
await publish_event({"type": "cdc_change", "entry": entry})
except Exception:
pass
except Exception as exc:
_state["connected"] = False
_state["last_error"] = str(exc)
await asyncio.sleep(10) # backoff then reconnect
finally:
if consumer is not None:
try:
await consumer.stop()
except Exception:
pass
# ── Endpoints ────────────────────────────────────────────────────────────────
@router.get("")
async def list_changes(
limit: int = Query(100, le=500),
source: str | None = None,
op: str | None = None,
table: str | None = None,
) -> JSONResponse:
items = list(_ring)
if source:
items = [c for c in items if c["source"] == source]
if op:
items = [c for c in items if c["op"] == op]
if table:
items = [c for c in items if c["table"] == table]
items = list(reversed(items))[:limit]
return JSONResponse({
"ok": True,
"changes": items,
"buffered": len(_ring),
"connected": _state["connected"],
"consumed": _state["consumed"],
"last_error": _state["last_error"],
})
@router.get("/stats")
async def change_stats(minutes: int = Query(15, le=240)) -> JSONResponse:
now = datetime.now(timezone.utc)
cutoff = now.timestamp() - minutes * 60
by_source: dict[str, int] = {}
by_op: dict[str, int] = {}
by_table: dict[str, int] = {}
buckets: dict[str, int] = {}
total = 0
for c in _ring:
try:
t = datetime.fromisoformat(c["ts"]).timestamp()
except Exception:
continue
if t < cutoff:
continue
total += 1
by_source[c["source"]] = by_source.get(c["source"], 0) + 1
by_op[c["op"]] = by_op.get(c["op"], 0) + 1
by_table[c["table"]] = by_table.get(c["table"], 0) + 1
bucket = datetime.fromtimestamp(t, timezone.utc).strftime("%H:%M")
buckets[bucket] = buckets.get(bucket, 0) + 1
return JSONResponse({
"ok": True,
"window_minutes": minutes,
"total": total,
"by_source": by_source,
"by_op": by_op,
"by_table": by_table,
"buckets": [{"t": k, "n": v} for k, v in sorted(buckets.items())],
"connected": _state["connected"],
"consumed": _state["consumed"],
})
@router.get("/status")
async def changes_status() -> JSONResponse:
return JSONResponse({"ok": True, "buffered": len(_ring), **_state})