メインコンテンツまでスキップ

Unity CatalogのバッチPythonユーザー定義関数 (UDF)

バッチ Unity Catalog Python UDFs は一般提供 (GA) されています。これらは、1行ずつではなく、データのバッチ単位で処理を行います。

必要条件​

クラシック コンピュートでは、バッチ Unity Catalog Python UDF には Databricks Runtime 16.3 以降が必要です。これらは、Serverless コンピュート、および Pro と Serverless の SQLウェアハウスでもサポートされています。

追加の機能には、独自のコンピュートおよびバージョンの要件があります。Python UDFの機能要件を参照してください。

バッチ Unity Catalog Python UDF の作成​

バッチ Unity Catalog Python UDF の作成は、通常の Unity Catalog UDF の作成と似ていますが、次の点が追加されています。

  • PARAMETER STYLE PANDAS: これは、 UDF が Pandas イテレータを使用してバッチでデータを処理することを指定します。
  • HANDLER 'handler_function':バッチを処理するハンドラー関数を指定します。

次の例では、Unity Catalog内に永続的なバッチ Python UDF を作成します。my_catalog および my_schema をカタログとスキーマに置き換えます:

Python
%sql
CREATE OR REPLACE FUNCTION my_catalog.my_schema.calculate_bmi_pandas(weight_kg DOUBLE, height_m DOUBLE)
RETURNS DOUBLE
LANGUAGE PYTHON
DETERMINISTIC
PARAMETER STYLE PANDAS
HANDLER 'handler_function'
AS $$
import pandas as pd
from typing import Iterator, Tuple

def handler_function(batch_iter: Iterator[Tuple[pd.Series, pd.Series]]) -> Iterator[pd.Series]:
for weight_series, height_series in batch_iter:
yield weight_series / (height_series ** 2)
$$;

登録したら、UDF SQLまたは を使用してPython を呼び出すことができます。

SQL
SELECT person_id, my_catalog.my_schema.calculate_bmi_pandas(weight_kg, height_m) AS bmi
FROM (
SELECT 1 AS person_id, CAST(70.0 AS DOUBLE) AS weight_kg, CAST(1.75 AS DOUBLE) AS height_m UNION ALL
SELECT 2 AS person_id, CAST(80.0 AS DOUBLE) AS weight_kg, CAST(1.80 AS DOUBLE) AS height_m
);

バッチ UDF ハンドラ関数​

Batch Unity Catalog Python UDF には、バッチを処理して結果を生成するハンドラー関数が必要です。HANDLER 句を使用して UDF を作成する場合は、ハンドラー関数の名前を指定する必要があります。

ハンドラ関数は、次の処理を行います。

  1. 1 つ以上の pandas.Series を反復処理するイテレータ引数を受け入れます。各 pandas.Series には、UDF の入力パラメーターが含まれています。
  2. ジェネレータを反復処理し、データを処理します。
  3. ジェネレータイテレータを返します。

バッチ Unity Catalog Python UDF は、入力と同じ数の行を返す必要があります。ハンドラー関数は、各バッチの入力シリーズと同じ長さの pandas.Series を生成することで、これを保証します。

カスタム依存関係のインストール​

バッチ Unity Catalog Python UDF の機能を Databricks Runtime 環境を超えて拡張するには、外部ライブラリのカスタム依存関係を定義します。

カスタム依存関係を使用した UDF の拡張を参照してください。

Unity Catalog のシークレットにアクセスする​

バッチ Unity Catalog Python UDF は、SECRETS 句で宣言されたシークレットにアクセスできます。environment_version を 6 以上に明示的に設定する必要があります。この句を使用する UDF の直接呼び出しは、専用アクセス モードのコンピュートではサポートされていません。列マスクの例外については、専用コンピュート上の列マスクでのシークレット対応 UDF の使用を参照してください。

バッチ UDF は、1 つまたは複数のパラメーターを受け入れることができます​

単一のパラメーター: ハンドラー関数が単一の入力パラメーターを使用する場合、各バッチの pandas.Series に対するイテレータを受け取ります。

