Nota
O acesso a esta página requer autorização. Pode tentar iniciar sessão ou alterar os diretórios.
O acesso a esta página requer autorização. Pode tentar alterar os diretórios.
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:
- Recolha todos os pontos de dados durante as
eval()chamadas - Treine o modelo K-Means em
terminate() - 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
-
Inicializar estado em
__init__: Configurar listas/variáveis vazias para acumular dados -
Acumular em
eval: Ainda não processar, apenas recolher dados -
Processo em
terminate: É aqui que acontece o verdadeiro trabalho -
Utilize
yieldpara devolver linhas: devolva os resultados um de cada vez a partir determinate - Lidar com casos limite: E se houver menos linhas do que clusters?
- Mantenha os tipos explícitos: os retornos do UDTF não podem referenciar os tipos de entrada