Tutorial: k-Means-Algorithmus

In diesem Tutorial erstellen Sie einen Python-benutzerdefinierten Tabellenfunktionsoperator (UDTF) für Lakeflow Designer, der K-Means-Clustering mit scikit-learn ausführt. UDTFs eignen sich gut für Maschinelle Lernaufgaben, die ganze Datasets verarbeiten. Hintergrundinformationen zu benutzerdefinierten Operatoren finden Sie unter Benutzerdefinierte Operatoren in Lakeflow Designer.

Überblick

Dieses Lernprogramm führt Sie durch das Erstellen eines benutzerdefinierten UDTF-Operators mithilfe von Python. Der Operator wendet K-Means-Clustering auf ausgewählte Spalten an und ermöglicht Benutzern:

  • Wählen Sie aus, welche Spalten als Features verwendet werden sollen.
  • Geben Sie die Anzahl der Cluster an.
  • Rufen Sie eine Tabelle mit Clusterzuweisungen für jede Zeile zurück.

Schritt 1: Grundlegendes zum UDTF-Handlermuster

Ein UDTF wird als Python Klasse mit drei Schlüsselmethoden implementiert:

  • __init__(): Wird einmal aufgerufen, bevor Zeilen verarbeitet werden, um den Zustand zu initialisieren.
  • eval(row, ...): Wird für jede Eingabezeile aufgerufen, um Daten zu sammeln.
  • terminate(): Wird aufgerufen, nachdem alle Zeilen verarbeitet wurden, um Ergebnisse zu erzielen.

Dieses Muster ermöglicht dem UDTF Folgendes:

  1. Alle Datenpunkte während eval()-Anrufen sammeln
  2. Trainieren Sie das K-Means-Modell in terminate()
  3. Gruppierte Ergebnisse zeilenweise ausgeben
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

Der row Parameter in eval() ist ein PySpark Row -Objekt. Verwenden Sie .asDict(), um es für einen einfacheren Zugriff in ein Wörterbuch umzuwandeln.

Schritt 2: Erstellen des YAML für den Operator

Die YAML-Konfiguration definiert, wie der Operator im Lakeflow-Designer angezeigt wird. Für diesen Operator:

  • Number-Parameter (k): Anzahl der zu erstellenden Cluster
  • Auswahl-Widget (id_column): Dropdown mit Spalten aus der Eingabetabelle ausgefüllt
  • Mehrfachauswahl-Widget (columns): Auswahl mehrerer Featurespalten
  • optionsSource: Füllt Dropdowns automatisch aus dem Eingabetabellenschema auf.
  • Eingabeport: Gibt an, dass dieser Operator tabellarische Daten akzeptiert.
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

Eine umfassende Anleitung zu allen verfügbaren Eigenschaften, Datentypen, Widgets und Optionen finden Sie in der YAML-Referenz für benutzerdefinierte Operatoren .

Schritt 3: Erstellen der Unity-Katalogfunktion

Kombinieren Sie die YAML-Konfiguration und Python Handlerklasse in einer einzelnen CREATE FUNCTION-Anweisung.

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

Schritt 4: Testen mit Beispieldaten

Erstellen Sie Beispiel-Kundendaten für 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);

Testen Sie das 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 diesem Fall möchten Sie die Clusterergebnisse mit ursprünglichen Daten verknüpfen, um Clusterzuweisungen anzuzeigen:

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

Schritt 5: Registrieren des Operators

Um den Operator in Lakeflow Designer zu verwenden, müssen Sie ihn registrieren, indem Sie ihn zu Ihrer .user_defined_operators.yaml Datei hinzufügen:

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

Note

Wenn Sie diese Datei in Ihrem Benutzerordner definieren, wird sie nur für Sie angezeigt. Weitere Informationen finden Sie unter Machen Sie Ihren Operator auffindbar.

Schritt 6: Einrichten von Berechtigungen

Gewähren Sie Zugriff auf Benutzer, die diesen Operator verwenden müssen:

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

Verwenden Sie den Operator in Lakeflow Designer

Nachdem der Operator registriert ist, erscheint er im Lakeflow Designer mit:

  • Ein Eingabeport zum Verbinden der Datenquelle
  • Eine Dropdownliste, um auszuwählen, welche Spalte Zeilen eindeutig identifiziert
  • Eine Mehrfachauswahl, um festzulegen, welche Spalten als Merkmale für das Clustering verwendet werden sollen
  • Eine Zahleneingabe für die gewünschte Anzahl von Clustern

Benutzer können Kunden, Produkte oder andere Daten in aussagekräftige Gruppen segmentieren, ohne Code zu schreiben.

Tipps zum Erstellen von UDTFs

  1. Initialisieren des Zustands in __init__: Einrichten leerer Listen/Variablen zum Ansammeln von Daten
  2. Ansammeln in eval: Verarbeiten Sie noch nicht, sammeln Sie einfach Daten.
  3. Prozess in terminate: Hier geschieht die eigentliche Arbeit
  4. yield zum Zurückgeben von Zeilen verwenden: Ergebnisse einzeln aus terminate zurückgeben
  5. Randfälle behandeln: Was passiert, wenn weniger Zeilen als Cluster vorhanden sind?
  6. Typen explizit beibehalten: UDTF-Rückgaben können nicht auf Eingabetypen verweisen.