Tutorial: cluster K-means

Neste tutorial, você criará uma função de tabela definida pelo usuário (UDTF) em Python no Lakeflow Designer que executa o agrupamento K-means com o scikit-learn. UDTFs são adequados para tarefas de aprendizado de máquina que processam conjuntos de dados inteiros. Para obter informações sobre operadores definidos pelo usuário, consulte operadores definidos pelo usuário no Lakeflow Designer.

Visão geral

Este tutorial explica como criar um operador definido pelo usuário da UDTF usando Python. O operador realiza o agrupamento K-Means nas colunas selecionadas, permitindo aos usuários:

  • Escolha quais colunas usar como recursos.
  • Especifique o número de clusters.
  • Retorne uma tabela com atribuições de cluster para cada linha.

Etapa 1: Entender o padrão de manipulador UDTF

Um UDTF é implementado como uma classe Python com três métodos principais:

  • __init__(): Chamado uma vez antes que qualquer linha seja processada para inicializar o estado.
  • eval(row, ...): chamado para cada linha de entrada para acumular dados.
  • terminate(): chamado depois que todas as linhas são processadas para produzir resultados.

Esse padrão permite que o UDTF:

  1. Coletar todos os pontos de dados durante eval() chamadas
  2. Treine o modelo K-Means em terminate()
  3. Produzir resultados clusterizados linha por linha
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

O parâmetro row em eval() é um objeto Row do PySpark. Use .asDict() para convertê-lo em um dicionário para facilitar o acesso.

Etapa 2: Criar o YAML para o operador

A configuração yaml define como o operador aparece no Lakeflow Designer. Para este operador:

  • Parâmetro de número (k): número de clusters a serem criados
  • Selecionar widget (id_column): lista suspensa preenchida com colunas da tabela de entrada
  • Widget de seleção múltipla (columns): seleção de várias colunas de recursos
  • optionsSource: preenche automaticamente as listas suspensas do esquema da tabela de entrada
  • Porta de entrada: especifica que esse operador aceita dados 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 a referência YAML do operador definido pelo usuário para obter um guia abrangente sobre todas as propriedades, tipos de dados, widgets e opções disponíveis.

Etapa 3: Criar a função catálogo do Unity

Combine a configuração YAML e a classe de manipulador Python em uma única instrução 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)
$$

Etapa 4: Testar com dados de exemplo

Crie dados de cliente de exemplo para teste:

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

Teste o 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')
)

Nesse caso, você deseja unir os resultados do clustering de volta com os dados originais para ver as atribuições de 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

Etapa 5: Registrar o operador

Para usar o operador no Lakeflow Designer, você deve registrá-lo adicionando-o ao seu .user_defined_operators.yaml arquivo:

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

Note

Se você definir esse arquivo em sua pasta de usuário, ele será exibido apenas para você. Para obter mais informações, consulte Tornar seu operador detectável.

Etapa 6: Configurar permissões

Conceda acesso aos usuários que precisam usar este operador:

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

Usar o operador no Designer do Lakeflow

Depois que ele é registrado, o operador aparece no Lakeflow Designer com:

  • Uma porta de entrada para conectar sua fonte de dados
  • Uma lista suspensa para selecionar qual coluna identifica cada linha de forma exclusiva
  • Uma seleção múltipla para escolher quais colunas usar como recursos de clustering
  • Uma entrada numérica para o número desejado de clusters

Os usuários podem segmentar clientes, produtos ou quaisquer outros dados em grupos significativos sem escrever código.

Dicas para criar UDTFs

  1. Inicializar o estado em __init__: Configurar listas/variáveis vazias para acumular dados
  2. Acumular em eval: Não processe ainda, apenas colete os dados
  3. Processo em terminate: É aqui que o trabalho real acontece
  4. Usar yield para retornar linhas: retornar resultados um de cada vez de terminate
  5. Lidar com casos extremos: e se houver menos linhas do que clusters?
  6. Manter os tipos explícitos: os retornos UDTF não podem referenciar tipos de entrada