チュートリアル: K-Means クラスタリング

このチュートリアルでは、scikit-learn を使用して K-means クラスタリングを実行する、Lakeflow Designer 用の Python ユーザー定義テーブル関数 (UDTF) 演算子を構築します。 UDF は、データセット全体を処理する機械学習タスクに適しています。 ユーザー定義演算子の背景については、 Lakeflow Designer の「ユーザー定義演算子」を参照してください。

開始する

このチュートリアルでは、Pythonを使用して UDTF ユーザー定義演算子を作成する手順について説明します。 オペレーターは、選択した列に対して K-Means クラスタリングを実行し、ユーザーは次のことを行うことができます。

  • 機能として使用する列を選択します。
  • クラスターの数を指定します。
  • 各行のクラスター割り当てを含むテーブルを取得します。

手順 1: UDTF ハンドラー パターンを理解する

UDTF は、次の 3 つの主要なメソッドを持つPython クラスとして実装されます。

  • __init__(): 状態を初期化するために行が処理される前に 1 回呼び出されます。
  • eval(row, ...): データを蓄積するために入力行ごとに呼び出されます。
  • terminate(): 結果を生成するためにすべての行が処理された後に呼び出されます。

このパターンにより、UDTF は次のことができます。

  1. eval()呼び出し中にすべてのデータ ポイントを収集する
  2. で K-Means モデルをトレーニングする terminate()
  3. クラスター化された結果を行ごとに生成する
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

roweval() パラメーターは PySpark Row オブジェクトです。 .asDict()を使用してディクショナリに変換すると、簡単にアクセスできます。

手順 2: 演算子の YAML を作成する

YAML 構成は、Lakeflow Designer でのオペレーターの表示方法を定義します。 この演算子の場合:

  • Number パラメーター (k): 作成するクラスターの数
  • ウィジェットを選択 (id_column): 入力テーブルの列が表示されるドロップダウン
  • 複数選択ウィジェット (columns): 複数の特徴列の選択
  • optionsSource: 入力テーブルのスキーマに基づいてドロップダウンを自動的に入力します
  • 入力ポート: この演算子が表形式データを受け入れることを指定します
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

使用可能なすべてのプロパティ、データ型、ウィジェット、およびオプションに関する包括的なガイドについては、 ユーザー定義演算子 YAML リファレンスを参照 してください。

手順 3: Unity カタログ関数を作成する

YAML 構成とPythonハンドラー クラスを 1 つの 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)
$$

手順 4: サンプル データを使用してテストする

テスト用のサンプル顧客データを作成します。

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

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

この場合は、クラスターの結果を元のデータと結合して、クラスターの割り当てを確認します。

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

手順 5: オペレーターを登録する

Lakeflow Designer でオペレーターを使用するには、 .user_defined_operators.yaml ファイルに追加して、オペレーターを登録する必要があります。

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

Note

このファイルをユーザー フォルダーに定義すると、そのファイルのみが表示されます。 詳細については、「 オペレーターを検出可能にする」を参照してください

手順 6: アクセス許可を設定する

このオペレーターを使用する必要があるユーザーにアクセス権を付与します。

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

Lakeflow Designer で演算子を使用する

登録が完了すると、オペレーターは Lakeflow Designer に次の状態で表示されます。

  • データ ソースを接続するための入力ポート
  • 行を一意に識別する列を選択するドロップダウン
  • クラスタリングの特徴量として使用する列を選択する複数選択
  • クラスターの目的の数に対する数値入力

ユーザーは、コードを記述することなく、顧客、製品、またはその他のデータを意味のあるグループに分割できます。

UDF の構築に関するヒント

  1. __init__で状態を初期化する: 空のリスト/変数を設定してデータを蓄積する
  2. evalに蓄積する: まだ処理せず、データを収集するだけです
  3. terminateのプロセス: これが実際の作業が行われる場所です
  4. 行を返すには yield を使用します: terminate から結果を 1 行ずつ返します
  5. エッジ ケースの処理: クラスターよりも行数が少ない場合はどうなりますか?
  6. 型を明示的に保持する: UDTF の戻り値は入力型を参照できません