세션 범위의 Scala 및 Java UDF

Important

거버넌스, 재사용 및 검색 기능을 위해 Unity 카탈로그에 Scala 및 Java UDF를 등록할 수 있습니다. Unity 카탈로그의 Scala 및 Java 사용자 정의 함수(UDF)를 참조하세요.

이 페이지에서는 Azure Databricks에서 세션 범위의 Scala 및 Java UDF를 만드는 방법을 설명합니다. 세션 범위 UDF는 Notebook 또는 작업에서 정의되며 현재 SparkSession에만 적용됩니다. SQL 언어 참조는 외부 UDF(사용자 정의 스칼라 함수)를 참조하세요.

접근 방식 선택

다음과 같은 방법으로 Scala 또는 Java UDF를 정의할 수 있습니다. 언어, 거버넌스 및 컴퓨팅에서 모든 UDF 유형을 비교하려면 Unity 카탈로그 관리 및 세션 범위 UDF를 참조하세요.

Approach Description
인라인 Scala UDF Scala 함수 또는 람다를 사용하여 Notebook에서 UDF를 정의합니다. 세션 범위로 한정됨. 서버리스 컴퓨팅에서는 지원되지 않습니다.
JAR의 Java UDF 를 사용하여 JAR에서 미리 컴파일된 UDF 클래스를 등록합니다 spark.udf.registerJavaFunction. 세션 범위로 한정됨. 서버리스 컴퓨팅에서 지원됩니다.
Unity 카탈로그로 관리되는 Scala 또는 Java UDF 거버넌스, 재사용 및 검색 가능성을 위해 Unity 카탈로그에 UDF를 등록합니다. 서버리스 컴퓨팅에서 지원됩니다.

요구 사항

  • 표준 액세스 모드를 사용하는 Unity 카탈로그 지원 컴퓨팅의 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 등록을 참조하세요.

Important

JAR을 실행하는 컴퓨팅 환경과 동일한 Scala 및 Apache Spark 버전에 맞춰 JAR을 빌드하세요. 일치하지 않는 경우 등록 또는 호출 시간에 UDF가 실패할 수 있습니다.

  • 클래식 컴퓨팅: Databricks 런타임 버전의 Scala 및 Spark 버전과 일치합니다. 사용 중인 버전의 Databricks Runtime 릴리스 노트 버전 및 호환성에서 시스템 환경 섹션을 참조하세요. 예를 들어 Databricks Runtime 18 LTS는 Scala 2.13.16 및 Apache Spark 4.0을 사용합니다.
  • 서버리스 컴퓨팅: 환경 버전의 Scala 버전과 일치시켜야 합니다. 환경 버전을 참조하세요.

Apache Spark 종속성을 provided로 표시하여 JAR 파일에 번들로 포함되지 않도록 하세요. UDF에서 사용하는 타사 종속성만 포함합니다.

UDF로 함수 등록

다음을 사용하여 spark.udf.registerScala 함수를 UDF로 등록합니다.

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

DataFrames에서 UDF 사용

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"))

UDF가 포함된 파일

Important

이 기능은 베타 버전으로 제공됩니다. 작업 영역 관리자는 미리 보기 페이지에서 이 기능에 대한 액세스를 제어할 수 있습니다. Azure Databricks 미리 보기 관리를 참조하세요.

파일의 스칼라 타입은 FileRef입니다. UDF에서 매개변수나 반환 타입으로 사용하며, 최상위 타입이나 중첩 형태로 사용할 수 있습니다. 타입과 중첩 규칙에 대해서는 타입을 참조하세요FILE.

UDF에서 파일 내용을 읽으려면 FileRef에서 다음 중 하나를 호출합니다:

  • asLocalFile(): 경로를 인자로 받는 모든 라이브러리에 전달할 수 있는 java.io.File를 반환합니다.
  • open(): 파일의 바이트를 읽으면 a java.io.InputStream 를 반환합니다. 발신자가 종료합니다.

UDF에서 새 FileRef 메서드를 생성하려면 다음 정적 메서드 중 하나를 호출하세요:

  • FileRef.create(uri): uri에 있는 파일에 대한 참조를 생성합니다.
  • FileRef.fromBytes(bytes, destinationPath, contentType): destinationPathbytes 볼륨 경로에 외부 파일로 업로드하고 참조를 반환합니다.
  • FileRef.fromLocalFile(localFile, destinationPath, contentType): 로컬 파일을 Volume 경로에 destinationPath 외부 파일 형식으로 업로드하고 참조를 반환합니다.

FileRef 열에 쓰는 UDF에서 FILE MANAGED을 반환하는 것은 지원되지 않습니다.

Python, Scala, SQL의 이미지 처리, 파일 유형 감지, 비디오 프레임 추출 등의 예시는 UDF가 포함된 프로세스 파일을 참조하세요.

JAR에서 Java UDF 등록

UDF를 JAR로 패키지하고, 세션에 spark.addArtifact추가하고, UDF 클래스를 등록합니다 spark.udf.registerJavaFunction.

참고

Databricks Runtime 18 LTS 이상에서 표준 액세스 모드 및 서버리스 컴퓨팅에서 지원됩니다. 등록된 함수는 세션 범위이며 Unity 카탈로그에 등록되지 않습니다.

다음 단계에서는 프로젝트를 만들고, UDF 클래스를 작성하고, fat JAR을 빌드하고, 등록하는 단계를 안내합니다.

1단계: 프로젝트 만들기

