Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
На этой странице описано, как создавать пользовательские функции Scala и Java (UDF), регистрировать их в Unity Catalog и использовать их в разных вычислительных средах. UDFS каталога Unity позволяют повторно использовать существующую логику JVM с управлением каталогом Unity и элементами управления доступом.
В отличие от сеансовых UDF Scala, которые ограничены одним ноутбуком или кластером, зарегистрированные UDF в Unity Catalog:
- Управляемый: управляется с помощью разрешений Unity Catalog и средств управления доступом.
- Повторное использование: совместное использование между командами, записными книжками, заданиями и хранилищами SQL.
- Обнаружение: отображается в обозревателе каталогов и системных таблицах.
- Изолированный: запуск в песочницах с однократной стоимостью холодного запуска на сеанс. Последующие вызовы выполняются быстро.
Requirements
Рабочая область должна быть активирована для Unity Catalog. Применяются следующие дополнительные требования.
Вычисления. Поддерживаются все типы вычислений, включая бессерверные записные книжки и задания, хранилища SQL и декларативные конвейеры Spark в Lakeflow. Для классических вычислений требуется Databricks Runtime 18.2 или более поздней версии. В бессерверных вычислительных ресурсах и хранилищах SQL определение UDF должно указывать среду версии 4 или более поздней в environment_version поле. Это требование относится к определению UDF, а не к ноутбуку, из которого она вызывается, или заданию. См. версии бессерверных сред .
Разработка:
- Scala: 2.13.16. Scala 2.12 не поддерживается.
- JDK: 17.
- Упаковка: толстый JAR-файл, содержащий все сторонние зависимости, используемые UDF.
Разрешения:
- Создайте UDF:
USAGEиCREATE FUNCTIONв схеме, аUSAGE— в каталоге. - Запустите UDF:
EXECUTEдля функции, иUSAGEдля схемы и каталога. - Откройте JAR-файл:
READ VOLUMEв томе, где хранится JAR.
Дополнительные сведения о разрешениях каталога Unity см. в разделе "Управление привилегиями" в каталоге Unity .
Создание JAR-файла UDF
Упаковайте скомпилированный код в виде JAR-файла и отправьте его в том каталога Unity перед регистрацией UDF. Выберите метод сборки:
Локальная сборка
Выполните следующие действия, чтобы создать fat JAR с помощью локальной среды разработки.
Настройка среды
Установите необходимые средства на локальном компьютере. Следующие команды предназначены для macOS. Для других платформ установите JDK 17 и sbt (Scala) или Maven (Java) с помощью диспетчера пакетов платформы.
Scala
Установите JDK 17 и sbt:
brew install openjdk@17
brew install sbt
Проверьте установку:
java -version # Should show Java 17
sbt --version # Should show sbt version
Java
Установите JDK 17 и Maven:
brew install openjdk@17
brew install maven
Проверьте установку:
java -version # Should show Java 17
mvn --version # Should show Maven version
Создание проекта
Настройте проект в Scala или Java.
Scala
Создание проекта Scala с помощью sbt:
sbt new scala/scala-seed.g8
При появлении запроса введите имя проекта (например, my-udf-project).
Настройте build.sbt
Замените содержимое build.sbt файла следующей конфигурацией:
scalaVersion := "2.13.16"
ThisBuild / organization := "com.example"
lazy val myUDF = (project in file("."))
.settings(
name := "my-udf"
)
Включите плагин sbt-assembly
Создайте или измените project/assembly.sbt и добавьте:
addSbtPlugin("com.eed3si9n" % "sbt-assembly" % "2.0.0")
Этот плагин создаёт толстый JAR-файл со всеми зависимостями.
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
В созданном pom.xml файле в <project></project> тегах добавьте <properties> блок со следующей конфигурацией:
<properties>
<maven.compiler.source>17</maven.compiler.source>
<maven.compiler.target>17</maven.compiler.target>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
</properties>
Кроме того, в <project></project> тегах добавьте блок со следующей <build> конфигурацией:
<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>
maven-shade-plugin создает толстый JAR-файл, содержащий все зависимости.
Напишите свою UDF
При написании UDF обратитесь к типам данных, чтобы узнать о поддерживаемых типах данных, и к сопоставлениям языков, чтобы увидеть, как типы в Scala и Java сопоставляются с типами SQL.
Обработчик UDF должен соответствовать следующим требованиям:
-
Scala: определение обработчика в качестве метода в объекте
object(а неclass). ЗначениеHANDLERсоответствует методу в Scalaobject. -
Java: Определите обработчик как метод
public static. -
Сигнатура: типы параметров метода, порядок и возвращаемый тип должны соответствовать списку аргументов и
RETURNSтипу в инструкцииCREATE FUNCTION. - Только скаляр: обработчик должен возвращать одно скалярное значение. Типы возвращаемых таблиц не поддерживаются.
- Автономный: обработчик должен работать только с его входными аргументами. Он не может использовать API Spark или зависеть от основных пакетов Spark. См. Ограничения.
Note
В Scala обработчик с примитивным типом параметра (например, Int) пропускается и возвращает NULL, если любой входной аргумент имеет значение SQL NULL. Чтобы получить и обработать значения, заключите NULL параметр в Option, например Option[Int].
Scala
Создайте объект Scala в src/main/scala/com/example/MyUDF.scala и определите функцию UDF.
Базовый пример
package com.example
object MyUDF {
def addOne(x: Int): Int = x + 1
}
Пример с внешней зависимостью
Чтобы использовать внешние библиотеки, добавьте их в build.sbt файл:
scalaVersion := "2.13.16"
ThisBuild / organization := "com.example"
lazy val myUDF = (project in file("."))
.settings(
name := "currency-udf",
libraryDependencies ++= Seq(
"org.apache.commons" % "commons-lang3" % "3.12.0"
)
)
Затем используйте зависимость в UDF:
package com.example
import org.apache.commons.lang3.StringUtils
object CurrencyUDF {
private val rates: Map[String, Double] = Map(
"USD" -> 1.0,
"EUR" -> 1.1,
"GBP" -> 1.3,
"JPY" -> 0.007
)
def convertToUSD(price: Double, currency: String): Double = {
require(currency != null, "Currency must not be null")
val normalizedCurrency = StringUtils.upperCase(currency)
rates.get(normalizedCurrency) match {
case Some(rate) => price * rate
case None => throw new IllegalArgumentException(s"Unsupported currency: $currency")
}
}
}
Перед развертыванием проверьте UDF с помощью модульных тестов. См. Локальное тестирование пользовательских функций.
Java
Создайте класс Java в src/main/java/com/example/MyUDF.java и определите UDF как открытый статический метод.
Базовый пример
package com.example;
public class MyUDF {
public static int addOne(int x) {
return x + 1;
}
}
Пример с внешней зависимостью
Чтобы использовать внешние библиотеки, добавьте их в <dependencies> раздел pom.xml файла:
<dependencies>
<dependency>
<groupId>org.apache.commons</groupId>
<artifactId>commons-lang3</artifactId>
<version>3.12.0</version>
</dependency>
</dependencies>
Затем используйте эту зависимость в UDF:
package com.example;
import org.apache.commons.lang3.StringUtils;
import java.util.Map;
import java.util.HashMap;
public class CurrencyUDF {
private static final Map<String, Double> rates = new HashMap<>();
static {
rates.put("USD", 1.0);
rates.put("EUR", 1.1);
rates.put("GBP", 1.3);
rates.put("JPY", 0.007);
}
public static double convertToUSD(double price, String currency) {
if (currency == null) {
throw new IllegalArgumentException("Currency must not be null");
}
String normalizedCurrency = StringUtils.upperCase(currency);
if (!rates.containsKey(normalizedCurrency)) {
throw new IllegalArgumentException("Unsupported currency: " + currency);
}
return price * rates.get(normalizedCurrency);
}
}
Перед развертыванием проверьте UDF с помощью модульных тестов. См. Локальное тестирование пользовательских функций.
Note
Ваша UDF выполняется в изолированной песочнице без активного сеанса Spark, поэтому она не может использовать API Spark внутри функции. Например, нельзя создавать объекты DataFrame или Dataset, запускать spark.sql(...) или получать доступ к SparkSession или SparkContext. UDF должен быть автономной логикой по его входным аргументам. Он также не может зависеть от основных пакетов Spark.
Создание толстых JAR-файлов
Создайте проект, чтобы создать толстый 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.
Загрузите JAR-файл в том Unity Catalog
Если у вас еще нет тома каталога Unity, создайте его:
CREATE VOLUME IF NOT EXISTS my_catalog.my_schema.udf_jars
COMMENT 'Storage for UDF JAR files';
Если другим пользователям необходимо запустить UDF, предоставьте их READ VOLUME на томе:
GRANT READ VOLUME ON VOLUME my_catalog.my_schema.udf_jars TO `user@example.com`;
Отправьте JAR-файл в том с помощью обозревателя каталогов:
- В рабочей области Azure Databricks щелкните
Catalog, чтобы открыть обозреватель каталогов.
- Выберите каталог, а затем выберите схему, содержащую том.
- Щелкните имя тома.
- Нажмите кнопку "Отправить в этот том " и выберите JAR-файл.
- Нажмите кнопку Отправить.
- После завершения загрузки нажмите на имя JAR-файла.
- Нажмите Копировать путь, чтобы скопировать путь тома в буфер обмена. Например,
/Volumes/my_catalog/my_schema/udf_jars/my-udf-assembly-0.1.0-SNAPSHOT.jar(Scala) или/Volumes/my_catalog/my_schema/udf_jars/my-udf-1.0-SNAPSHOT.jar(Java). Этот путь требуется при регистрации UDF.
Создать в блокноте
Вы можете скомпилировать UDF, упаковать его в JAR-файл и загрузить его в том Unity Catalog непосредственно из записной книжки Azure Databricks. Этот подход применим к небольшим пользовательским функциям без зависимостей. Для пользовательских функций со сторонними библиотеками используйте Собирать локально.
Следующая Python ячейка записывает Java UDF, которая очищает строку (обрезает пробелы, сверчивает повторяющиеся пробелы и строчные регистры), компилирует его с помощью JDK 17, упаковает его в виде JAR-файла и копирует его в том каталога Unity. Измените volume_path, чтобы он указывал на существующий том, на который у вас есть разрешение WRITE VOLUME.
import os
import subprocess
import shutil
build_dir = "/tmp/udf_build"
package_dir = f"{build_dir}/src/com/databricks/udf"
classes_dir = f"{build_dir}/classes"
os.makedirs(package_dir, exist_ok=True)
os.makedirs(classes_dir, exist_ok=True)
# The UDF handler: a public static method on a plain Java class.
# The doubled backslashes produce a single backslash in the Java source (\\s+).
udf_code = """package com.databricks.udf;
public class StringCleanUDF {
public static String clean(String input) {
if (input == null) return null;
return input.trim().replaceAll("\\\\s+", " ").toLowerCase();
}
}
"""
with open(f"{package_dir}/StringCleanUDF.java", "w") as f:
f.write(udf_code)
# Compile with JDK 17 to match Environment Version 4.
subprocess.run(
["javac", "--release", "17", "-d", classes_dir, f"{package_dir}/StringCleanUDF.java"],
check=True,
)
# Package the compiled class into a JAR.
jar_path = f"{build_dir}/string_clean_udf.jar"
subprocess.run(["jar", "cf", jar_path, "-C", classes_dir, "."], check=True)
# Copy the JAR to a Unity Catalog volume.
volume_path = "/Volumes/my_catalog/my_schema/udf_jars/string_clean_udf.jar"
os.makedirs(os.path.dirname(volume_path), exist_ok=True)
shutil.copy2(jar_path, volume_path)
print(f"JAR uploaded to: {volume_path}")
После того как JAR-файл находится в томе, зарегистрируйте UDF. Используйте LANGUAGE JAVA и задайте для HANDLER полное имя метода, например com.databricks.udf.StringCleanUDF.clean.
Регистрация UDF в каталоге Unity
После сборки и отправки JAR-файла используйте инструкцию CREATE FUNCTION для регистрации UDF в каталоге Unity.
Scala
CREATE OR REPLACE FUNCTION my_catalog.my_schema.add_one(x INT)
RETURNS INT
LANGUAGE SCALA
DETERMINISTIC
ENVIRONMENT (
java_dependencies = '["/Volumes/my_catalog/my_schema/udf_jars/my-udf-assembly-0.1.0-SNAPSHOT.jar"]',
environment_version = '4'
)
HANDLER 'com.example.MyUDF.addOne';
Java
CREATE OR REPLACE FUNCTION my_catalog.my_schema.add_one(x INT)
RETURNS INT
LANGUAGE JAVA
DETERMINISTIC
ENVIRONMENT (
java_dependencies = '["/Volumes/my_catalog/my_schema/udf_jars/my-udf-1.0-SNAPSHOT.jar"]',
environment_version = '4'
)
HANDLER 'com.example.MyUDF.addOne';
Оператор CREATE FUNCTION использует следующие параметры:
LANGUAGE: язык UDF.HANDLER: полный путь к методу в формате'package.Object.method'(Scala) или'package.ClassName.method'(Java).DETERMINISTIC: объявляет, что функция всегда возвращает одинаковые выходные данные для одних и того же входных данных, что позволяет оптимизировать запросы.Note
Удалите,
DETERMINISTICесли функция вызывает внешние API или имеет другое недетерминированное поведение.ENVIRONMENT: определяет среду выполнения для UDF.-
java_dependencies: массив JSON, содержащий пути к JAR-файлам в томах Unity Catalog. Это путь к файлу, скопированный на предыдущем шаге. Используйте одинарные кавычки вокруг массива и двойные кавычки вокруг путей. -
environment_version: должно быть не ниже'4'для пользовательских функций Scala и Java. Среда выполнения версии 4 использует Scala 2.13.16 и JDK 17. См. версии бессерверных сред .
-
Вызов UDF в SQL и записных книжках
После регистрации можно вызвать UDF в sql-запросах, записных книжках и представлениях:
-- Simple select
SELECT my_catalog.my_schema.add_one(5) AS result;
-- With table data
SELECT
id,
price,
currency,
my_catalog.my_schema.convert_to_usd(price, currency) AS price_usd
FROM my_catalog.my_schema.transactions;
-- Filtering
SELECT *
FROM my_catalog.my_schema.products
WHERE my_catalog.my_schema.convert_to_usd(price, currency) > 100;
-- Aggregation
SELECT
category,
SUM(my_catalog.my_schema.convert_to_usd(price, currency)) AS total_usd
FROM my_catalog.my_schema.sales
GROUP BY category;
Управление и совместное использование
Используйте разрешения каталога Unity, чтобы контролировать, кто может запускать UDF и сделать его обнаруживаемым в вашей организации.
Предоставление разрешений
Используйте Catalog Explorer или SQL, чтобы предоставить другим пользователям необходимые разрешения на выполнение ваших пользовательских функций.
Обозреватель каталогов
- На боковой панели щелкните
Каталог.
- Выберите каталог, а затем выберите схему, содержащую функцию.
- Щелкните имя функции.
- На вкладке "Разрешения" нажмите кнопку "Предоставить".
- Выберите субъекты, к которым вы хотите предоставить доступ, и выберите
EXECUTEразрешение. - Нажмите кнопку "Подтвердить".
SQL
Выполните следующую команду в записной книжке или редакторе Databricks SQL, чтобы предоставить EXECUTE разрешения пользователю или группе.
-- Grant to a specific user
GRANT EXECUTE ON FUNCTION my_catalog.my_schema.add_one TO `user@example.com`;
-- Grant to a group
GRANT EXECUTE ON FUNCTION my_catalog.my_schema.add_one TO `data-engineers`;
Отмена разрешений
Используйте обозреватель каталогов или SQL для отзыва разрешений от других пользователей.
Обозреватель каталогов
- На боковой панели щелкните
Каталог.
- Выберите каталог, а затем выберите схему, содержащую функцию.
- Щелкните имя функции.
- На вкладке "Разрешения" установите флажок рядом с субъектом, из которого требуется отозвать доступ. Нажмите Отменить.
- В уведомлении нажмите кнопку "Отозвать".
SQL
Выполните следующую команду в записной книжке или редакторе SQL Databricks, чтобы отозвать EXECUTE разрешения пользователя или группы.
-- Revoke from specific user
REVOKE EXECUTE ON FUNCTION my_catalog.my_schema.add_one FROM `user@example.com`;
-- Revoke from a group
REVOKE EXECUTE ON FUNCTION my_catalog.my_schema.add_one FROM `data-engineers`;
Найти пользовательские функции
Чтобы найти пользовательские функции, которыми управляет Unity Catalog, выполните запрос к таблице information_schema.routines, заменив значения my_catalog и my_schema:
SELECT
routine_catalog,
routine_schema,
routine_name,
routine_definition,
created
FROM system.information_schema.routines
WHERE routine_catalog = 'my_catalog'
AND routine_schema = 'my_schema';
Обновите UDF
Чтобы обновить существующую UDF в Unity Catalog новым кодом:
- Внесите изменения в код локально.
- Перестройте JAR-файл с новым номером версии.
- Scala:
sbt clean assembly(например,my-udf-assembly-0.2.0-SNAPSHOT.jar) - Java:
mvn clean package(например,my-udf-2.0-SNAPSHOT.jar)
- Scala:
- Загрузите новый JAR-файл в том Unity Catalog.
- Используйте
CREATE OR REPLACE FUNCTIONс тем же именем функции, чтобы обновить UDF. Убедитесь, что вы ссылаетесь на последнюю версию JAR-файла в вашемjava_dependencies.
Azure Databricks использует новый код для следующего вызова. Вам не нужно перезапустить кластер.
Оптимизация производительности
Задержка холодного запуска
Первый вызов UDF в сеансе инициализирует изолированную песочницу, что приводит к задержке. Последующие вызовы в том же сеансе выполняются быстрее. Учитывайте это при тестировании или проектировании рабочих нагрузок с учетом задержки.
Кэширование дорогостоящих вычислений
Если UDF выполняет дорогостоящие инициализации или вычисления, кэшируйте результат, чтобы вычислить его только один раз.
Scala
val Используйте поле в объекте Scala для кэширования результата:
package example
object CachedUDF {
// Computed once and cached
val expensiveData: Map[String, Double] = {
// Load data from somewhere expensive
Map("key1" -> 1.0, "key2" -> 2.0)
}
def lookup(key: String): Double = {
expensiveData.getOrElse(key, 0.0)
}
}
Java
Используйте поле static и блок статической инициализации, чтобы кэшировать результат:
package example;
import java.util.Map;
import java.util.HashMap;
public class CachedUDF {
// Computed once and cached
private static Map<String, Double> expensiveData;
static {
// Load data from somewhere expensive
expensiveData = new HashMap<>();
expensiveData.put("key1", 1.0);
expensiveData.put("key2", 2.0);
}
public static double lookup(String key) {
return expensiveData.getOrDefault(key, 0.0);
}
}
При необходимости используйте DETERMINISTIC
Пометьте UDF как DETERMINISTIC, если она всегда возвращает один и тот же результат для одних и тех же входных данных. Это позволяет оптимизатору запросов кэшировать результаты и повысить производительность.
Ограничения
- Поддерживаются только скалярные определяемые пользователем функции. Определяемые пользователем агрегатные функции (UDAFs) и определяемые пользователем функции таблиц (UDTFs) не поддерживаются.
- Пользовательские функции выполняются в изолированной песочнице без активной сессии Spark. API Spark (
SparkSession,SparkContext,spark.sql(...), операции DataFrame и Dataset) недоступны. - UDF не могут зависеть от пакетов ядра Spark.
- Пользовательские функции не имеют доступа к файлам рабочей области или томам Unity Catalog во время выполнения.
Лучшие практики
Databricks рекомендует следующие методики.
- Управляйте версиями JAR-файлов. Например,
my-udf-0.1.0.jar,my-udf-0.2.0.jar. - Проверка сопоставлений типов SQL перед развертыванием. См. сопоставления языков.
- Предоставляйте разрешения
READ VOLUMEиEXECUTEтолько тем пользователям, которым необходимо запускать UDF. Используйте групповое владение для пользовательских функций, общих для разных команд.
Локальное тестирование пользовательских функций
Протестируйте UDF с помощью модульных тестов перед развертыванием в рабочей среде.
Scala
Чтобы проверить src/main/scala/example/MyUDF.scala, создайте тестовый файл в src/test/scala/example/MyUDFTest.scala:
package example
import org.scalatest.funsuite.AnyFunSuite
class MyUDFTest extends AnyFunSuite {
test("addOne should add 1 to input") {
assert(MyUDF.addOne(5) == 6)
}
test("addOne should handle negative numbers") {
assert(MyUDF.addOne(-1) == 0)
}
}
Добавьте зависимость для тестов в build.sbt:
libraryDependencies += "org.scalatest" %% "scalatest" % "3.2.15" % Test
Чтобы выполнить тесты, выполните следующие действия:
sbt test
Java
Чтобы проверить src/main/java/com/example/MyUDF.java, создайте тестовый файл в src/test/java/com/example/MyUDFTest.java:
package com.example;
import org.junit.jupiter.api.Test;
import static org.junit.jupiter.api.Assertions.*;
public class MyUDFTest {
@Test
public void testAddOne() {
assertEquals(6, MyUDF.addOne(5));
}
@Test
public void testAddOneWithNegativeNumbers() {
assertEquals(0, MyUDF.addOne(-1));
}
}
Добавьте зависимость JUnit в раздел <dependencies> вашего pom.xml:
<dependency>
<groupId>org.junit.jupiter</groupId>
<artifactId>junit-jupiter</artifactId>
<version>5.10.0</version>
<scope>test</scope>
</dependency>
Чтобы выполнить тесты, выполните следующие действия:
mvn test