Python
%sql
CREATE OR REPLACE TEMPORARY FUNCTION one_parameter_udf(value INT)
RETURNS STRING
LANGUAGE PYTHON
DETERMINISTIC
PARAMETER STYLE PANDAS
HANDLER 'handler_func'
AS $$
import pandas as pd
from typing import Iterator
def handler_func(batch_iter: Iterator[pd.Series]) -> Iterator[pd.Series]:
for value_batch in batch_iter:
d = {"min": value_batch.min(), "max": value_batch.max()}
yield pd.Series([str(d)] * len(value_batch))
$$;
SELECT one_parameter_udf(id), count(*) from range(0, 100000, 3, 8) GROUP BY ALL;

複数のパラメーター: 複数の入力パラメーターの場合、ハンドラー関数は、複数の pandas.Seriesを反復処理するイテレータを受け取ります。 系列の値は、入力パラメーターと同じ順序です。

Python
%sql
CREATE OR REPLACE TEMPORARY FUNCTION two_parameter_udf(p1 INT, p2 INT)
RETURNS INT
LANGUAGE PYTHON
DETERMINISTIC
PARAMETER STYLE PANDAS
HANDLER 'handler_function'
AS $$
import pandas as pd
from typing import Iterator, Tuple

def handler_function(batch_iter: Iterator[Tuple[pd.Series, pd.Series]]) -> Iterator[pd.Series]:
for p1, p2 in batch_iter: # same order as arguments above
yield p1 + p2
$$;
SELECT two_parameter_udf(id , id + 1) from range(0, 100000, 3, 8);

コストのかかる操作を分離することでパフォーマンスを最適化​

計算コストの高い操作を最適化するには、これらの操作をハンドラー関数から分離します。これにより、データのバッチに対するすべての反復ではなく、一度だけ実行されるようになります。

次の例は、負荷の高い計算が 1 回だけ実行されるようにする方法を示しています。

Python
%sql
CREATE OR REPLACE TEMPORARY FUNCTION expensive_computation_udf(value INT)
RETURNS INT
LANGUAGE PYTHON
DETERMINISTIC
PARAMETER STYLE PANDAS
HANDLER 'handler_func'
AS $$
def compute_value():
# expensive computation...
return 1

expensive_value = compute_value()
def handler_func(batch_iter):
for batch in batch_iter:
yield batch * expensive_value
$$;
SELECT expensive_computation_udf(id), count(*) from range(0, 100000, 3, 8) GROUP BY ALL

環境分離​

注記

共有分離環境には、Databricks Runtime 17.1 以上が必要です。以前のバージョンでは、すべての バッチUnity Catalog Python UDF は厳密な分離モードで実行されました。

同じ所有者とセッションを持つバッチUnity Catalog Python UDFは、安全に分離環境を共有できます。 これにより、起動する必要がある個別の環境の数が減り、パフォーマンスが向上し、メモリ使用量が削減されます。

厳密な分離​

UDF が常に独自の完全に分離された環境で実行されるようにするには、 STRICT ISOLATION 特性句を追加します。

ほとんどの UDF は厳密な分離を必要としません。標準のデータ処理 UDF は、デフォルトの共有分離環境の恩恵を受け、より少ないメモリ消費でより高速に実行されます。

次の STRICT ISOLATION 特性句を UDF に追加します。

  • eval()、exec()、または同様の関数を使用して入力をコードとして実行します
  • ローカルファイルシステムへのファイルの書き込み
  • グローバル変数またはシステム状態の変更
  • 環境変数を変更する

次の例は、入力をコードとして実行し、厳密な分離を必要とする UDF を示しています。

SQL
CREATE OR REPLACE TEMPORARY FUNCTION eval_string(input STRING)
RETURNS STRING
LANGUAGE PYTHON
PARAMETER STYLE PANDAS
HANDLER 'handler_func'
STRICT ISOLATION
AS $$
import pandas as pd
from typing import Iterator

def handler_func(batch_iter: Iterator[pd.Series]) -> Iterator[pd.Series]:
for code_series in batch_iter:
def eval_func(code):
try:
return str(eval(code))
except Exception as e:
return f"Error: {e}"
yield code_series.apply(eval_func)
$$;

バッチ Unity Catalog Python UDFs のサービス認証情報​

バッチ Unity Catalog Python UDF では、Unity Catalog サービス資格情報を使用して外部クラウドサービスにアクセスできます。これは、セキュリティ トークナイザーなどのクラウド関数をデータ処理ワークフローに統合する場合に特に便利です。

