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