Didacticiel : clustering K-means
Dans ce tutoriel, vous créez un opérateur de fonction de table défini par l'utilisateur (UDTF) Python pour Lakeflow Designer qui exécute un clustering K-means avec scikit-learn. Les UDTF sont bien adaptées aux tâches de Machine Learning qui traitent des datasets entiers. Pour plus d'informations sur les opérateurs définis par l'utilisateur, consultez Opérateurs définis par l'utilisateur dans Lakeflow Designer.
Présentation
Ce tutoriel vous guide à travers la création d'un opérateur défini par l'utilisateur UDTF à l'aide de Python. L'opérateur effectue un clustering K-Means sur les colonnes sélectionnées, permettant aux utilisateurs de :
- Choisissez les colonnes à utiliser comme caractéristiques.
- Spécifiez le nombre de clusters.
- Obtenez un tableau avec les affectations de clusters pour chaque ligne.
Étape 1 : Comprendre le modèle de gestionnaire UDTF
Une UDTF est implémentée comme une classe Python avec trois méthodes clés :
__init__(): Appelé une fois avant le traitement des lignes pour initialiser l’état.eval(row, ...)** ** : Appelée pour chaque ligne d'entrée afin d'accumuler les données.terminate(): Appélé après que toutes les lignes ont été traitées pour produire des résultats.
Ce modèle permet à l'UDTF de :
- Collectez tous les points de données pendant
eval()appels - Entraîner le modèle K-Means dans
terminate() - Produire des résultats en cluster ligne par ligne
class SklearnKMeans:
def __init__(self):
self.id_col = None
self.feature_cols = None
self.k = None
self.rows = []
self.features = []
def eval(self, row, id_column, columns, k):
"""Called one time per input row - accumulate data here."""
# Initialize configuration on first row
if self.id_col is None:
self.id_col = id_column
if self.feature_cols is None:
self.feature_cols = columns
if self.k is None:
self.k = max(1, int(k))
# Convert row to dictionary and store
row_dict = row.asDict(recursive=False)
self.rows.append(row_dict)
# Extract numeric features
feats = []
for c in self.feature_cols:
v = row_dict.get(c)
if v is None:
v = 0.0
feats.append(float(v))
self.features.append(feats)
def terminate(self):
"""Called after all rows - train model and yield results."""
import numpy as np
from sklearn.cluster import KMeans
if not self.rows:
return
X = np.asarray(self.features, dtype=float)
n_samples = X.shape[0]
n_clusters = min(self.k, n_samples)
model = KMeans(
n_clusters=n_clusters,
n_init=10,
random_state=42
)
labels = model.fit_predict(X)
# Yield results row by row
for row_dict, label in zip(self.rows, labels):
yield str(row_dict[self.id_col]), int(label)
Le row paramètre dans eval() est un objet PySpark Row. Utilisez .asDict() pour le convertir en dictionnaire pour un accès plus facile.
Étape 2 : Créez le YAML pour l'opérateur
La configuration YAML définit la manière dont l'opérateur apparaît dans Lakeflow Designer. Pour cet opérateur :
- Paramètre numérique (
k) : nombre de clusters à créer - Sélectionner le widget (
id_column) : Liste déroulante renseignée avec les colonnes de la table d'entrée - **Widget
columnsà sélection multiple** () : sélection de plusieurs colonnes de fonctionnalités. optionsSource: Remplit automatiquement les listes déroulantes à partir du schéma de la table d'entrée- Port d'entrée : Spécifie que cet opérateur accepte les données tabulaires
schema: user-defined-operator-v0.1.0
type: uc-udtf
name: K-Means Clustering
id: kmeans
version: '1.0.0'
description: Perform K-Means clustering on selected columns
config:
type: object
properties:
k:
type: number
title: Number of Clusters
default: 3
minimum: 1
maximum: 100
x-ui:
widget: number
id_column:
type: string
title: ID Column
x-ui:
widget: select
optionsSource:
type: inputColumns
port: input_data
columns:
type: array
items:
type: string
title: Feature Columns
x-ui:
widget: multi-select
optionsSource:
type: inputColumns
port: input_data
required:
- k
- id_column
- columns
additionalProperties: false
ports:
input:
- name: input_data
title: Input Data
output:
- name: output
title: Clustered Data
Voir référence YAML de l'opérateur défini par l'utilisateur pour un guide complet de toutes les propriétés, types de données, widgets et options disponibles.
Étape 3 : créez la fonction Unity Catalog
Combinez la configuration YAML et la classe du gestionnaire Python en une seule instruction CREATE FUNCTION.
CREATE OR REPLACE FUNCTION main.my_schema.k_means(
input_data TABLE,
id_column STRING,
columns ARRAY<STRING>,
k INT
)
RETURNS TABLE (
id STRING,
cluster_id INT
)
LANGUAGE PYTHON
HANDLER 'SklearnKMeans'
AS $$
"""
schema: user-defined-operator-v0.1.0
type: uc-udtf
name: K-Means Clustering
id: kmeans
version: "1.0.0"
description: Perform K-Means clustering on selected columns
config:
type: object
properties:
k:
type: number
title: Number of Clusters
default: 3
minimum: 1
maximum: 100
x-ui:
widget: number
id_column:
type: string
title: ID Column
x-ui:
widget: select
optionsSource:
type: inputColumns
port: input_data
columns:
type: array
items:
type: string
title: Feature Columns
x-ui:
widget: multi-select
optionsSource:
type: inputColumns
port: input_data
required:
- k
- id_column
- columns
additionalProperties: false
ports:
input:
- name: input_data
title: Input Data
output:
- name: output
title: Clustered Data
"""
class SklearnKMeans:
def __init__(self):
self.id_col = None
self.feature_cols = None
self.k = None
self.rows = []
self.features = []
def eval(self, row, id_column, columns, k):
if self.id_col is None:
self.id_col = id_column
if self.feature_cols is None:
self.feature_cols = columns
if self.k is None:
self.k = max(1, int(k))
row_dict = row.asDict(recursive=False)
self.rows.append(row_dict)
feats = []
for c in self.feature_cols:
v = row_dict.get(c)
if v is None:
v = 0.0
feats.append(float(v))
self.features.append(feats)
def terminate(self):
import numpy as np
from sklearn.cluster import KMeans
if not self.rows:
return
X = np.asarray(self.features, dtype=float)
n_samples = X.shape[0]
n_clusters = min(self.k, n_samples)
model = KMeans(
n_clusters=n_clusters,
n_init=10,
random_state=42
)
labels = model.fit_predict(X)
for row_dict, label in zip(self.rows, labels):
yield str(row_dict[self.id_col]), int(label)
$$
Étape 4 : Tester avec des exemples de données
Créer un échantillon de données client pour les tests :
-- Create sample customer data
CREATE OR REPLACE TEMP VIEW customers AS
SELECT * FROM VALUES
('C001', 25, 35000, 20),
('C002', 45, 85000, 80),
('C003', 35, 55000, 50),
('C004', 50, 95000, 90),
('C005', 23, 30000, 15),
('C006', 40, 75000, 70),
('C007', 60, 100000, 95),
('C008', 30, 45000, 40)
AS t(customer_id, age, annual_income, spending_score);
Tester le K-Means UDTF :
-- Run K-Means clustering with 3 clusters
SELECT * FROM main.my_schema.k_means(
input_data => TABLE(SELECT * FROM customers) WITH SINGLE PARTITION,
k => 3,
id_column => 'customer_id',
columns => array('age', 'annual_income', 'spending_score')
)
Dans ce cas, vous souhaitez joindre les résultats du clustering avec les données d'origine pour voir les attributions de clusters :
-- Join cluster results with original data
SELECT
c.*,
k.cluster_id
FROM customers c
INNER JOIN main.my_schema.k_means(
input_data => TABLE(SELECT * FROM customers) WITH SINGLE PARTITION,
k => 3,
id_column => 'customer_id',
columns => array('age', 'annual_income', 'spending_score')
) k
ON c.customer_id = k.id
ORDER BY k.cluster_id, c.customer_id
Étape 5 : Enregistrez l'opérateur
Pour utiliser l'opérateur dans Lakeflow Designer, vous devez l'enregistrer en l'ajoutant à votre fichier .user_defined_operators.yaml :
operators:
- catalog: main
schema: my_schema
functionName: k_means
Si vous définissez ce fichier dans votre dossier utilisateur, il n'apparaît que pour vous. Pour plus d'informations, consultez Rendez votre opérateur visible.
Étape 6 : Configurez les autorisations
Accordez l’accès aux utilisateurs qui doivent utiliser cet opérateur :
GRANT USE SCHEMA ON SCHEMA main.my_schema TO `<user>`;
GRANT EXECUTE ON FUNCTION main.my_schema.k_means TO `<user>`;
Utilisez l'opérateur dans Lakeflow Designer
Une fois enregistré, l'opérateur apparaît dans Lakeflow Designer avec :
- Un port d'entrée pour connecter votre source de données
- Une liste déroulante pour sélectionner la colonne qui identifie de manière unique les lignes.
- Une sélection multiple pour choisir les colonnes à utiliser comme fonctionnalités de clustering
- Une entrée numérique pour le nombre de clusters souhaité
Les utilisateurs peuvent segmenter les clients, les produits ou toute autre donnée en groupes significatifs sans écrire de code.
Conseils pour la création d’UDTF
- Initialiser l'état dans
__init__: Configurez des listes/variables vides pour accumuler des données. - Accumuler dans
eval: Ne traitez pas encore, collectez simplement les données - Traiter en
terminate: c’est là que le vrai travail a lieu - Utilisez
yieldpour retourner les lignes : Retournez les résultats un par un determinate - Gérer les cas limites : Que faire s'il y a moins de lignes que de clusters ?
- Gardez les types explicites : les retours UDTF ne peuvent pas référencer les types d’entrée