Tutoriel : clustering k-moyennes

Dans ce tutoriel, vous créez pour Lakeflow Designer un opérateur de fonction de table définie par l’utilisateur (UDTF) en Python qui exécute un clustering K-means avec scikit-learn. Les UDTFs sont particulièrement adaptées aux tâches d’apprentissage automatique qui traitent des ensembles de données entiers. Pour plus d’informations sur les opérateurs définis par l’utilisateur, consultez Les opérateurs définis par l’utilisateur dans Lakeflow Designer.

Vue d’ensemble

Ce tutoriel vous guide tout au long de la création d’un opérateur défini par l’utilisateur UDTF à l’aide de Python. L’opérateur réalise un partitionnement K-means sur des colonnes sélectionnées, ce qui permet aux utilisateurs de :

  • Choisissez les colonnes à utiliser en tant que fonctionnalités.
  • Spécifiez le nombre de clusters.
  • Renvoie une table avec les affectations de cluster pour chaque ligne.

Étape 1 : Comprendre le modèle de gestionnaire UDTF

Une UDTF est implémentée en tant que classe Python avec trois méthodes clés :

  • __init__(): appelé une fois avant que toutes les lignes soient traitées pour initialiser l’état.
  • eval(row, ...): appelé pour chaque ligne d’entrée pour accumuler des données.
  • terminate(): appelé après que toutes les lignes soient traitées pour générer des résultats.

Ce modèle permet à l’UDTF de :

  1. Collectez tous les points de données lors des appels eval()
  2. Entraîner le modèle K-Means avec terminate()
  3. Cesser temporairement l'exécution de résultats 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)

Note

Le row paramètre dans eval() est un objet PySpark Row. Permet .asDict() de le convertir en dictionnaire pour faciliter l’accès.

Étape 2 : Créer le YAML pour l’opérateur

La configuration YAML définit la façon dont l’opérateur apparaît dans Lakeflow Designer. Pour cet opérateur :

  • Paramètre number (k) : nombre de clusters à créer
  • Sélectionner un widget (id_column) : liste déroulante remplie avec des colonnes de la table d’entrée
  • Widget à sélection multiple (columns) : sélection de plusieurs colonnes de fonctionnalité
  • optionsSource: remplit automatiquement les listes déroulantes du schéma de 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

Consultez la référence YAML de l’opérateur défini par l’utilisateur pour obtenir un guide complet sur toutes les propriétés, types de données, widgets et options disponibles.

Étape 3 : Créer la fonction Catalogue Unity

Regroupez la configuration YAML et la classe de 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éez des exemples 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);

Testez l’UDTF K-Means :

-- 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 de clustering avec les données d’origine pour afficher les affectations de cluster :

-- 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 : Inscrire l’opérateur

Pour utiliser l’opérateur dans Lakeflow Designer, vous devez l’inscrire en l’ajoutant à votre .user_defined_operators.yaml fichier :

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

Note

Si vous définissez ce fichier dans votre dossier utilisateur, il s’affiche uniquement pour vous. Pour plus d’informations, consultez Rendre votre opérateur détectable.

Étape 6 : Configurer 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>`;

Utiliser l’opérateur dans Lakeflow Designer

Une fois inscrit, l’opérateur apparaît dans Lakeflow Designer avec :

  • Port d’entrée pour connecter votre source de données
  • Liste déroulante pour sélectionner la colonne qui identifie de manière unique les lignes
  • Sélection multiple pour choisir les colonnes à utiliser comme fonctionnalités de clustering
  • Entrée numérique pour le nombre souhaité de clusters

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 de UDTF

  1. Initialiser l’état dans __init__: configurer des listes/variables vides pour accumuler des données
  2. Accumuler dans eval: Ne pas encore traiter, il suffit de collecter des données
  3. Processus dans terminate: c’est là que le travail réel se produit
  4. Utiliser yield pour retourner des lignes : renvoyer les résultats un par un à la fois à partir de terminate
  5. Gérer les cas de périphérie : Que se passe-t-il s’il y a moins de lignes que de clusters ?
  6. Conserver les types explicites : les retours UDTF ne peuvent pas référencer les types d’entrée