#!/usr/bin/env python3 """Create Superset dashboard for Debezium, Kafka CDC changes, and Spark.""" import json import os import requests BASE = "http://127.0.0.1:8088" DASH_TITLE = "ATC Lakehouse · Pipeline & CDC" CSS_PATH = "/tmp/palantir_dashboard.css" MONITOR_URI = "trino://mo@10.0.21.50:8089/postgres_sales/monitor" KAFKA_URI = "trino://mo@10.0.21.50:8089/kafka/default" def session(): s = requests.Session() r = s.post( f"{BASE}/api/v1/security/login", json={"username": "admin", "password": "admin", "provider": "db", "refresh": True}, ) r.raise_for_status() h = {"Authorization": "Bearer " + r.json()["access_token"], "Content-Type": "application/json"} h["X-CSRFToken"] = s.get(f"{BASE}/api/v1/security/csrf_token/", headers=h).json()["result"] h["Referer"] = BASE return s, h def get_db(s, h, name, uri): r = s.get(f"{BASE}/api/v1/database/", headers=h) for d in r.json().get("result", []): if d["database_name"] == name: return d["id"] r = s.post( f"{BASE}/api/v1/database/", headers=h, json={"database_name": name, "sqlalchemy_uri": uri, "expose_in_sqllab": True}, ) r.raise_for_status() return r.json()["id"] def ds_table(s, h, db_id, schema, table): r = s.get(f"{BASE}/api/v1/dataset/", headers=h) for d in r.json().get("result", []): if d.get("table_name") == table and d.get("schema") == schema and d.get("database", {}).get("id") == db_id: return d["id"] r = s.post( f"{BASE}/api/v1/dataset/", headers=h, json={"database": db_id, "schema": schema, "table_name": table}, ) r.raise_for_status() return r.json()["id"] def ds_sql(s, h, db_id, name, sql): r = s.get(f"{BASE}/api/v1/dataset/", headers=h) for d in r.json().get("result", []): if d.get("table_name") == name: return d["id"] r = s.post( f"{BASE}/api/v1/dataset/", headers=h, json={"database": db_id, "table_name": name, "sql": sql}, ) r.raise_for_status() return r.json()["id"] def chart(s, h, name, ds_id, viz, params): r = s.get(f"{BASE}/api/v1/chart/", headers=h) for c in r.json().get("result", []): if c.get("slice_name") == name: return c["id"] p = {"datasource": f"{ds_id}__table", "viz_type": viz, "row_limit": 1000, **params} r = s.post( f"{BASE}/api/v1/chart/", headers=h, json={ "slice_name": name, "viz_type": viz, "datasource_id": ds_id, "datasource_type": "table", "params": json.dumps(p), "owners": [1, 2, 3], }, ) r.raise_for_status() return r.json()["id"] def layout(items): L = { "DASHBOARD_VERSION": "v2", "ROOT_ID": {"type": "ROOT", "id": "ROOT_ID", "children": ["GRID_ID"]}, "GRID_ID": {"type": "GRID", "id": "GRID_ID", "children": [], "parents": ["ROOT_ID"]}, } ri = 0 def row(): nonlocal ri ri += 1 rid = f"ROW-{ri}" L["GRID_ID"]["children"].append(rid) L[rid] = { "type": "ROW", "id": rid, "children": [], "parents": ["ROOT_ID", "GRID_ID"], "meta": {"background": "BACKGROUND_TRANSPARENT"}, } return rid for item in items: if item[0] == "md": rid = row() mid = f"MD-{ri}" L[rid]["children"].append(mid) L[mid] = { "type": "MARKDOWN", "id": mid, "children": [], "parents": ["ROOT_ID", "GRID_ID", rid], "meta": {"width": 12, "height": 10, "code": f"## {item[1]}\n\n{item[2]}"}, } else: cid, name = item rid = row() charts_in_row = [k for k in L[rid]["children"] if k.startswith("CHART-")] if len(charts_in_row) >= 3: rid = row() key = f"CHART-{cid}" L[rid]["children"].append(key) hgt = 70 if "Recent" in name or "table" in name.lower() else 55 L[key] = { "type": "CHART", "id": key, "children": [], "parents": ["ROOT_ID", "GRID_ID", rid], "meta": {"width": 4 if "Recent" not in name else 12, "height": hgt, "chartId": cid, "sliceName": name}, } return L def main(): s, h = session() db_mon = get_db(s, h, "Trino · Pipeline Monitor", MONITOR_URI) db_kfk = get_db(s, h, "Trino · Kafka CDC", KAFKA_URI) ds_conn = ds_table(s, h, db_mon, "monitor", "debezium_connectors") ds_topics = ds_table(s, h, db_mon, "monitor", "kafka_topics") ds_ops = ds_table(s, h, db_mon, "monitor", "cdc_operations") ds_recent = ds_table(s, h, db_mon, "monitor", "cdc_recent_events") ds_spark = ds_table(s, h, db_mon, "monitor", "spark_applications") LIVE_CDC_SQL = """ SELECT 'PostgreSQL' AS source_system, CASE json_extract_scalar(_message, '$.payload.op') WHEN 'c' THEN 'INSERT' WHEN 'u' THEN 'UPDATE' WHEN 'd' THEN 'DELETE' WHEN 'r' THEN 'SNAPSHOT' ELSE 'OTHER' END AS change_type, json_extract_scalar(_message, '$.payload.source.table') AS table_name, json_extract_scalar(_message, '$.payload.after.order_id') AS record_key, _timestamp AS event_time FROM kafka.default."postgres-sales.public.sales_orders" WHERE _timestamp > current_timestamp - INTERVAL '7' DAY LIMIT 500 """ ds_live = ds_sql(s, h, db_kfk, "live_cdc_postgres_sample", LIVE_CDC_SQL) charts = [] charts.append(chart(s, h, "Debezium · Connector Status", ds_conn, "table", {})) charts.append( chart( s, h, "Debezium · RUNNING vs State", ds_conn, "pie", { "metric": {"expressionType": "SQL", "sqlExpression": "COUNT(*)", "label": "COUNT(*)"}, "groupby": ["state"], }, ) ) charts.append( chart( s, h, "Kafka · Topic Offsets", ds_topics, "echarts_timeseries_bar", { "metrics": [ {"expressionType": "SQL", "sqlExpression": "SUM(end_offset)", "label": "Messages"} ], "groupby": ["topic"], "row_limit": 20, }, ) ) charts.append( chart( s, h, "Kafka · Partitions per Topic", ds_topics, "echarts_timeseries_bar", { "metrics": [ {"expressionType": "SQL", "sqlExpression": "SUM(end_offset)", "label": "Offset"} ], "groupby": ["topic", "partition_id"], "row_limit": 30, }, ) ) charts.append( chart( s, h, "CDC · Changes by Type", ds_ops, "pie", { "metric": {"expressionType": "SQL", "sqlExpression": "SUM(event_count)", "label": "Events"}, "groupby": ["operation_label"], }, ) ) charts.append( chart( s, h, "CDC · Changes per Source", ds_ops, "echarts_timeseries_bar", { "metrics": [ {"expressionType": "SQL", "sqlExpression": "SUM(event_count)", "label": "Events"} ], "groupby": ["source_system", "operation_label"], "row_limit": 20, }, ) ) charts.append( chart( s, h, "CDC · Recent Changes (sampled)", ds_recent, "table", { "all_columns": [ "source_system", "operation_label", "table_name", "record_key", "detail", "event_ts", ], "row_limit": 50, }, ) ) charts.append( chart( s, h, "CDC · Live Stream Sample (PostgreSQL)", ds_live, "table", { "all_columns": ["source_system", "change_type", "table_name", "record_key", "event_time"], "row_limit": 100, }, ) ) charts.append( chart( s, h, "Spark · Applications", ds_spark, "table", {"all_columns": ["app_id", "app_name", "state", "cores", "memory_mb", "duration_sec"]}, ) ) charts.append( chart( s, h, "Pipeline · Total Kafka Messages", ds_topics, "big_number_total", { "metric": { "expressionType": "SQL", "sqlExpression": "SUM(end_offset)", "label": "Total Offset", } }, ) ) chart_specs = [ ("md", "Pipeline & Change Data Capture", "Debezium → Kafka → Spark · Live CDC visibility"), ("md", "Debezium Connect", "Connector health on kafka01 :8083"), (charts[0], "Debezium · Connector Status"), (charts[1], "Debezium · RUNNING vs State"), ("md", "Apache Kafka", "Topic volume & CDC streams on kafka01"), (charts[2], "Kafka · Topic Offsets"), (charts[3], "Kafka · Partitions per Topic"), (charts[9], "Pipeline · Total Kafka Messages"), ("md", "Data Changes (CDC)", "INSERT / UPDATE / DELETE / SNAPSHOT — gewijzigde data"), (charts[4], "CDC · Changes by Type"), (charts[5], "CDC · Changes per Source"), (charts[6], "CDC · Recent Changes (sampled)"), (charts[7], "CDC · Live Stream Sample (PostgreSQL)"), ("md", "Apache Spark", "Batch & streaming jobs · lake01:8080"), (charts[8], "Spark · Applications"), ] items = chart_specs cids = [c[0] for c in chart_specs if c[0] != "md"] css = open(CSS_PATH).read() if os.path.exists(CSS_PATH) else "" payload = { "dashboard_title": DASH_TITLE, "published": True, "position_json": json.dumps(layout(items)), "css": css, "json_metadata": json.dumps( { "color_scheme": "palantir_ops", "refresh_frequency": 120, "chart_configuration": { str(c): {"id": c, "crossFilters": {"scope": "global", "chartsInScope": cids}} for c in cids }, } ), "owners": [1, 2, 3], } r = s.get(f"{BASE}/api/v1/dashboard/", headers=h) dash_id = None for d in r.json().get("result", []): if d.get("dashboard_title") == DASH_TITLE: dash_id = d["id"] break if dash_id: r = s.put(f"{BASE}/api/v1/dashboard/{dash_id}", headers=h, json=payload) else: r = s.post(f"{BASE}/api/v1/dashboard/", headers=h, json=payload) dash_id = r.json()["id"] print("Dashboard", dash_id, r.status_code) for cid in cids: s.put(f"{BASE}/api/v1/chart/{cid}", headers=h, json={"dashboards": [dash_id], "owners": [1, 2, 3]}) # query contexts os.system("python3 /tmp/fix_charts_qc.py 2>/dev/null || true") print(f"URL: {BASE}/superset/dashboard/{dash_id}/") if __name__ == "__main__": main()