Download src/tools/analytics/cluster.py from DataEyond/Agentic-Service-Data-Eyond-Catalog: direct link, hf CLI and curl.
- Browser
- Download file 10.8 kB
-
https://huggingface.co/spaces/DataEyond/Agentic-Service-Data-Eyond-Catalog/resolve/main/src/tools/analytics/cluster.py
- Command line
-
hf download hf://spaces/DataEyond/Agentic-Service-Data-Eyond-Catalog/src/tools/analytics/cluster.py
-
curl -L -o cluster.py https://huggingface.co/spaces/DataEyond/Agentic-Service-Data-Eyond-Catalog/resolve/main/src/tools/analytics/cluster.py
10.8 kB
| """analyze_cluster — discover groups that aren't already known (§0.9). | |
| Fourth ML analytical tool. Scales the selected numeric columns, fits KMeans for every | |
| k in a small range, keeps the k with the best silhouette, and returns one profile row | |
| per cluster plus a scatter coloured by cluster ("sepaket"). | |
| **Cluster vs segmentation — the distinction the 2026-09-15 discussion drew, and the | |
| reason both tools exist.** Segmentation is *by rules*: the criteria are already known | |
| ("split spend into tiers"), so the buckets are deterministic. Clustering is for when | |
| *"kita enggak tahu sama sekali"* — the groups are discovered from the data. Picking the | |
| wrong one of the two is the failure mode the tool descriptions have to prevent, which | |
| is why `dont_use_when` names the other explicitly on both sides. | |
| **k is chosen, not guessed.** Silhouette is computed for each candidate k and the best | |
| wins; the score is returned so the answer can say how well-separated the clusters | |
| actually are, and a weak score is a caveat rather than a hidden footnote. Clusters are | |
| descriptive — they describe how rows group, not why. | |
| Follows docs/tools/ANALYTICAL_TOOLS_CONVENTIONS.md: compute is synchronous except for | |
| the `await run_fit` offload; `data` is never self-fetched (Pattern A); errors escape to | |
| the invoker's never-throw seam; output is deterministic (fixed random_state). Mirrors | |
| analyze_anomaly's shape so the ML tools stay uniform. | |
| """ | |
| from __future__ import annotations | |
| import numpy as np | |
| import pandas as pd | |
| from pydantic import BaseModel, Field | |
| from src.tools.analytics.charts import cluster_scatter_chart | |
| from src.tools.analytics.descriptive import ColumnNotFoundError | |
| from src.tools.analytics.temporal import _clean | |
| from src.tools.ml.executor import run_fit | |
| from src.tools.spec_builder import build_description | |
| # KMeans on a handful of rows describes noise, not structure. | |
| _MIN_ROWS = 20 | |
| # Candidate k range (§0.9): small enough to stay interpretable in a narrative. | |
| _K_MIN = 2 | |
| _K_MAX = 8 | |
| # Feature cap (§0.9 R-5), consistent with analyze_importance. | |
| _MAX_FEATURES = 20 | |
| # Fixed so the same data always yields the same clusters (family-B determinism). | |
| _RANDOM_STATE = 42 | |
| # A silhouette at or below this means the groups barely separate — say so. | |
| _WEAK_SILHOUETTE = 0.25 | |
| # Chart column name (internal; never collides with user data). | |
| _CLUSTER_COL = "cluster" | |
| class NotEnoughDataError(ValueError): | |
| """Too few complete rows to cluster (error_code NOT_ENOUGH_DATA).""" | |
| class EmptyColumnSelectionError(ValueError): | |
| """column_ids was empty (error_code EMPTY_COLUMN_SELECTION).""" | |
| class ClusterInput(BaseModel): | |
| # `data` FIRST — property order follows field order (planner-prompt readability). | |
| data: str = Field(..., description="Placeholder ${t<id>} of the upstream table (Pattern A).") | |
| column_ids: list[str] = Field( | |
| ..., description="NUMERIC columns describing each row — the traits to group on." | |
| ) | |
| k: int | None = Field( | |
| None, | |
| description=f"Number of clusters; omit to choose automatically between " | |
| f"{_K_MIN} and {_K_MAX} by silhouette score.", | |
| ) | |
| DESCRIPTION = build_description( | |
| summary="Discover natural groups in the rows and profile each one.", | |
| use_when="the user asks what kinds/groups/segments exist WITHOUT saying how to " | |
| "split them — 'ada berapa kelompok pelanggan', 'what types of units do we have', " | |
| "'kelompokkan berdasarkan pola pemakaian'.", | |
| dont_use_when=[ | |
| "the user already states the split rule ('tier by spend', 'bucket by age') " | |
| "-> analyze_segmentation (that is by rules; this discovers the groups)", | |
| "the user wants unusual individual rows -> analyze_anomaly", | |
| "the user wants totals per an EXISTING category column -> analyze_aggregate", | |
| ], | |
| output="the chosen k with its silhouette score, one profile row per cluster " | |
| "(size, share, per-column means), and a scatter coloured by cluster.", | |
| examples=[ | |
| "ada berapa kelompok pelanggan berdasarkan demografi?", | |
| "what natural groups are in this data?", | |
| "kelompokkan unit berdasarkan pola pemakaiannya", | |
| ], | |
| ) | |
| def _fit_kmeans(x: np.ndarray, k_values: list[int], random_state: int): | |
| """Scale once, fit KMeans for each candidate k, return the best by silhouette. | |
| Imported lazily so app startup doesn't pay for scikit-learn unless a clustering | |
| analysis runs. Runs inside run_fit's thread — this is the whole search, so the | |
| per-fit timeout covers every candidate k, not just one. | |
| Returns (labels, best_k, best_score, {k: score}). | |
| """ | |
| from sklearn.cluster import KMeans | |
| from sklearn.metrics import silhouette_score | |
| from sklearn.preprocessing import StandardScaler | |
| # Unsupervised: there is no holdout, so scaling on all rows leaks nothing. | |
| x_scaled = StandardScaler().fit_transform(x) | |
| scores: dict[int, float] = {} | |
| best = (None, 0, -1.0) # (labels, k, score) | |
| for k in k_values: | |
| labels = KMeans(n_clusters=k, n_init=10, random_state=random_state).fit_predict(x_scaled) | |
| if len(set(labels)) < 2: # degenerate fit — silhouette is undefined | |
| continue | |
| score = float(silhouette_score(x_scaled, labels)) | |
| scores[k] = round(score, 4) | |
| if score > best[2]: | |
| best = (labels, k, score) | |
| if best[0] is None: | |
| # Every candidate collapsed to one cluster: fall back to the smallest k so the | |
| # caller still gets a grouping, with the weak-separation caveat attached. | |
| labels = KMeans(n_clusters=_K_MIN, n_init=10, random_state=random_state).fit_predict( | |
| x_scaled | |
| ) | |
| return labels, _K_MIN, 0.0, scores | |
| return best[0], best[1], best[2], scores | |
| def _profiles(frame: pd.DataFrame, labels: np.ndarray, column_ids: list[str]) -> list[dict]: | |
| """One row per cluster: size, share of rows, and the mean of each column.""" | |
| total = len(labels) | |
| out: list[dict] = [] | |
| for label in sorted(set(labels.tolist())): | |
| members = frame[labels == label] | |
| profile: dict[str, object] = { | |
| "cluster": int(label), | |
| "size": int(len(members)), | |
| "share": round(len(members) / total, 4) if total else 0.0, | |
| } | |
| for col in column_ids: | |
| profile[col] = _clean(members[col].mean()) | |
| out.append(profile) | |
| return out | |
| async def analyze_cluster( | |
| df: pd.DataFrame, | |
| column_ids: list[str], | |
| k: int | None = None, | |
| ) -> dict[str, object]: | |
| """Group rows by `column_ids` with KMeans, choosing k by silhouette when unset. | |
| Async ONLY so the fit search can be offloaded via `run_fit`; everything else is | |
| synchronous compute. `data` is materialised upstream (Pattern A) and arrives as | |
| `df`. | |
| Returns a dict with: method, column_ids, n_rows, k, requested_k (what the caller | |
| asked for; differs from k only when a fallback overrode an explicit k — D-10), | |
| `quality` (silhouette + floor + weak flag), silhouette_by_k, | |
| `clusters` (one profile per group), `caveats`, and `charts` (the reserved | |
| auto-chart key). | |
| Raises (all wrapped by the invoker's never-throw seam): | |
| ColumnNotFoundError, EmptyColumnSelectionError, NotEnoughDataError. | |
| """ | |
| if not column_ids: | |
| raise EmptyColumnSelectionError("column_ids must name at least one numeric column") | |
| missing = [c for c in column_ids if c not in df.columns] | |
| if missing: | |
| raise ColumnNotFoundError(f"columns not found: {missing}") | |
| capped = column_ids[:_MAX_FEATURES] | |
| features = df[capped].apply(pd.to_numeric, errors="coerce") | |
| frame = features.dropna() | |
| if frame.shape[0] < _MIN_ROWS: | |
| raise NotEnoughDataError( | |
| f"need >= {_MIN_ROWS} rows with numeric values in all of {capped} " | |
| f"to cluster, got {frame.shape[0]}" | |
| ) | |
| # Never ask for more clusters than rows, and honour an explicit k as-is. | |
| ceiling = min(_K_MAX, frame.shape[0] - 1) | |
| k_values = [k] if k else list(range(_K_MIN, max(_K_MIN, ceiling) + 1)) | |
| labels, chosen_k, silhouette, by_k = await run_fit( | |
| _fit_kmeans, frame.to_numpy(dtype=float), k_values, _RANDOM_STATE, label="analyze_cluster" | |
| ) | |
| # D-10: an explicit k that every candidate collapsed on is silently replaced by the | |
| # fallback's _K_MIN. Detect it (explicit k, but a different k came back) and say so, | |
| # instead of the old line that claimed the fallback's k "was given". | |
| overridden = k is not None and chosen_k != k | |
| caveats = [ | |
| "Clusters are descriptive — they show how rows group together, not why, and " | |
| "they are not business categories until someone names them.", | |
| ] | |
| if overridden: | |
| caveats.append( | |
| f"Requested k={k} could not form {k} separable clusters (every candidate " | |
| f"collapsed to a single group), so it fell back to k={chosen_k}; " | |
| f"silhouette = {silhouette:.3f}." | |
| ) | |
| else: | |
| caveats.append( | |
| f"k={chosen_k} " + ("was given" if k else "chosen by silhouette") | |
| + f"; silhouette = {silhouette:.3f} (1.0 = perfectly separated)." | |
| ) | |
| if silhouette <= _WEAK_SILHOUETTE: | |
| caveats.append( | |
| "Separation is weak — these groups overlap substantially, so treat the " | |
| "split as tentative." | |
| ) | |
| if len(column_ids) > _MAX_FEATURES: | |
| caveats.append( | |
| f"Only the first {_MAX_FEATURES} of {len(column_ids)} columns were used." | |
| ) | |
| chart_df = frame.assign(**{_CLUSTER_COL: labels}) | |
| # Two real columns, not PCA components: a user can read "PA vs MTTR" and cannot | |
| # read "PC1 vs PC2". With a single column, plot it against the row index. | |
| if len(capped) >= 2: | |
| chart = cluster_scatter_chart(chart_df, capped[0], capped[1], _CLUSTER_COL) | |
| else: | |
| chart_df = chart_df.assign(__row__=range(len(chart_df))) | |
| chart = cluster_scatter_chart(chart_df, "__row__", capped[0], _CLUSTER_COL) | |
| return { | |
| "method": "kmeans", | |
| "column_ids": capped, | |
| "n_rows": int(frame.shape[0]), | |
| "k": int(chosen_k), | |
| # What the caller asked for (None = auto) vs what `k` ended up — they differ only | |
| # when the fallback overrode an explicit k (D-10). | |
| "requested_k": k, | |
| # Reserved `quality` block (conventions §3) — see analyze_importance. | |
| "quality": { | |
| "metric": "silhouette", | |
| "value": round(float(silhouette), 4), | |
| "floor": _WEAK_SILHOUETTE, | |
| "weak": silhouette <= _WEAK_SILHOUETTE, | |
| }, | |
| "silhouette_by_k": by_k, | |
| "clusters": _profiles(frame, labels, capped), | |
| "caveats": caveats, | |
| "charts": [chart], | |
| } | |