From e2c64fd06801856a94a9088659951127b8e406f5 Mon Sep 17 00:00:00 2001 From: PM-pinou <2504420230@qq.com> Date: Fri, 24 Jul 2026 11:35:18 +0800 Subject: [PATCH] =?UTF-8?q?fix:=20raw=20row=20clustering=20=E2=80=94=20rem?= =?UTF-8?q?ove=20entity=20concept,=20use=20row=20index,=20all=20numeric=20?= =?UTF-8?q?features?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- analysis/views/auto.py | 2 +- analysis/views/clustering.py | 37 ++++++++++++++++++------------------ analysis/views/manual.py | 8 ++++---- 3 files changed, 23 insertions(+), 24 deletions(-) diff --git a/analysis/views/auto.py b/analysis/views/auto.py index 3244b13..527a77d 100644 --- a/analysis/views/auto.py +++ b/analysis/views/auto.py @@ -274,7 +274,7 @@ def run_llm_analysis_view(request): _run_clustering_pipeline( run=run, store=store, - entity_ds_id=entity_ds_id, + ds_id=entity_ds_id, feature_columns=None, algorithm='agglomerative', min_cluster_size=5, diff --git a/analysis/views/clustering.py b/analysis/views/clustering.py index f2c5798..e375071 100644 --- a/analysis/views/clustering.py +++ b/analysis/views/clustering.py @@ -196,16 +196,19 @@ def cluster_detail(request, display_id, cluster_label): @profile -def _run_clustering_pipeline(run, store, entity_ds_id, feature_columns, algorithm, +def _run_clustering_pipeline(run, store, ds_id, feature_columns, algorithm, min_cluster_size, run_umap=True, head=None): - """Unified clustering pipeline: clustering → feature extraction → UMAP → ORM save. + """Clustering pipeline on raw rows: clustering → feature extraction → UMAP → ORM save. + + Operates directly on raw data rows (no entity aggregation). + Per-row results stored in EntityProfile with row index as identifier. Args: run: AnalysisRun ORM object. store: SessionStore instance. - entity_ds_id: dataset ID in SessionStore pointing to entity-aggregated data. - feature_columns: list of feature column names (None = auto-detect). - algorithm: 'hdbscan' or 'kmeans'. + ds_id: dataset ID in SessionStore pointing to raw data. + feature_columns: list of feature column names (None = auto-detect all numeric). + algorithm: clustering algorithm. min_cluster_size: int for HDBSCAN. run_umap: bool, whether to compute and save UMAP-2D embeddings. head: optional int, downsample to at most this many rows (min 100). @@ -217,10 +220,10 @@ def _run_clustering_pipeline(run, store, entity_ds_id, feature_columns, algorith from analysis.tool_registry import _handle_run_clustering, _handle_extract_features import polars as pl - entry = store.get_dataset(entity_ds_id) + entry = store.get_dataset(ds_id) if entry is None: run.status = 'failed' - run.error_message = 'Entity dataset not found in session store' + run.error_message = '数据集在会话中不存在' run.save(update_fields=['status', 'error_message']) return @@ -241,11 +244,11 @@ def _run_clustering_pipeline(run, store, entity_ds_id, feature_columns, algorith run.save(update_fields=['progress_msg']) lf = lf.head(head) store.store_dataset( - entity_ds_id, lf, + ds_id, lf, schema=schema, metadata=entry.get('metadata', {}), ) - entry = store.get_dataset(entity_ds_id) # refresh + entry = store.get_dataset(ds_id) row_count = head # Resolve feature columns — use actual LazyFrame dtypes (not stored schema dict @@ -268,7 +271,7 @@ def _run_clustering_pipeline(run, store, entity_ds_id, feature_columns, algorith feature_cols = [ live_names[i] for i, dt in enumerate(live_dtypes) if dt in numeric_types and not live_names[i].startswith('_') - ][:10] + ] # ── Step 1: Clustering ──────────────────────────────────────────────── # If no numeric columns detected, pass empty list to let @@ -279,7 +282,7 @@ def _run_clustering_pipeline(run, store, entity_ds_id, feature_columns, algorith run.save(update_fields=['status', 'progress_pct', 'progress_msg']) result = asyncio.run(_handle_run_clustering( - dataset_id=entity_ds_id, + dataset_id=ds_id, cluster_columns=feature_cols, algorithm=algorithm, params={'min_cluster_size': min_cluster_size}, @@ -307,7 +310,7 @@ def _run_clustering_pipeline(run, store, entity_ds_id, feature_columns, algorith run.save(update_fields=['status', 'progress_pct', 'progress_msg']) feat_result = asyncio.run(_handle_extract_features( - dataset_id=entity_ds_id, + dataset_id=ds_id, cluster_result_id=cluster_id, top_k=10, method='zscore', save_to_db=True, run_id=run.id, @@ -431,10 +434,6 @@ def _run_clustering_pipeline(run, store, entity_ds_id, feature_columns, algorith if coords is None: raise Exception('UMAP produced no output') - non_num = [c for c in df_umap.columns - if c not in num_cols and not c.startswith('_')] - ent_col = non_num[0] if non_num else df_umap.columns[0] - cluster_entry = store.get_cluster_result(cluster_id) labels_list = cluster_entry.get('labels', []) if cluster_entry else [] @@ -443,9 +442,9 @@ def _run_clustering_pipeline(run, store, entity_ds_id, feature_columns, algorith profiles = [] cr_cache = {} - for i, row in enumerate(df_umap.iter_rows(named=True)): - ev = str(row.get(ent_col, '')) or 'entity' - ev = f'{ev}_{i}' + for i in range(len(df_umap)): + # Use row index as identifier — no entity column + ev = f'row_{i}' lbl = labels_list[i] if i < len(labels_list) else -1 cache_key = f'{run.id}_{lbl}' if cache_key not in cr_cache: diff --git a/analysis/views/manual.py b/analysis/views/manual.py index 6ea0ffd..1ef10eb 100644 --- a/analysis/views/manual.py +++ b/analysis/views/manual.py @@ -23,9 +23,9 @@ def manual_page(request): # Get all tools metadata tools = get_tools_meta() - core_names = {'profile_data', 'build_entity_profiles', 'compute_scores', - 'run_clustering', 'extract_features', 'detect_anomalies', - 'visualize_anomalies'} + core_names = {'profile_data', 'run_clustering', 'extract_features', + 'filter_data', 'preprocess_data', + 'detect_anomalies', 'visualize_anomalies'} diag_names = {'validate_data', 'explore_distributions', 'find_outliers', 'diagnose_clustering', 'compare_datasets', 'export_debug_sample', 'repair_schema'} @@ -165,7 +165,7 @@ def manual_run_analysis(request): # Use the unified clustering pipeline (clustering → extraction → UMAP) _run_clustering_pipeline( - run=run, store=store, entity_ds_id=ds_id, + run=run, store=store, ds_id=ds_id, feature_columns=feature_columns, algorithm=algorithm, min_cluster_size=min_cluster_size, run_umap=True, head=head,