Notitie
Voor toegang tot deze pagina is autorisatie vereist. U kunt proberen u aan te melden of de directory te wijzigen.
Voor toegang tot deze pagina is autorisatie vereist. U kunt proberen de mappen te wijzigen.
In deze tutorial bouwt u een door de gebruiker gedefinieerde Python-tabelfunctie (UDTF) voor Lakeflow Designer die K-means-clustering uitvoert met scikit-learn. UDDF's zijn geschikt voor machine learning-taken die volledige gegevenssets verwerken. Zie Door de gebruiker gedefinieerde operators in Lakeflow Designer voor achtergrondinformatie over door de gebruiker gedefinieerde operators.
Overview
In deze zelfstudie leert u hoe u een door de gebruiker gedefinieerde UDTF-operator maakt met behulp van Python. De operator voert K-Means-clustering uit op geselecteerde kolommen, zodat gebruikers het volgende kunnen doen:
- Kies welke kolommen u wilt gebruiken als functies.
- Geef het aantal clusters op.
- Haal een tabel terug met clustertoewijzingen voor elke rij.
Stap 1: Inzicht in het UDTF-handlerpatroon
Een UDTF wordt geïmplementeerd als een Python klasse met drie belangrijke methoden:
-
__init__(): Eenmaal aangeroepen voordat rijen worden verwerkt om de status te initialiseren. -
eval(row, ...): Wordt voor elke invoerregel aangeroepen om gegevens te accumuleren. -
terminate(): Aangeroepen nadat alle rijen zijn verwerkt om resultaten te opleveren.
Met dit patroon kan de UDTF het volgende doen:
- Alle gegevenspunten verzamelen tijdens
eval()aanroepen - Het K-Means-model trainen in
terminate() - Geclusterde resultatenrij per rij opleveren
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
De row parameter in eval() is een PySpark Row-object. Gebruik .asDict() deze om deze te converteren naar een woordenlijst voor eenvoudigere toegang.
Stap 2: de YAML voor de operator maken
De YAML-configuratie definieert hoe de operator wordt weergegeven in Lakeflow Designer. Voor deze operator:
-
Nummerparameter (
k): Aantal clusters dat moet worden gemaakt -
Widget selecteren (
id_column): Vervolgkeuzelijst gevuld met kolommen uit invoertabel -
Widget voor meervoudige selectie (
columns): Selectie van meerdere functiekolommen -
optionsSource: vervolgkeuzelijsten uit het invoertabelschema automatisch vullen - Invoerpoort: hiermee geeft u aan dat deze operator tabellaire gegevens accepteert
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
Zie door de gebruiker gedefinieerde operator YAML-referentie voor een uitgebreide handleiding voor alle beschikbare eigenschappen, gegevenstypen, widgets en opties.
Stap 3: de unity-catalogusfunctie maken
Combineer de YAML-configuratie en Python handlerklasse in één CREATE FUNCTION instructie.
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)
$$
Stap 4: Testen met voorbeeldgegevens
Voorbeeld van klantgegevens maken voor testen:
-- 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);
Test de 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 dit geval wilt u de clusterresultaten weer samenvoegen met de oorspronkelijke gegevens om clustertoewijzingen weer te geven:
-- 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
Stap 5: De operator registreren
Als u de operator in Lakeflow Designer wilt gebruiken, moet u deze registreren door deze toe te voegen aan uw .user_defined_operators.yaml bestand:
operators:
- catalog: main
schema: my_schema
functionName: k_means
Note
Als u dit bestand in uw gebruikersmap definieert, wordt het alleen voor u weergegeven. Zie Uw operator detecteerbaar maken voor meer informatie.
Stap 6: Machtigingen instellen
Toegang verlenen aan gebruikers die deze operator moeten gebruiken:
GRANT USE SCHEMA ON SCHEMA main.my_schema TO `<user>`;
GRANT EXECUTE ON FUNCTION main.my_schema.k_means TO `<user>`;
Operator gebruiken in Lakeflow Designer
Nadat deze is geregistreerd, wordt de operator weergegeven in Lakeflow Designer met:
- Een invoerpoort om verbinding te maken met uw gegevensbron
- Een vervolgkeuzelijst om te selecteren welke kolom rijen uniek identificeert
- Een meervoudige selectie om te kiezen welke kolommen moeten worden gebruikt als clusteringfuncties
- Een numeriek invoerveld voor het gewenste aantal clusters
Gebruikers kunnen klanten, producten of andere gegevens segmenteren in zinvolle groepen zonder code te schrijven.
Tips voor het bouwen van UDDF's
-
Status initialiseren in
__init__: Lege lijsten/variabelen instellen om gegevens te verzamelen -
Verzamelen in
eval: Nog niet verwerken, gewoon gegevens verzamelen -
Proces in
terminate: Dit is waar het echte werk gebeurt -
Gebruiken
yieldom rijen te retourneren: resultaten één voor één retourneren vanterminate - Randgevallen afhandelen: Wat gebeurt er als er minder rijen zijn dan clusters?
- Typen expliciet houden: UDTF retourneert geen verwijzing naar invoertypen