重要
Scala 和 Java UDF 可以在 Unity 目錄中註冊,以便治理、重用及可發現。 請參閱 Unity 目錄中的 Scala 與 Java 使用者定義函式(UDF)。
本頁說明如何在 Azure Databricks 中建立有會話範圍的 Scala 與 Java UDF。 會話範圍的 UDF 定義於筆記本或工作中,且僅適用於目前的 SparkSession。 關於 SQL 語言的參考,請參見外部使用者定義純量函數(UDFs)。
選擇您的方法
你可以用以下方式定義 Scala 或 Java UDF。 若要比較不同語言、治理和運算環境中的所有 UDF 類型,請參閱 由 Unity Catalog 管理的 UDF 與工作階段範圍 UDF 的比較。
| Approach | Description |
|---|---|
| Inline Scala UDF | 使用 Scala 函式或 lambda 定義工作階段範圍的 UDF。 支援在傳統運算上的 Scala 筆記本中,以及使用環境版本 4 或以上的 無伺服器 JAR 作業 中。 |
| 從 JAR 轉化的 Java UDF | 使用 spark.udf.registerJavaFunction 從 JAR 檔註冊預先編譯的工作階段範圍 UDF 類別。 支援 無伺服器 JAR 作業。 |
| 由 Unity Catalog 控管的 Scala 或 Java UDF | 在 Unity 目錄中註冊 UDF,以便治理、重用及可發現。 支援無伺服器運算。 |
需求
- 在啟用 Unity Catalog 的運算資源上執行 Scala UDF 時,若使用標準存取模式,則需要 Databricks Runtime 14.2 或更新版本。
- 在支援 Unity 目錄的叢集上支援 Scala UDF 的 ARM 實例,需要 Databricks Runtime 15.2 或以上版本。
- 從 JAR
spark.udf.registerJavaFunction註冊 Java UDF 需要 Databricks Runtime 18 LTS 或以上版本。 請參閱 從 JAR 註冊 Java UDF。 - 若要從工作階段範圍的 Scala UDF 存取 Unity Catalog 秘密,必須在具有標準存取模式的經典運算上使用 Databricks Runtime 19 或更新版本。 在無伺服器 JAR 工作中,秘密存取需要環境 版本 6 或以上。 請參閱 UDF 的要求與權限。
重要
請用與執行 JAR 的計算相同的 Scala 和 Apache Spark 版本來編譯你的 JAR。 不匹配可能導致 UDF 在註冊或呼叫時失效。
- 傳統運算:使 Scala 和 Spark 版本與您的 Databricks Runtime 版本相符。 請參閱 Databricks 執行時版本說明中的系統環境部分,了解版本與相容性。 例如,Databricks Runtime 18 LTS 使用 Scala 2.13.16 和 Apache Spark 4.0。
- 無伺服器運算:請使 Scala 版本與環境版本一致。 請參閱 環境版本。
將 Apache Spark 相依標記為 provided ,這樣它就不會被綁在你的 JAR 裡。 只包含您的 UDF 所使用的第三方相依性。
將函式註冊為 UDF
將 Scala 函式註冊為 UDF,使用:spark.udf.register
val squared = (s: Long) => {
s * s
}
spark.udf.register("square", squared)
在 Spark SQL 中呼叫 UDF
建立一個臨時視圖,然後在 SQL 查詢中呼叫 UDF:
spark.range(1, 20).createOrReplaceTempView("test")
%sql select id, square(id) as id_squared from test
使用 UDF 與 DataFrames 搭配
你也可以使用 DataFrame API 呼叫 UDF:
import org.apache.spark.sql.functions.{col, udf}
val squared = udf((s: Long) => s * s)
display(spark.range(1, 20).select(squared(col("id")) as "id_squared"))
從會話範圍的 Scala UDF 存取 Unity 目錄的秘密
工作階段範圍的 Scala UDF 可以透過 com.databricks.Secrets.get(catalog, schema, key) 擷取 Unity Catalog 密鑰。 秘密存取則是利用呼叫者的權限。 請參閱 UDF 要求與權限 以取得所需權限。
以下範例會檢查是否已設定秘密,而不會回傳其值:
import org.apache.spark.sql.functions.udf
val secretIsConfigured = udf((_: Int) => {
val apiKey = com.databricks.Secrets.get("main", "default", "api_key")
apiKey != null && apiKey.nonEmpty
})
spark.udf.register("secret_is_configured", secretIsConfigured)
spark.sql("SELECT secret_is_configured(1)").show()
請勿從 UDF 回傳機密值。 詳見 Unity 目錄中的秘密。
帶有 UDF 的檔案
重要
這項功能位於 測試版 (Beta) 中。 工作區管理員可以從 「預覽 」頁面控制對此功能的存取。 請參閱 管理 Azure Databricks 預覽。
Scala 檔案的類型為 FileRef。 可將其用作 UDF 中的參數或回傳型別,無論是頂層型態或巢狀型態。 關於類型及其巢狀規則,請參見 FILE 類型。
要在 UDF 中讀取檔案內容,請呼叫以下 FileRef其中一個:
-
asLocalFile():回傳一個java.io.File,可將其傳遞給任何接受路徑的程式庫。 -
open():回傳一個java.io.InputStream,用來讀取檔案的位元組。 來電者關上了門。
要在 UDF 中產生新資料 FileRef ,請呼叫以下靜態方法之一:
-
FileRef.create(uri):在uri建立該檔案的參考。 -
FileRef.fromBytes(bytes, destinationPath, contentType):將bytes作為外部檔案上傳至destinationPath磁碟區路徑,並傳回參考。 -
FileRef.fromLocalFile(localFile, destinationPath, contentType): 將本地檔案上傳到destinationPath卷路徑,作為外部檔案並回傳參考。
從寫入FILE MANAGED欄位的 UDF 回傳 a FileRef 是不被支援的。
關於 Python、Scala 和 SQL 中的範例,包括影像處理、檔案類型偵測及影片影格擷取,請參見帶有 UDF 的 Process 檔案。
從 JAR 註冊 Java UDF
將 UDF 打包成 JAR,使用 spark.addArtifact 將其加入你的工作階段,並使用 spark.udf.registerJavaFunction 註冊 UDF 類別。
注意
支援標準存取模式及 Databricks Runtime 18 LTS 以上的無伺服器運算。 已註冊的函式屬於工作階段範圍,且未在 Unity Catalog 中註冊。
接下來的步驟會帶你建立專案、撰寫 UDF 類別、建立 fat JAR,以及註冊它。
步驟一:建立你的專案
用 Scala 或 Java 建立專案。
Scala
請使用以下方式 sbt建立新的 Scala 專案:
sbt new scala/scala-seed.g8
將您的 build.sbt 檔案內容替換為以下內容。 設定 scalaVersion,並將 spark-sql 的版本設為與你的運算環境相符:
scalaVersion := "2.13.16"
ThisBuild / organization := "com.example"
lazy val myUDF = (project in file("."))
.settings(
name := "my-udf",
libraryDependencies += "org.apache.spark" %% "spark-sql" % "4.0.0" % "provided"
)
啟用 sbt-assembly 外掛程式以建置胖 JAR 檔。 建立或編輯 project/assembly.sbt 並新增:
addSbtPlugin("com.eed3si9n" % "sbt-assembly" % "2.0.0")
JAVA
使用快速入門原型建立一個新的 Maven 專案:
mvn archetype:generate \
-DgroupId=com.example \
-DartifactId=my-udf \
-DarchetypeArtifactId=maven-archetype-quickstart \
-DinteractiveMode=false
此指令建立標準的 Maven 專案結構,包含 src/main/java 和 src/test/java 目錄。
在產生的 pom.xml 中,於 <project></project> 標籤內加入 <properties> 區塊,並設定 maven-shade-plugin 以建置 fat JAR:
<properties>
<maven.compiler.source>17</maven.compiler.source>
<maven.compiler.target>17</maven.compiler.target>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
</properties>
<build>
<plugins>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-shade-plugin</artifactId>
<version>3.5.0</version>
<executions>
<execution>
<phase>package</phase>
<goals>
<goal>shade</goal>
</goals>
</execution>
</executions>
</plugin>
</plugins>
</build>
步驟 2:撰寫您的 UDF 類別
你的 UDF 類別必須實作其中一個 org.apache.spark.sql.api.java.UDF 介面(UDF1 透過 UDF22),其中數字代表 UDF 接收的輸入參數數量。 用你的邏輯實作這個 call() 方法。
處理常式必須是 Java 類別。
spark.udf.registerJavaFunction 透過反射載入類別,因此它必須是頂層(或 static 巢狀)的公共類別,並有一個公開的無 arg 建構子。 Scala class 或 object 不符合此要求,並在呼叫時失敗。 你可以用 SBT 編譯 JAR,但 UDF 類別本身必須用 Java 撰寫。
創建 src/main/java/com/example/MyIntegerUDF.java:
package com.example;
import org.apache.spark.sql.api.java.UDF1;
public class MyIntegerUDF implements UDF1<Integer, Integer> {
@Override
public Integer call(Integer x) {
return x + 1;
}
}
步驟 3:建立 fat JAR
將編譯好的 UDF 打包成 fat JAR 檔。
Scala
從你的專案根目錄執行:
sbt clean assembly
fat JAR 會在 target/scala-2.13/ 中建立,名稱類似 my-udf-assembly-0.1.0-SNAPSHOT.jar。
JAVA
從你的專案根目錄執行:
mvn clean package
fat JAR 會在 target/ 中建立,名稱類似 my-udf-1.0-SNAPSHOT.jar。
步驟 4:將您的 JAR 上傳到 Unity 目錄卷
把 JAR 上傳到 Unity 目錄卷 ,讓你的電腦能存取它。 如果你還沒有一個磁碟區,請建立一個:
CREATE VOLUME IF NOT EXISTS my_catalog.my_schema.udf_jars
COMMENT 'Storage for UDF JAR files';
使用 目錄瀏覽器將您的 JAR 檔案上傳至卷中:
- 在您的 Azure Databricks 工作區中,按一下
以開啟目錄總管。
- 選擇目錄,再選擇包含你卷的結構。
- 點擊卷名。
- 點擊 「上傳到此卷 」並選擇你的 JAR 檔案。
- 按一下 [上傳] 。
- 上傳完成後,點選你的 JAR 檔案名稱,然後點選 複製 路徑以複製磁碟區路徑。 例如:
/Volumes/my_catalog/my_schema/udf_jars/my-udf-assembly-0.1.0-SNAPSHOT.jar。 你下一步需要這條路。
步驟五:註冊並致電UDF
將 JAR 以其磁碟區路徑新增至工作階段、註冊 UDF 類別,然後從 Spark SQL 中呼叫該 UDF:
# Add the JAR containing your UDF class to the session
spark.addArtifact("/Volumes/my_catalog/my_schema/udf_jars/my-udf-assembly-0.1.0-SNAPSHOT.jar")
# Register the UDF class, providing the SQL function name,
# the fully qualified class name, and the return type
from pyspark.sql.types import IntegerType
spark.udf.registerJavaFunction(
"my_udf",
"com.example.MyIntegerUDF",
IntegerType(),
)
# Call the UDF from Spark SQL
spark.sql("SELECT my_udf(21)").show()
在無伺服器和標準存取模式運算中,你必須明確指定回傳型別。 省略回傳類型則失敗。UC_COMMAND_NOT_SUPPORTED_IN_SHARED_ACCESS_MODE 使用者自訂彙總函數(UDAFs)不支援與 registerJavaFunction 搭配使用。
查詢會回傳 UDF 輸出,確認函式已註冊且可呼叫:
+----------+
| my_udf(21)|
+----------+
| 22|
+----------+
評估順序和 Null 檢查
Spark SQL(包括 SQL 以及 DataFrame 和 Dataset API)並不保證子表達式評估的順序。 Spark 不會從左到右評估運算子或函數的輸入。 邏輯 AND 與 OR 表達式沒有從左到右短路的語意。
不要依賴布林運算式的副作用、評估順序,或 WHERE 子句和 HAVING 子句的順序。 查詢優化器可以重新排序這些表達式和子句。 若 UDF 依賴短路語意來進行空檢查,Spark 並不保證空檢查會在 UDF 之前執行。 例如:
spark.udf.register("strlen", (s: String) => s.length)
spark.sql("select s from test1 where s is not null and strlen(s) > 1") // no guarantee
此 WHERE 條款並不保證 Spark 會在篩除空值後呼叫 strlen UDF。
Databricks 建議使用下列任一方式來處理空值檢查:
- 讓 UDF 本身具備 null 感知,並在 UDF 內部進行 null 檢查
- 使用
IF或CASE WHEN運算式執行 Null 檢查,並在條件分支中叫用 UDF
spark.udf.register("strlen_nullsafe", (s: String) => if (s != null) s.length else -1)
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
型別化資料集 API
注意
在 Databricks Runtime 15.4 和更新版本具有標準存取模式的 Unity 目錄啟用叢集上支援此功能。
使用類型化的資料集 API,在資料集上執行映射、篩選及彙整等轉換,並搭配使用者自訂函式。
以下範例使用 map() API 將結果欄位中的數字修改為前綴字串:
spark.range(3).map(f => s"row-$f").show()
此範例使用 map(),但同樣的模式也適用於其他型別化的資料集 API,如 filter()、 mapPartitions()、 foreach()foreachPartition()reduce()flatMap()。
Scala UDF 功能和 Databricks 執行環境兼容性
下列功能在已啟用 Unity Catalog 且使用標準(共用)存取模式的叢集上,要求 Databricks Runtime 的最低版本。
| 特徵 / 功能 | Databricks 執行時間最低版本 |
|---|---|
| 純量 UDF | Databricks Runtime 14.2(Databricks 執行環境 14.2) |
Dataset.map、Dataset.mapPartitions、Dataset.filter、Dataset.reduce、Dataset.flatMap |
Databricks 執行環境 15.4 |
KeyValueGroupedDataset.flatMapGroups、KeyValueGroupedDataset.mapGroups |
Databricks 執行環境 15.4 |
(串流) foreachWriter Sink |
Databricks 執行環境 15.4 |
(串流) foreachBatch |
Databricks 執行環境 16.1 |
(串流) KeyValueGroupedDataset.flatMapGroupsWithState |
Databricks 執行環境 16.2 |
spark.udf.registerJavaFunction(來自 JAR 檔的 Java UDF) |
Databricks Runtime 18 LTS |