Пользовательские функции Scala и Java в Unity Catalog

На этой странице описано, как создавать пользовательские функции 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 соответствует методу в Scala object.
  • 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-файл в том с помощью обозревателя каталогов:

  1. В рабочей области Azure Databricks щелкните Data icon.Catalog, чтобы открыть обозреватель каталогов.
  2. Выберите каталог, а затем выберите схему, содержащую том.
  3. Щелкните имя тома.
  4. Нажмите кнопку "Отправить в этот том " и выберите JAR-файл.
  5. Нажмите кнопку Отправить.
  6. После завершения загрузки нажмите на имя JAR-файла.
  7. Нажмите Копировать путь, чтобы скопировать путь тома в буфер обмена. Например, /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, чтобы предоставить другим пользователям необходимые разрешения на выполнение ваших пользовательских функций.

Обозреватель каталогов

  1. На боковой панели щелкните значок Каталог.
  2. Выберите каталог, а затем выберите схему, содержащую функцию.
  3. Щелкните имя функции.
  4. На вкладке "Разрешения" нажмите кнопку "Предоставить".
  5. Выберите субъекты, к которым вы хотите предоставить доступ, и выберите EXECUTE разрешение.
  6. Нажмите кнопку "Подтвердить".

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 для отзыва разрешений от других пользователей.

Обозреватель каталогов

  1. На боковой панели щелкните значок Каталог.
  2. Выберите каталог, а затем выберите схему, содержащую функцию.
  3. Щелкните имя функции.
  4. На вкладке "Разрешения" установите флажок рядом с субъектом, из которого требуется отозвать доступ. Нажмите Отменить.
  5. В уведомлении нажмите кнопку "Отозвать".

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 новым кодом:

  1. Внесите изменения в код локально.
  2. Перестройте JAR-файл с новым номером версии.
    • Scala: sbt clean assembly (например, my-udf-assembly-0.2.0-SNAPSHOT.jar)
    • Java: mvn clean package (например, my-udf-2.0-SNAPSHOT.jar)
  3. Загрузите новый JAR-файл в том Unity Catalog.
  4. Используйте 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

Дополнительные ресурсы