Files

96 lines
6.8 KiB
Python
Raw Permalink Normal View History

2026-06-28 17:08:57 +02:00
#!/usr/bin/env python3
import json, sys, requests
BASE = "http://127.0.0.1:8088"
HOST = "10.0.21.50:8089"
USER = "mo"
SOURCES = [
("Trino · PostgreSQL Sales", f"trino://{USER}@{HOST}/postgres_sales/public", "public", "sales_orders", [
("PostgreSQL · Total Orders", "big_number_total", {"metric":{"expressionType":"SQL","sqlExpression":"COUNT(*)","label":"COUNT(*)"}}),
("PostgreSQL · Revenue by Region", "pie", {"metric":{"expressionType":"SQL","sqlExpression":"SUM(amount)","label":"Revenue"},"groupby":["region"],"row_limit":20}),
("PostgreSQL · Orders per Month", "echarts_timeseries_bar", {"metrics":[{"expressionType":"SQL","sqlExpression":"COUNT(*)","label":"Orders"}],"groupby":[{"expressionType":"SQL","sqlExpression":"date_trunc('month', order_ts)","label":"Month"}],"row_limit":24}),
]),
("Trino · MySQL HR", f"trino://{USER}@{HOST}/mysql_hr/hr", "hr", "employee_events", [
("MySQL HR · Total Events", "big_number_total", {"metric":{"expressionType":"SQL","sqlExpression":"COUNT(*)","label":"COUNT(*)"}}),
("MySQL HR · By Department", "pie", {"metric":{"expressionType":"SQL","sqlExpression":"COUNT(*)","label":"COUNT(*)"},"groupby":["department"],"row_limit":15}),
("MySQL HR · By Region", "echarts_timeseries_bar", {"metrics":[{"expressionType":"SQL","sqlExpression":"COUNT(*)","label":"Events"}],"groupby":["region"],"row_limit":20}),
]),
("Trino · MongoDB Supply Chain", f"trino://{USER}@{HOST}/mongodb_supplychain/supplychain", "supplychain", "events", [
("MongoDB · Total Events", "big_number_total", {"metric":{"expressionType":"SQL","sqlExpression":"COUNT(*)","label":"COUNT(*)"}}),
("MongoDB · By Type", "pie", {"metric":{"expressionType":"SQL","sqlExpression":"COUNT(*)","label":"COUNT(*)"},"groupby":["type"],"row_limit":10}),
("MongoDB · By Region", "echarts_timeseries_bar", {"metrics":[{"expressionType":"SQL","sqlExpression":"COUNT(*)","label":"Events"}],"groupby":["region"],"row_limit":10}),
]),
("Trino · Cassandra Telemetry", f"trino://{USER}@{HOST}/cassandra_telemetry/telemetry", "telemetry", "device_metrics", [
("Cassandra · Total Metrics", "big_number_total", {"metric":{"expressionType":"SQL","sqlExpression":"COUNT(*)","label":"COUNT(*)"}}),
("Cassandra · By Metric Type", "pie", {"metric":{"expressionType":"SQL","sqlExpression":"COUNT(*)","label":"COUNT(*)"},"groupby":["metric_type"],"row_limit":20}),
("Cassandra · Avg by Device", "echarts_timeseries_bar", {"metrics":[{"expressionType":"SQL","sqlExpression":"AVG(metric_value)","label":"Avg"}],"groupby":["device_id"],"row_limit":20}),
]),
]
OVERVIEW_SQL = """SELECT 'PostgreSQL sales_orders' AS source, COUNT(*) AS records FROM postgres_sales.public.sales_orders
UNION ALL SELECT 'MySQL employee_events', COUNT(*) FROM mysql_hr.hr.employee_events
UNION ALL SELECT 'MongoDB events', COUNT(*) FROM mongodb_supplychain.supplychain.events
UNION ALL SELECT 'Cassandra device_metrics', COUNT(*) FROM cassandra_telemetry.telemetry.device_metrics"""
def headers(s):
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 h
def get_db(s,h,name,uri):
r=s.get(f"{BASE}/api/v1/database/",headers=h); r.raise_for_status()
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,"allow_run_async":True})
if r.status_code not in (200,201): raise SystemExit(r.text)
print("DB",name,r.json()["id"]); return r.json()["id"]
def get_ds(s,h,db,schema,table,sql=None):
r=s.get(f"{BASE}/api/v1/dataset/",headers=h); r.raise_for_status()
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: return d["id"]
p={"database":db,"table_name":table,"schema":schema} if not sql else {"database":db,"table_name":table,"sql":sql}
r=s.post(f"{BASE}/api/v1/dataset/",headers=h,json=p)
if r.status_code not in (200,201): raise SystemExit(r.text)
print(" DS",table,r.json()["id"]); return r.json()["id"]
def mk_chart(s,h,name,ds,viz,params):
r=s.get(f"{BASE}/api/v1/chart/",headers=h); r.raise_for_status()
for c in r.json().get("result",[]):
if c.get("slice_name")==name: return c["id"]
p={"datasource":f"{ds}__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,"datasource_type":"table","params":json.dumps(p)})
if r.status_code not in (200,201): raise SystemExit(f"chart {name}: {r.text[:300]}")
print(" chart",name,r.json()["id"]); return r.json()["id"]
def mk_dash(s,h,title,cids):
layout={"DASHBOARD_VERSION":"v2","ROOT_ID":{"type":"ROOT","id":"ROOT_ID","children":["GRID_ID"]},"GRID_ID":{"type":"GRID","id":"GRID_ID","children":[],"parents":["ROOT_ID"]}}
row=col=0
for cid in cids:
k=f"CHART-{cid}"; x=(col%3)*4; y=row*12
layout[k]={"type":"CHART","id":k,"children":[],"meta":{"width":4,"height":10,"chartId":cid},"parents":["ROOT_ID","GRID_ID"]}
layout["GRID_ID"]["children"].append(k); col+=1
if col%3==0: row+=1
payload={"dashboard_title":title,"published":True,"position_json":json.dumps(layout),"json_metadata":"{}"}
r=s.get(f"{BASE}/api/v1/dashboard/",headers=h); r.raise_for_status()
for d in r.json().get("result",[]):
if d.get("dashboard_title")==title:
did=d["id"]; s.put(f"{BASE}/api/v1/dashboard/{did}",headers=h,json=payload)
for cid in cids: s.put(f"{BASE}/api/v1/chart/{cid}",headers=h,json={"dashboards":[did]})
print("Dashboard",did); return did
r=s.post(f"{BASE}/api/v1/dashboard/",headers=h,json=payload)
if r.status_code not in (200,201): raise SystemExit(r.text)
did=r.json()["id"]
for cid in cids: s.put(f"{BASE}/api/v1/chart/{cid}",headers=h,json={"dashboards":[did]})
print("Dashboard",did); return did
def main():
s=requests.Session(); h=headers(s)
cids=[]
odb=get_db(s,h,"Trino · Lakehouse Overview",f"trino://{USER}@{HOST}/postgres_sales/public")
ods=get_ds(s,h,odb,None,"lakehouse_counts",OVERVIEW_SQL)
cids.append(mk_chart(s,h,"Lakehouse · Records per Source",ods,"pie",{"metric":{"expressionType":"SQL","sqlExpression":"SUM(records)","label":"Records"},"groupby":["source"],"row_limit":10}))
for dbn,uri,sch,tbl,charts in SOURCES:
db=get_db(s,h,dbn,uri); ds=get_ds(s,h,db,sch,tbl)
for nm,viz,pr in charts: cids.append(mk_chart(s,h,nm,ds,viz,pr))
mk_dash(s,h,"ATC Lakehouse · Trino Federated",cids)
if __name__=="__main__": main()