注記

サービス資格情報用の UDF 固有の API:
UDF では、 databricks.service_credentials.getServiceCredentialsProvider()を使用してサービス資格情報にアクセスします。

これは、UDF 実行コンテキストでは使用できません、ノートブックで使用されるdbutils.credentials.getServiceCredentialsProvider()関数とは異なります。

サービス資格情報を作成するには、「 サービス資格情報の作成」を参照してください。

使用するサービス資格情報を UDF 定義の CREDENTIALS 句で指定します。

SQL
CREATE OR REPLACE TEMPORARY FUNCTION example_udf(data STRING)
RETURNS STRING
LANGUAGE PYTHON
PARAMETER STYLE PANDAS
HANDLER 'handler_function'
CREDENTIALS (
`credential-name` DEFAULT,
`complicated-credential-name` AS short_name,
`simple-cred`,
cred_no_quotes
)
AS $$
# Python code here
$$;

サービス資格情報のアクセス許可​

コンピュートタイプ全体の作成および呼び出し元の権限要件については、Python UDF でのサービス認証情報の使用を参照してください。

デフォルトの資格情報とエイリアス​

CREDENTIALS 句には複数の認証情報を含めることができますが、DEFAULTとしてマークできるのは 1 つだけです。デフォルト以外の認証情報のエイリアスは、 AS キーワードを使用して作成できます。各資格情報には、一意のエイリアスが必要です。

パッチが適用された Cloud SDK は、デフォルトの認証情報を自動的に取得します。デフォルトの資格情報 は、コンピュートの Spark 設定で指定されたデフォルトよりも優先され、 Unity Catalog UDF 定義に保持されます。

サービス認証情報の例 - Google Cloud Storage​

次の例では、サービス認証情報を使用して、 バッチUnity Catalog Python UDFからGoogle Cloud Storageアクセスします。

Python
%sql
CREATE OR REPLACE FUNCTION main.test.read_gcs_blob(blob_name STRING) RETURNS STRING LANGUAGE PYTHON
PARAMETER STYLE PANDAS
HANDLER 'batchhandler'
CREDENTIALS (
`batch-udf-service-creds-example-cred` DEFAULT
)
ENVIRONMENT (
dependencies = '["google-auth", "google-cloud-storage"]', environment_version = 'None'
)
AS $$
import google.auth # This import is required to enable SDK credential integration
import pandas as pd
from google.cloud import storage


def batchhandler(it):
# The client automatically uses the DEFAULT service credential
client = storage.Client(project="your-project")
bucket = client.bucket("your-bucket")

for blob_names in it:
results = []
for name in blob_names:
blob = bucket.blob(name)
try:
content = blob.download_as_text()
results.append(content)
except Exception as e:
results.append(f"Error: {e}")
yield pd.Series(results)
$$;

登録後に UDF を呼び出します。

SQL
SELECT main.test.read_gcs_blob(blob_name)
FROM VALUES
('config/settings.json'),
('data/input.txt')
AS t(blob_name)

タスク実行コンテキストを取得する​

TaskContext PySpark API を使用して、ユーザーの ID、クラスタータグ、spark ジョブ ID などのコンテキスト情報を取得します。 UDF でタスク コンテキストを取得するを参照してください。

関数が一貫した結果を生成する場合に DETERMINISTIC を設定します​

同じ入力に対して同じ出力を生成する場合は、関数定義に DETERMINISTIC を追加します。これにより、クエリの最適化によりパフォーマンスが向上します。

defaultでは、バッチ Unity Catalog Python UDTF は、明示的に宣言されない限り、非決定論的であると見なされます。非決定論的な関数の例としては、ランダムな値の生成、現在の日時の取得、外部 API 呼び出しの実行などが挙げられます。

CREATE FUNCTION (SQL、Python、Scala、および Java)を参照してください。

制限​

  • Python 関数は NULL 値を個別に処理する必要があり、すべての型マッピングは Databricks SQL 言語マッピングに従う必要があります。
  • バッチ Unity Catalog Python UDF は、セキュリティで保護された分離された環境で実行され、共有ファイル システムや内部サービスにはアクセスできません。
  • ステージ内の複数の UDF 呼び出しはシリアル化され、中間結果がマテリアライズされ、ディスクにスピルする可能性があります。