Tutorial: K-significa agrupamento

Neste tutorial, vai criar um operador de função de tabela definida pelo utilizador (UDTF) em Python para o Lakeflow Designer que executa agrupamento K-means com o scikit-learn. Os UDTFs são bem adequados para tarefas de aprendizagem automática que processam conjuntos de dados inteiros. Para contexto sobre operadores definidos pelo utilizador, veja Operadores definidos pelo utilizador no Lakeflow Designer.

Overview

Este tutorial orienta-o na criação de um operador UDTF definido pelo utilizador usando Python. O operador realiza agrupamento K-Means em colunas selecionadas, permitindo aos utilizadores:

  • Escolhe que colunas usar como funcionalidades.
  • Especifique o número de clusters.
  • Recebe uma tabela com atribuições de clusters para cada linha.

Passo 1: Compreender o padrão de handler UDTF

Uma UDTF é implementada como uma classe Python com três métodos-chave:

  • __init__(): Chamado uma vez antes de quaisquer linhas serem processadas para inicializar o estado.
  • eval(row, ...): É chamado para cada linha de entrada para acumular dados.
  • terminate(): Chamado após todas as linhas serem processadas para obter resultados.

Este padrão permite à UDTF:

  1. Recolha todos os pontos de dados durante as eval() chamadas
  2. Treine o modelo K-Means em terminate()
  3. Obter resultados agrupados linha a 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 row parâmetro em eval() é um objeto PySpark Row. Usa .asDict() para converter para um dicionário para facilitar o acesso.

Passo 2: Criar o YAML para o operador

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

  • Parâmetro numérico (k): Número de clusters a criar
  • Selecionar o widget (id_column): Lista pendente preenchida com as colunas da tabela de entrada
  • Widget de seleção múltipla (columns): Seleção de várias colunas de atributos
  • optionsSource: Preenche automaticamente as listas pendentes com base no esquema da tabela de entrada
  • Porta de entrada: Especifica que este 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 utilizador para um guia abrangente de todas as propriedades, tipos de dados, widgets e opções disponíveis.

Passo 3: Criar a função Unity Catalog

Combine a configuração YAML e a classe handler Python numa ú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)
$$

Passo 4: Testar com dados de amostra

Crie dados de amostra de clientes para testes:

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

Neste caso, pretende voltar a associar os resultados do agrupamento aos dados originais para ver as atribuições aos clusters:

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

Passo 5: Registar o operador

Para usar o operador no Lakeflow Designer, deve registá-lo, adicionando-o ao seu .user_defined_operators.yaml ficheiro:

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

Note

Se definir este ficheiro na sua pasta de utilizador, ele só aparece para si. Para mais informações, consulte Tornar o seu operador descobrável.

Passo 6: Configurar permissões

Conceda acesso a utilizadores que necessitem de usar este operador:

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

Utilize o operador no Lakeflow Designer

Depois de registado, o operador aparece no Lakeflow Designer com:

  • Uma porta de entrada para ligar a sua fonte de dados
  • Um menu suspenso para selecionar qual a coluna que identifica de forma única as linhas
  • Uma multi-seleção para escolher quais as colunas a usar como funcionalidades de agrupamento
  • Uma entrada numérica para o número desejado de clusters

Os utilizadores podem segmentar clientes, produtos ou quaisquer outros dados em grupos significativos sem necessidade de escrever código.

Dicas para criar UDTFs

  1. Inicializar estado em __init__: Configurar listas/variáveis vazias para acumular dados
  2. Acumular em eval: Ainda não processar, apenas recolher dados
  3. Processo em terminate: É aqui que acontece o verdadeiro trabalho
  4. Utilize yield para devolver linhas: devolva os resultados um de cada vez a partir de terminate
  5. Lidar com casos limite: E se houver menos linhas do que clusters?
  6. Mantenha os tipos explícitos: os retornos do UDTF não podem referenciar os tipos de entrada