Zelfstudie: K-means-clustering

In deze tutorial bouwt u een door de gebruiker gedefinieerde Python-tabelfunctie (UDTF) voor Lakeflow Designer die K-means-clustering uitvoert met scikit-learn. UDDF's zijn geschikt voor machine learning-taken die volledige gegevenssets verwerken. Zie Door de gebruiker gedefinieerde operators in Lakeflow Designer voor achtergrondinformatie over door de gebruiker gedefinieerde operators.

Overview

In deze zelfstudie leert u hoe u een door de gebruiker gedefinieerde UDTF-operator maakt met behulp van Python. De operator voert K-Means-clustering uit op geselecteerde kolommen, zodat gebruikers het volgende kunnen doen:

  • Kies welke kolommen u wilt gebruiken als functies.
  • Geef het aantal clusters op.
  • Haal een tabel terug met clustertoewijzingen voor elke rij.

Stap 1: Inzicht in het UDTF-handlerpatroon

Een UDTF wordt geïmplementeerd als een Python klasse met drie belangrijke methoden:

  • __init__(): Eenmaal aangeroepen voordat rijen worden verwerkt om de status te initialiseren.
  • eval(row, ...): Wordt voor elke invoerregel aangeroepen om gegevens te accumuleren.
  • terminate(): Aangeroepen nadat alle rijen zijn verwerkt om resultaten te opleveren.

Met dit patroon kan de UDTF het volgende doen:

  1. Alle gegevenspunten verzamelen tijdens eval() aanroepen
  2. Het K-Means-model trainen in terminate()
  3. Geclusterde resultatenrij per rij opleveren
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

De row parameter in eval() is een PySpark Row-object. Gebruik .asDict() deze om deze te converteren naar een woordenlijst voor eenvoudigere toegang.

Stap 2: de YAML voor de operator maken

De YAML-configuratie definieert hoe de operator wordt weergegeven in Lakeflow Designer. Voor deze operator:

  • Nummerparameter (k): Aantal clusters dat moet worden gemaakt
  • Widget selecteren (id_column): Vervolgkeuzelijst gevuld met kolommen uit invoertabel
  • Widget voor meervoudige selectie (columns): Selectie van meerdere functiekolommen
  • optionsSource: vervolgkeuzelijsten uit het invoertabelschema automatisch vullen
  • Invoerpoort: hiermee geeft u aan dat deze operator tabellaire gegevens accepteert
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

Zie door de gebruiker gedefinieerde operator YAML-referentie voor een uitgebreide handleiding voor alle beschikbare eigenschappen, gegevenstypen, widgets en opties.

Stap 3: de unity-catalogusfunctie maken

Combineer de YAML-configuratie en Python handlerklasse in één CREATE FUNCTION instructie.

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

Stap 4: Testen met voorbeeldgegevens

Voorbeeld van klantgegevens maken voor testen:

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

Test de 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')
)

In dit geval wilt u de clusterresultaten weer samenvoegen met de oorspronkelijke gegevens om clustertoewijzingen weer te geven:

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

Stap 5: De operator registreren

Als u de operator in Lakeflow Designer wilt gebruiken, moet u deze registreren door deze toe te voegen aan uw .user_defined_operators.yaml bestand:

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

Note

Als u dit bestand in uw gebruikersmap definieert, wordt het alleen voor u weergegeven. Zie Uw operator detecteerbaar maken voor meer informatie.

Stap 6: Machtigingen instellen

Toegang verlenen aan gebruikers die deze operator moeten gebruiken:

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

Operator gebruiken in Lakeflow Designer

Nadat deze is geregistreerd, wordt de operator weergegeven in Lakeflow Designer met:

  • Een invoerpoort om verbinding te maken met uw gegevensbron
  • Een vervolgkeuzelijst om te selecteren welke kolom rijen uniek identificeert
  • Een meervoudige selectie om te kiezen welke kolommen moeten worden gebruikt als clusteringfuncties
  • Een numeriek invoerveld voor het gewenste aantal clusters

Gebruikers kunnen klanten, producten of andere gegevens segmenteren in zinvolle groepen zonder code te schrijven.

Tips voor het bouwen van UDDF's

  1. Status initialiseren in __init__: Lege lijsten/variabelen instellen om gegevens te verzamelen
  2. Verzamelen in eval: Nog niet verwerken, gewoon gegevens verzamelen
  3. Proces in terminate: Dit is waar het echte werk gebeurt
  4. Gebruiken yield om rijen te retourneren: resultaten één voor één retourneren van terminate
  5. Randgevallen afhandelen: Wat gebeurt er als er minder rijen zijn dan clusters?
  6. Typen expliciet houden: UDTF retourneert geen verwijzing naar invoertypen