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

La méthode __init__ s’exécute une fois avant que les lignes ne soient traitées. Elle configure l’état qui persiste lors des appels à eval() : la configuration (colonne d’ID, colonnes de caractéristiques, nombre de clusters) et les accumulateurs pour les données de lignes et les caractéristiques.

Python
class SklearnKMeans:
def __init__(self):
# State set here persists across eval() calls for use in terminate()
self.id_col = None
self.feature_cols = None
self.k = None
self.rows = []
self.features = []

eval() est appelée une fois par ligne d’entrée. Il capture la configuration de la première ligne, puis accumule chaque ligne et ses valeurs de caractéristiques numériques.

Python
    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)
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.

terminate() s’exécute après l’accumulation de toutes les lignes. Il entraîne le modèle K-Means sur les caractéristiques collectées et produit une ligne de sortie par ligne d’entrée, en associant l’ID de chaque ligne à son cluster assigné.

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

L’exemple crée la fonction dans main.example_output. Créez d’abord le schéma s’il n’existe pas :

SQL
CREATE SCHEMA IF NOT EXISTS main.example_output
SQL
CREATE OR REPLACE FUNCTION main.example_output.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 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);

Testez l'UDTF K-Means avec trois clusters :

SQL
SELECT * FROM main.example_output.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
SELECT
c.*,
k.cluster_id
FROM customers c
INNER JOIN main.example_output.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: example_output
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.example_output TO `<user>`;
GRANT EXECUTE ON FUNCTION main.example_output.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