Python スカラーユーザー定義関数 (UDF)
Python スカラー UDF を使用すると、Databricks 上の SQL クエリー内でカスタムの Python ロジックを実行できます。このページでは、それらの登録と呼び出し方法、サービス認証情報とシークレットの使用方法、および Spark SQL での部分式の評価順序に関する注意事項の処理方法について説明します。
要件
-
Databricks Runtime14.0 以前では、標準アクセスPython PandasUnity Catalogモードを使用する クラスターでは、 UDF と UDF はサポートされていません。
-
スカラー Python UDF と Pandas UDF は、Databricks Runtime 14.1 以降のすべてのアクセス モードでサポートされています。
-
Databricks Runtime 14.1 以降では、SQL構文を使用してスカラーPython UDFをUnity Catalogに登録できます。Unity Catalog の SQL および Python ユーザー定義関数 (UDF)を参照してください。
-
Unity Catalog 対応クラスター上の Python UDF の ARM インスタンス サポートには、Databricks Runtime 15.2 以上が必要です。
Databricks Runtime14.0 以前では、標準アクセスPython PandasUnity Catalogモードを使用する クラスターでは、 UDF と UDF はサポートされていません。スカラー Python UDF と Pandas UDF は、Databricks Runtime 14.1 以降のすべてのアクセス モードでサポートされています。
Databricks Runtime 14.1 以降では、SQL構文を使用してスカラーPython UDFをUnity Catalogに登録できます。Unity Catalog の SQL および Python ユーザー定義関数 (UDF)を参照してください。
関数を UDF として登録する
def squared(s):
return s * s
spark.udf.register("squaredWithPython", squared)
オプションで、UDF の戻り値の型を設定できます。 デフォルトの戻り値の型は StringTypeです。
from pyspark.sql.types import LongType
def squared_typed(s):
return s * s
spark.udf.register("squaredWithPython", squared_typed, LongType())
Spark SQL で UDF を呼び出す
spark.range(1, 20).createOrReplaceTempView("test")
%sql select id, squaredWithPython(id) as id_squared from test
UDF と DataFrames の併用
from pyspark.sql.functions import udf
from pyspark.sql.types import LongType
squared_udf = udf(squared, LongType())
df = spark.table("test")
display(df.select("id", squared_udf("id").alias("id_squared")))
または、アノテーション構文を使用して同じ UDF を宣言することもできます。
from pyspark.sql.functions import udf
@udf("long")
def squared_udf(s):
return s * s
df = spark.table("test")
display(df.select("id", squared_udf("id").alias("id_squared")))
UDF を使用したバリアント
バリアントの PySpark 型はVariantTypeであり、値の型はVariantValです。バリアントに関する情報については、 「バリアント データのクエリ」を参照してください。
from pyspark.sql.functions import col, lit, udf
from pyspark.sql.types import VariantType, VariantVal
# Return Variant
@udf(returnType = VariantType())
def toVariant(jsonString):
return VariantVal.parseJson(jsonString)
spark.range(1).select(lit('{"a" : 1}').alias("json")).select(toVariant(col("json"))).display()
+---------------+
|toVariant(json)|
+---------------+
| {"a":1}|
+---------------+
from pyspark.sql.functions import col, lit, udf
from pyspark.sql.types import StructField, StructType, VariantType, VariantVal
# Return Struct<Variant>
@udf(returnType = StructType([StructField("v", VariantType(), True)]))
def toStructVariant(jsonString):
return {"v": VariantVal.parseJson(jsonString)}
spark.range(1).select(lit('{"a" : 1}').alias("json")).select(toStructVariant(col("json"))).display()
+---------------------+
|toStructVariant(json)|
+---------------------+
| {"v":{"a":1}}|
+---------------------+
from pyspark.sql.functions import col, lit, udf
from pyspark.sql.types import ArrayType, VariantType, VariantVal
# Return Array<Variant>
@udf(returnType = ArrayType(VariantType()))
def toArrayVariant(jsonString):
return [VariantVal.parseJson(jsonString)]
spark.range(1).select(lit('{"a" : 1}').alias("json")).select(toArrayVariant(col("json"))).display()
+--------------------+
|toArrayVariant(json)|
+--------------------+
| [{"a":1}]|
+--------------------+
from pyspark.sql.functions import col, lit, udf
from pyspark.sql.types import MapType, StringType, VariantType, VariantVal
# Return Map<String, Variant>
@udf(returnType = MapType(StringType(), VariantType(), True))
def toMapVariant(jsonString):
return {"v1": VariantVal.parseJson(jsonString), "v2": VariantVal.parseJson("[" + jsonString + "]")}
spark.range(1).select(lit('{"a" : 1}').alias("json")).select(toMapVariant(col("json"))).display()
+-----------------------------+
| toMapVariant(json)|
+-----------------------------+
|{"v2":[{"a":1}],"v1":{"a":1}}|
+-----------------------------+
UDF を含むファイル
ベータ版
この機能はベータ版です。ワークスペース管理者は、 プレビュー ページからこの機能へのアクセスを制御できます。Databricksのプレビューを管理するを参照してください。
ファイルの PySpark タイプは FileType です。UDF のパラメーターまたは戻り値の型として、トップレベルの型またはネストされた型として使用します。型、そのネスト規則、および FileRef API については、ファイルタイプを参照してください。
UDF でファイルの内容を読み取るには、file.as_local_file() を呼び出して開くことができるローカルパスを取得するか、file.open() を呼び出してそのバイトをストリームとして読み取ります。画像処理、ファイルタイプ検出、ビデオフレーム抽出など、Python、Scala、SQL の例については、UDF を使用したファイルの処理を参照してください。タイプのリファレンスについては、FILE タイプを参照してください。
Unity Catalogに登録された UDF (CREATE FUNCTION) は FILE のメタデータを読み取ることができますが、コンテンツは読み取れず、ファイルを作成することもできません。セッションスコープの UDF を使用して、ファイルの内容 (open, as_local_file) を読み取るか、ファイルを作成します (from_bytes, from_local_file)。
UDF 本文では、各 FILE 値は FileRef オブジェクトです。
from pyspark.sql.functions import col, udf
from pyspark.sql.types import FileRef, StringType
@udf(returnType=StringType())
def file_content_type(file: FileRef) -> str:
return file.content_type
df = spark.table("documents")
display(df.select(col("file").uri, file_content_type(col("file"))))
評価順序と null チェック
Spark SQL ( SQL 、 データフレーム 、データセット APIを含む)は、
部分式の評価。 特に、演算子または関数の入力はそうではありません
必ず左から右に評価するか、その他の固定された順序で評価されます。 たとえば、論理 AND
また、 OR 式には、左から右への "短絡" セマンティクスはありません。
そのため、Booleanの副作用や評価の順番に頼るのは危険です
式、および WHERE 節と HAVING 節の順序は、
クエリの最適化と計画中に並べ替えられます。具体的には、UDF が null チェックのために SQL の短絡セマンティクスに依存している場合、
UDF を呼び出す前にヌル チェックが行われることを保証します。例えば
spark.udf.register("strlen", lambda s: len(s), "int")
spark.sql("select s from test1 where s is not null and strlen(s) > 1") # no guarantee
この WHERE 句は、null をフィルターで除外した後に strlen UDF が呼び出されることを保証するものではありません。
適切な null チェックを実行するには、次のいずれかを実行することをお勧めします。
- UDF 自体を null 対応にし、UDF 自体の内部で null チェックを行います
IF式またはCASE WHEN式を使用してヌルチェックを行い、条件分岐でUDFを呼び出します
spark.udf.register("strlen_nullsafe", lambda s: len(s) if not s is None else -1, "int")
spark.sql("select s from test1 where s is not null and strlen_nullsafe(s) > 1") # ok
spark.sql("select s from test1 where if(s is not null, strlen(s), null) > 1") # ok
Unity Catalog のシークレットへのアクセス
セッションスコープの Python UDF から Unity Catalog のシークレットにアクセスするには、セッションスコープの Python UDF でシークレットを使用するを参照してください。スカラーまたはバッチ Unity Catalog Python UDF から宣言されたシークレットにアクセスするには、Python UDF でシークレットを使用するを参照してください。
Python UDF におけるサービス資格情報
セッションスコープの Python スカラー UDF およびスカラー Unity Catalog Python UDF では、Unity Catalog サービス認証情報を使用して外部クラウドサービスに安全にアクセスできます。これは、クラウドベースのトークン化、暗号化、シークレット管理などの操作をデータ変換に直接統合する場合に便利です。
要件は、UDF の種類とコンピュートによって異なります。Python UDF でのサービス資格情報の使用を参照してください。
サービス資格情報を作成するには、「 サービス資格情報の作成」を参照してください。
セッションスコープのスカラー Python UDF でのサービス認証情報の使用
サービス認証情報にアクセスするには、UDF ロジックの databricks.service_credentials.getServiceCredentialsProvider() ユーティリティを使用して、適切な認証情報でクラウド SDK を初期化します。すべてのコードは UDF 本体にカプセル化する必要があります。
@udf
def use_service_credential():
from google.cloud import storage
# Assuming there is a service credential named 'testcred' set up in Unity Catalog
client = storage.Client(project='your-project', credentials=getServiceCredentialsProvider('testcred'))
# Use the client to perform operations
サービス資格情報のアクセス許可
セッションスコープの UDF は、呼び出し元の権限を使用します。必要な権限については、Python UDF でのサービス認証情報の使用を参照してください。
セッションスコープの UDF のコンピュートレベルのdefault認証情報
スカラー Python UDF で使用される場合、Databricks はコンピュート環境変数からdefaultのサービス認証情報を自動的に使用します。この動作により、UDF コード内で認証情報のエイリアスを明示的に管理することなく、外部サービスを安全に参照できます。コンピュートリソースのdefaultサービス認証情報の指定を参照してください
Default資格情報のサポートは、StandardおよびDedicatedアクセスモードのクラスターでのみ利用できます。SQLウェアハウスでは使用できません。
@udf
def use_service_credential():
import google.auth # This import is required to enable SDK credential integration
from google.cloud import storage
# The client automatically uses the default service credential
client = storage.Client(project='your-project')
# Use the client to perform operations
スカラー Unity Catalog Python UDF でのサービス資格情報の使用
UDF 定義の CREDENTIALS 句でサービス資格情報を指定します。1 つの資格情報を DEFAULT としてマークし、パッチ適用済みクラウド SDK がそれを自動的に使用するようにすることができます。クラシック コンピュートでは、この機能には Databricks Runtime 18.1 以降が必要です。サーバーレス コンピュート、ならびに Pro および Serverless SQL Warehouse では、UDF の environment_version を 6 以上に明示的に設定します。コンピュート、ネットワーク、および権限の要件の詳細については、Python UDF でのサービス資格情報の使用を参照してください。
タスク実行コンテキストを取得する
TaskContext PySpark API を使用して、ユーザーの ID、クラスタータグ、spark ジョブ ID などのコンテキスト情報を取得します。 UDF でタスク コンテキストを取得するを参照してください。
制限
PySpark UDF には、次の制限が適用されます。
-
ファイルアクセス制限: Databricks Runtime 14.2 以前では、共有クラスタリング上の PySpark UDF は、Git フォルダ、ワークスペース ファイル、または Unity Catalog ボリュームにアクセスできません。
-
ブロードキャスト変数: 標準アクセス モード クラスター と サーバレス コンピュートの PySpark UDF は、ブロードキャスト変数をサポートしていません。
-
サーバレスのメモリ制限 : サーバレス コンピュートの PySpark UDF には、 PySpark UDFあたり 1GB のメモリ制限があります。 この制限を超えると、タイプ UDF_PYSPARK_USER_CODE_ERROR のエラーが発生します 。MEMORY_LIMIT_SERVERLESS。
-
標準アクセスモードのメモリ制限 : 標準アクセスモードの PySpark UDF には、選択したインスタンスタイプの使用可能なメモリに基づいてメモリ制限があります。使用可能なメモリを超えると、タイプ UDF_PYSPARK_USER_CODE_ERROR のエラーが発生します 。MEMORY_LIMIT。
-
サーバレスSQLウェアハウスのネットワークアクセス : 一応、サーバレスSQLウェアハウスのPython UDFは送信ネットワークリクエストを行うことができず、ネットワーク呼び出しを試みるクエリは無期限にハングします。 アウトバウンド ネットワーク アクセスを有効にするには、ワークスペースの [プレビュー]ページ で、パブリック プレビュー機能[サーバーレスSQL ウェアハウスの分離されたワークロードのネットワークを有効にする] を有効にします。それ以外の場合は、ネットワーク アクセスを必要とする UDF にはサーバレス コンピュートまたはクラシック コンピュートを使用します。