368 lines
11 KiB
Python
368 lines
11 KiB
Python
#!/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()
|