Scala 또는 Java 프로젝트를 설정합니다.

Scala

다음을 사용하여 sbt새 Scala 프로젝트를 만듭니다.

sbt new scala/scala-seed.g8

파일의 build.sbt 내용을 다음으로 바꿉니다. scalaVersionspark-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 플러그 인을 사용하도록 설정하여 fat 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

이 명령은 src/main/javasrc/test/java 디렉터리가 있는 표준 Maven 프로젝트 구조를 만듭니다.

생성된 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 클래스는 UDF가 사용하는 입력 인수 수를 org.apache.spark.sql.api.java.UDF 나타내는 인터페이스 중 하나(UDF1 통과 UDF22)를 구현해야 합니다. 사용자의 로직으로 call() 메서드를 구현합니다.

처리기는 Java 클래스여야 합니다. spark.udf.registerJavaFunction 는 리플렉션으로 클래스를 로드하므로 공용 no-arg 생성자가 있는 최상위(또는 static 중첩된) 공용 클래스여야 합니다. 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단계: Unity 카탈로그 볼륨에 JAR 업로드

컴퓨팅에서 액세스할 수 있도록 JAR을 Unity 카탈로그 볼륨 에 업로드합니다. 아직 볼륨이 없다면 볼륨을 생성하세요:

CREATE VOLUME IF NOT EXISTS my_catalog.my_schema.udf_jars
COMMENT 'Storage for UDF JAR files';

카탈로그 탐색기를 사용하여 볼륨에 JAR 파일을 업로드합니다.

  1. Azure Databricks 작업 영역에서 데이터 아이콘을 클릭한 후 카탈로그를 선택하여 카탈로그 탐색기를 엽니다.
  2. 카탈로그를 선택한 다음 볼륨이 포함된 스키마를 선택합니다.
  3. 볼륨 이름을 클릭합니다.
  4. 이 볼륨에 업로드를 클릭하고 JAR 파일을 선택합니다.
  5. 업로드를 클릭합니다.
  6. 업로드가 완료되면 JAR 파일의 이름을 클릭한 다음 복사 경로를 클릭하여 볼륨 경로를 복사합니다. /Volumes/my_catalog/my_schema/udf_jars/my-udf-assembly-0.1.0-SNAPSHOT.jar을 예로 들 수 있습니다. 다음 단계에서 이 경로가 필요합니다.

5단계: UDF 등록 및 호출

볼륨 경로를 사용하여 세션에 JAR을 추가하고, UDF 클래스를 등록하고, Spark SQL에서 호출합니다.

# 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에서 실패합니다. 사용자 정의 집계 함수(UDAF)는 registerJavaFunction에서 지원되지 않습니다.

쿼리는 UDF 출력을 반환하여 함수가 등록되고 호출 가능한지 확인합니다.

+----------+
| my_udf(21)|
+----------+
|        22|
+----------+

평가 순서 및 Null 검사

Spark SQL(SQL 및 DataFrame 및 데이터 세트 API 포함)은 하위 식 평가 순서를 보장하지 않습니다. Spark는 연산자 또는 함수의 입력을 왼쪽에서 오른쪽으로 평가하지 않습니다. 논리 ANDOR 식에는 왼쪽에서 오른쪽 단락 의미 체계가 없습니다.

부울 식의 부작용, 평가 순서 또는 WHEREHAVING 절의 순서에 의존하지 마세요. 쿼리 최적화 프로그램은 이러한 식 및 절의 순서를 변경할 수 있습니다. UDF가 null 검사에 단락 의미 체계를 사용하는 경우 Spark는 UDF 전에 null 검사가 실행되도록 보장하지 않습니다. 다음은 그 예입니다.

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가 null을 걸러낸 후 strlen UDF를 호출한다고 보장하지 않습니다.

Null 검사를 처리하기 위해 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를 사용하여 사용자 정의 함수를 사용하여 데이터 세트에서 맵, 필터 및 집계와 같은 변환을 실행합니다.

다음 예제에서는 API를 map() 사용하여 결과 열의 숫자를 접두사 문자열로 수정합니다.

spark.range(3).map(f => s"row-$f").show()

이 예제에서는 동일한 패턴을 사용map()하지만 , filter(),mapPartitions()foreach()foreachPartition(), 및 reduce()같은 형식의 다른 데이터 세트 APIflatMap()에도 동일한 패턴이 적용됩니다.

Scala UDF 기능 및 Databricks 런타임 호환성

다음 기능을 사용하려면 표준(공유) 액세스 모드에서 Unity 카탈로그 사용 클러스터의 최소 Databricks 런타임 버전이 필요합니다.

특징 최소 Databricks 런타임 버전
스칼라 사용자 정의 함수 데이터브릭스 런타임 14.2
Dataset.map, Dataset.mapPartitions, Dataset.filter, Dataset.reduceDataset.flatMap Databricks 런타임 15.4
KeyValueGroupedDataset.flatMapGroups, KeyValueGroupedDataset.mapGroups Databricks 런타임 15.4
(스트리밍) foreachWriter Sink Databricks 런타임 15.4
(스트리밍) foreachBatch Databricks Runtime 16.1 (데이터브릭스 런타임 16.1)
(스트리밍) KeyValueGroupedDataset.flatMapGroupsWithState Databricks Runtime 16.2 (데이터브릭스 런타임 16.2)
spark.udf.registerJavaFunction (JAR의 Java UDF) Databricks 런타임 18 LTS