Esercitazione: Clustering K-means

In questa esercitazione viene creato un operatore UDTF (User-Defined Table Function) Python per Lakeflow Designer che esegue il clustering K-means con scikit-learn. Le UDTF sono particolarmente adatte alle attività di apprendimento automatico che elaborano interi insiemi di dati. Per informazioni di base sugli operatori definiti dall'utente, vedere Operatori definiti dall'utente in Lakeflow Designer.

Informazioni generali

Questa esercitazione illustra come creare un operatore definito dall'utente UDTF usando Python. L'operatore esegue il clustering K-Means sulle colonne selezionate, consentendo agli utenti di:

  • Scegliere le colonne da usare come funzionalità.
  • Specificare il numero di cluster.
  • Restituisce una tabella contenente le assegnazioni ai cluster per ogni riga.

Passaggio 1: Comprendere il modello di gestore UDTF

Una UDTF viene implementata come una classe Python con tre metodi chiave:

  • __init__(): Chiamato una volta prima che tutte le righe vengano elaborate per inizializzare lo stato.
  • eval(row, ...): Chiamato per ogni riga di input per l'accumulo dei dati.
  • terminate(): chiamato dopo l'elaborazione di tutte le righe per produrre risultati.

Questo schema consente alla UDTF di:

  1. Raccogliere tutti i punti dati durante eval() le chiamate
  2. Addestra il modello K-Means in terminate()
  3. Restituisce i risultati raggruppati riga per riga
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)

Annotazioni

Il row parametro in eval() è un oggetto PySpark Row. Usare .asDict() per convertirlo in un dizionario per facilitare l'accesso.

Passaggio 2: Creare il file YAML per l'operatore

La configurazione YAML definisce la modalità di visualizzazione dell'operatore in Lakeflow Designer. Per questo operatore:

  • Parametro number (k): numero di cluster da creare
  • Seleziona widget (id_column): menu a discesa popolato con le colonne della tabella di input
  • Controllo a selezione multipla (columns): selezione di più colonne delle funzionalità
  • optionsSource: popola automaticamente gli elenchi a discesa dallo schema della tabella di input
  • Porta di input: specifica che questo operatore accetta dati tabulari
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

Per una guida completa a tutte le proprietà, ai tipi di dati, ai widget e alle opzioni disponibili, vedere il riferimento YAML per gli operatori definiti dall'utente.

Passaggio 3: Creare la funzione Catalogo Unity

Combina la configurazione YAML e la classe handler Python in un'unica istruzione 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)
$$

Passaggio 4: Testare con dati di esempio

Creare dati dei clienti di esempio per i test:

-- 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);

Testare il formato definito dall'utente 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')
)

In questo caso, si vuole unire i risultati del clustering con i dati originali per visualizzare le assegnazioni del 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

Passaggio 5: Registrare l'operatore

Per usare l'operatore in Lakeflow Designer, è necessario registrarlo aggiungendolo al .user_defined_operators.yaml file:

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

Annotazioni

Se si definisce questo file nella cartella utente, viene visualizzato solo per l'utente. Per altre informazioni, vedere Rendere individuabile l'operatore.

Passaggio 6: Configurare le autorizzazioni

Concedere l'accesso agli utenti che devono usare questo operatore:

GRANT USE SCHEMA ON SCHEMA main.my_schema TO `<user>`;
GRANT EXECUTE ON FUNCTION main.my_schema.k_means TO `<user>`;

Usare l'operatore in Lakeflow Designer

Dopo la registrazione, l'operatore viene visualizzato in Lakeflow Designer con:

  • Una porta di ingresso per collegare l'origine dei dati
  • Elenco a discesa per selezionare la colonna che identifica in modo univoco le righe
  • Selezione multipla per scegliere le colonne da usare come funzionalità di clustering
  • Input numerico per il numero desiderato di cluster

Gli utenti possono segmentare clienti, prodotti o altri dati in gruppi significativi senza scrivere codice.

Suggerimenti per la creazione di UDTF

  1. Inizializzare lo stato in __init__: Configurare elenchi/variabili vuoti per accumulare dati
  2. Accumula in eval: non elaborare ancora, raccogli solo i dati
  3. Processo in terminate: questo è il percorso in cui si verifica il lavoro reale
  4. Usare yield per restituire righe: restituisce risultati uno alla volta da terminate
  5. Gestire i casi perimetrali: cosa accade se sono presenti meno righe rispetto ai cluster?
  6. Mantieni espliciti i tipi: i valori restituiti da UDTF non possono fare riferimento ai tipi di input