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 をカタログとスキーマに置き換えます:
%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 を呼び出すことができます。
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 つ以上の
pandas.Seriesを反復処理するイテレータ引数を受け入れます。各pandas.Seriesには、UDF の入力パラメーターが含まれています。 - ジェネレータを反復処理し、データを処理します。
- ジェネレータイテレータを返します。
バッチ 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 に対するイテレータを受け取ります。
%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を反復処理するイテレータを受け取ります。 系列の値は、入力パラメーターと同じ順序です。
%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 回だけ実行されるようにする方法を示しています。
%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 を示しています。
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 句で指定します。
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アクセスします。
%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 を呼び出します。
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 呼び出しはシリアル化され、中間結果がマテリアライズされ、ディスクにスピルする可能性があります。