Tutorial: Agrupación en clústeres k-means

En este tutorial, crearás un operador Python de función tabular definida por el usuario (UDTF) para Lakeflow Designer que ejecuta la agrupación K-means con scikit-learn. Las UDF son adecuadas para las tareas de aprendizaje automático que procesan conjuntos de datos completos. Para obtener información sobre los operadores definidos por el usuario, consulte Operadores definidos por el usuario en Lakeflow Designer.

Información general

Este tutorial le guiará a través de la creación de un operador definido por el usuario UDTF mediante Python. El operador realiza la agrupación en clústeres K-Means en columnas seleccionadas, lo que permite a los usuarios:

  • Elija las columnas que se van a usar como características.
  • Especifique el número de clústeres.
  • Vuelva a obtener una tabla con asignaciones de clúster para cada fila.

Paso 1: Descripción del patrón de controlador UDTF

Un UDTF se implementa como una clase Python con tres métodos clave:

  • __init__(): Se llama una sola vez antes de procesar ninguna fila para inicializar el estado.
  • eval(row, ...): se llama para cada fila de entrada para acumular datos.
  • terminate(): se llama después de que se procesen todas las filas para generar resultados.

Este patrón permite que el UDTF:

  1. Recopilación de todos los puntos de datos durante las llamadas a eval()
  2. Entrene el modelo K-Means en terminate()
  3. Devolver resultados agrupados fila por fila
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

El row parámetro de eval() es un objeto Row de PySpark. Use .asDict() para convertirlo en un diccionario para facilitar el acceso.

Paso 2: Creación del código YAML para el operador

La configuración de YAML define cómo aparece el operador en Lakeflow Designer. Para este operador:

  • Parámetro number (k): número de clústeres que se van a crear
  • Seleccionar widget (id_column): lista desplegable rellenada con columnas de la tabla de entrada
  • Widget de selección múltiple (columns): selección de varias columnas de características
  • optionsSource: rellena automáticamente las listas desplegables del esquema de tabla de entrada.
  • Puerto de entrada: especifica que este operador acepta datos tabulares.
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

Consulte la referencia de YAML del operador definido por el usuario para obtener una guía exhaustiva de todas las propiedades, tipos de datos, widgets y opciones disponibles.

Paso 3: Crear la función catálogo de Unity

Combine la configuración YAML y la clase controladora de Python en una sola instrucción 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)
$$

Paso 4: Prueba con datos de ejemplo

Cree datos de cliente de ejemplo para realizar pruebas:

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

Pruebe el 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')
)

En este caso, querrá unir de nuevo los resultados de la agrupación en clústeres con los datos originales para ver las asignaciones a clústeres:

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

Paso 5: Registrar el operador

Para usar el operador en Lakeflow Designer, debe registrarlo, agregándolo al .user_defined_operators.yaml archivo:

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

Note

Si define este archivo en la carpeta de usuario, solo aparece para usted. Para obtener más información, consulte Hacer que el operador sea reconocible.

Paso 6: Configuración de permisos

Conceda acceso a los usuarios que necesitan usar este operador:

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

Use el operador en Lakeflow Designer

Una vez registrado, el operador aparece en Lakeflow Designer con:

  • Un puerto de entrada para conectar el origen de datos
  • Lista desplegable para seleccionar qué columna identifica de forma única las filas.
  • Selección múltiple para elegir las columnas que se van a usar como características de agrupación en clústeres
  • Una entrada de número para el número deseado de clústeres

Los usuarios pueden segmentar clientes, productos o cualquier otro dato en grupos significativos sin escribir código.

Sugerencias para crear UDTF

  1. Inicializar el estado en __init__: configurar listas o variables vacías para acumular datos
  2. Acumular en eval: No procesar todavía, simplemente recopilar datos
  3. Proceso en terminate: aquí es donde se produce el trabajo real
  4. Uso yield para devolver filas: devuelve resultados de uno a uno desde terminate
  5. Gestionar casos perimetrales: ¿Qué ocurre si hay menos filas que clústeres?
  6. Mantener los tipos explícitos: las devoluciones UDTF no pueden hacer referencia a los tipos de entrada