Självstudiekurs: K-means-klustring

I den här självstudien skapar du en användardefinierad tabellfunktionsoperator i Python (UDTF) för Lakeflow Designer som kör K-means-klustring med scikit-learn. UDF:er passar bra för maskininlärningsuppgifter som bearbetar hela datamängder. Bakgrund om användardefinierade operatorer finns i Användardefinierade operatorer i Lakeflow Designer.

Overview

Den här självstudien beskriver hur du skapar en användardefinierad UDTF-operator med hjälp av Python. Operatorn utför K-Means-klustring på valda kolumner så att användarna kan:

  • Välj vilka kolumner som ska användas som funktioner.
  • Ange antalet kluster.
  • Hämta en tabell med klustertilldelningar för varje rad.

Steg 1: Förstå UDTF-hanterarens mönster

En UDTF implementeras som en Python-klass med tre viktiga metoder:

  • __init__(): Anropas en gång innan några rader bearbetas för att initiera tillståndet.
  • eval(row, ...): Anropas för varje indatarad för att ackumulera data.
  • terminate(): Anropas när alla rader har bearbetats för att ge resultat.

Med det här mönstret kan UDTF:

  1. Samla in alla datapunkter under eval() anrop
  2. Träna K-Means-modellen i terminate()
  3. Returnera grupperade resultat rad för rad
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

Parametern row i eval() är ett PySpark Row-objekt. Använd .asDict() för att konvertera den till en ordlista för enklare åtkomst.

Steg 2: Skapa YAML för operatorn

YAML-konfigurationen definierar hur operatorn visas i Lakeflow Designer. För den här operatorn:

  • Talparameter (k): Antal kluster som ska skapas
  • Välj widget (id_column): Listrutan fylls med kolumner från indatatabellen
  • Multi-select widget (columns): Val av flera funktionskolumner
  • optionsSource: Fyller listrutor automatiskt utifrån indatatabellens schema
  • Ingångsport: Anger att den här operatorn accepterar data i tabellform
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

Mer information om alla tillgängliga egenskaper, datatyper, widgetar och alternativ finns i YAML-referens för användardefinierad operatör .

Steg 3: Skapa funktionen Unity Catalog

Kombinera YAML-konfigurationen och Python-hanterarklassen till en enda CREATE FUNCTION-instruktion.

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

Steg 4: Testa med exempeldata

Skapa exempelkunddata för testning:

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

Testa 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')
)

I det här fallet vill du koppla tillbaka klustringsresultaten med ursprungliga data för att se klustertilldelningar:

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

Steg 5: Registrera operatorn

Om du vill använda operatorn i Lakeflow Designer måste du registrera den genom att lägga till den i din .user_defined_operators.yaml-fil:

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

Note

Om du definierar den här filen i användarmappen visas den bara för dig. Mer information finns i Gör operatören identifierbar.

Steg 6: Konfigurera behörigheter

Bevilja åtkomst till användare som behöver använda den här operatorn:

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

Använda operatorn i Lakeflow Designer

När den har registrerats visas operatorn i Lakeflow Designer med:

  • En indataport för anslutning av din datakälla
  • En listruta för att välja vilken kolumn som unikt identifierar rader
  • Ett flerval för att välja vilka kolumner som ska användas som klustringsfunktioner
  • Ett inmatningsfält för tal för det önskade antalet kluster

Användare kan segmentera kunder, produkter eller andra data i meningsfulla grupper utan att skriva kod.

Tips för att bygga UDTF:er

  1. Initiera tillståndet i __init__: Konfigurera tomma listor/variabler för att ackumulera data
  2. Ackumulera i eval: Bearbeta inte ännu, samla bara in data
  3. Process i terminate: Det är här det verkliga arbetet sker
  4. Använd yield för att returnera rader: Returnera resultat en i taget från terminate
  5. Hantera kantfall: Vad händer om det finns färre rader än kluster?
  6. Håll typerna explicita: UDTF-returer kan inte referera till indatatyper