""" ClickHouse -> public/data/metrics.json Writes a labelled document, not a bare array: { "source": "clickhouse" | "demo", "fetchedAt": "...", "table": "...", "windowDays": 30, "events": [ { eventKey, totalEvents, uniqueUsers, ... } ] } The `source` label is the point. Every number the simulator shows is tagged with where it came from, and hotspots whose event key is absent from this file get no metrics at all rather than a plausible-looking guess. Configuration (environment, or a .env file next to this project): CLICKHOUSE_HOST=http://10.0.0.5:8123 CLICKHOUSE_USER=default CLICKHOUSE_PASSWORD=... CLICKHOUSE_DB=default CLICKHOUSE_TABLE=amplitude_events """ import argparse import json import os import sys import urllib.error import urllib.parse import urllib.request from datetime import datetime, timezone from pathlib import Path sys.path.insert(0, str(Path(__file__).resolve().parent)) import telecom_cdp as T def load_dotenv(): env_path = T.PROJECT_ROOT / ".env" if not env_path.exists(): return for line in env_path.read_text(encoding="utf-8").splitlines(): line = line.strip() if not line or line.startswith("#") or "=" not in line: continue key, _, value = line.partition("=") os.environ.setdefault(key.strip(), value.strip().strip('"').strip("'")) load_dotenv() HOST = os.getenv("CLICKHOUSE_HOST", "") USER = os.getenv("CLICKHOUSE_USER", "default") PASSWORD = os.getenv("CLICKHOUSE_PASSWORD", "") DATABASE = os.getenv("CLICKHOUSE_DB", "default") TABLE = os.getenv("CLICKHOUSE_TABLE", "amplitude_events") class ClickHouseError(RuntimeError): pass def run_query(sql, params=None, host=None, timeout=60): """ Execute SQL over the HTTP interface using ClickHouse server-side parameters (param_ + {name:Type} placeholders) so no user value is ever pasted into the query text. """ host = host or HOST if not host: raise ClickHouseError( "CLICKHOUSE_HOST is not set. Put it in .env or the environment, e.g. " "CLICKHOUSE_HOST=http://10.0.0.5:8123" ) query = {"query": sql + " FORMAT JSON", "database": DATABASE} for key, value in (params or {}).items(): query["param_" + key] = str(value) request = urllib.request.Request( host.rstrip("/") + "/?" + urllib.parse.urlencode(query), headers={"X-ClickHouse-User": USER, "X-ClickHouse-Key": PASSWORD}, ) try: with urllib.request.urlopen(request, timeout=timeout) as response: return json.loads(response.read().decode("utf-8")).get("data", []) except urllib.error.HTTPError as exc: raise ClickHouseError( "ClickHouse returned HTTP " + str(exc.code) + ": " + exc.read().decode("utf-8", "replace")[:400] ) from exc except Exception as exc: raise ClickHouseError("Cannot reach ClickHouse at " + host + ": " + str(exc)) from exc def fetch_event_metrics(table=None, window_days=30): table = table or TABLE if not table.replace("_", "").replace(".", "").isalnum(): raise ClickHouseError("Refusing to query a non-identifier table name: " + repr(table)) # `shareOfClicks` is each event's share of all click events in the window, so # the number in the UI has a defined meaning instead of being decorative. sql = ( "SELECT event_type AS eventKey," " count() AS totalEvents," " uniqExact(user_id) AS uniqueUsers," " uniqExact(device_id) AS uniqueDevices," " round(count() / nullIf(uniqExact(user_id), 0), 2) AS avgEventsPerUser," " round(100 * count() / nullIf(sum(count()) OVER (), 0), 2) AS shareOfClicks" " FROM " + table + " WHERE event_time >= now() - INTERVAL {window:UInt32} DAY" " GROUP BY event_type" " ORDER BY totalEvents DESC" ) return run_query(sql, {"window": window_days}) def fetch_user_session(user_id, table=None, limit=200): table = table or TABLE sql = ( "SELECT event_type AS eventKey, event_time AS timestamp, event_properties AS properties" " FROM " + table + " WHERE user_id = {uid:String} OR device_id = {uid:String}" " ORDER BY event_time ASC LIMIT {lim:UInt32}" ) return run_query(sql, {"uid": user_id, "lim": limit}) def write_metrics(events, source, table, window_days): T.DATA_DIR.mkdir(parents=True, exist_ok=True) document = { "source": source, "fetchedAt": datetime.now(timezone.utc).isoformat(timespec="seconds"), "table": table, "windowDays": window_days, "events": events, } tmp = T.METRICS_PATH.with_suffix(".json.tmp") with open(tmp, "w", encoding="utf-8") as f: json.dump(document, f, ensure_ascii=False, indent=2) tmp.replace(T.METRICS_PATH) return document def main(): parser = argparse.ArgumentParser(description="Fetch Amplitude event metrics from ClickHouse.") parser.add_argument("--table", default=TABLE) parser.add_argument("--days", type=int, default=30) parser.add_argument("--user", default=None, help="Print one user's event timeline and exit") args = parser.parse_args() if args.user: try: for row in fetch_user_session(args.user, args.table): print(json.dumps(row, ensure_ascii=False)) except ClickHouseError as exc: print(json.dumps({"success": False, "error": str(exc)}, ensure_ascii=False)) return 1 return 0 try: events = fetch_event_metrics(args.table, args.days) except ClickHouseError as exc: print( json.dumps( { "success": False, "error": str(exc), "hint": "metrics.json was left untouched; the UI keeps showing its current source label.", }, ensure_ascii=False, ) ) return 1 doc = write_metrics(events, "clickhouse", args.table, args.days) print( json.dumps( {"success": True, "source": doc["source"], "events": len(events), "table": args.table}, ensure_ascii=False, ) ) return 0 if __name__ == "__main__": sys.exit(main())