Aller au contenu principal

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 :

  1. Collectez tous les points de données pendant eval() appels
  2. Entraîner le modèle K-Means dans terminate()
  3. Produire des résultats en cluster ligne par ligne
Python
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)
remarque

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
YAML
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.

SQL
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 :

SQL
-- 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 :

SQL
-- 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 :

SQL
-- 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 :

YAML
operators:
- catalog: main
schema: my_schema
functionName: k_means
remarque

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 :

SQL
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

  1. Initialiser l'état dans __init__ : Configurez des listes/variables vides pour accumuler des données.
  2. Accumuler dans eval : Ne traitez pas encore, collectez simplement les données
  3. Traiter en terminate : c’est là que le vrai travail a lieu
  4. Utilisez yield pour retourner les lignes : Retournez les résultats un par un de terminate
  5. Gérer les cas limites : Que faire s'il y a moins de lignes que de clusters ?
  6. Gardez les types explicites : les retours UDTF ne peuvent pas référencer les types d’entrée