ishaq101's picture sofhiaazzhr's picture
/feat ds tools (#22)
e201ae4
Raw History Blame Contribute Delete
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],
}