Scala dengan cakupan sesi dan UDF Java

Important

Scala dan Java UDF dapat didaftarkan di Unity Catalog untuk tata kelola, penggunaan kembali, dan penemuan. Lihat fungsi yang ditentukan pengguna (UDF) Scala dan Java di Unity Catalog.

Halaman ini menjelaskan cara membuat UDF Scala dan Java yang bercakupan sesi di Azure Databricks. UDF dengan cakupan sesi ditentukan dalam buku catatan atau pekerjaan dan hanya berlaku untuk SparkSession saat ini. Untuk referensi bahasa SQL, lihat Fungsi skalar yang ditentukan pengguna eksternal (UDF).

Pilih pendekatan Anda

Anda dapat menentukan Scala atau Java UDF dengan cara berikut. Untuk membandingkan semua jenis UDF di berbagai bahasa, tata kelola, dan sumber daya komputasi, lihat UDF yang dikelola Unity Catalog vs. UDF berlingkup sesi.

Approach Deskripsi
Inline Scala UDF Tentukan UDF di notebook menggunakan fungsi Scala atau lambda. Berlaku dalam lingkup sesi. Tidak didukung pada komputasi tanpa server.
Java UDF dari JAR Daftarkan kelas UDF yang telah dikompilasi sebelumnya dari JAR menggunakan spark.udf.registerJavaFunction. Berlaku dalam lingkup sesi. Didukung pada komputasi tanpa server.
UDF Scala atau Java yang diatur oleh Unity Catalog Daftarkan UDF di Unity Catalog untuk tata kelola, penggunaan kembali, dan penemuan. Didukung pada komputasi tanpa server.

Persyaratan

  • UDF Scala pada komputasi yang mendukung Unity Catalog dengan mode akses standar memerlukan Databricks Runtime 14.2 atau lebih tinggi.
  • Dukungan instans ARM untuk Scala UDF pada kluster yang mendukung Unity Catalog memerlukan Databricks Runtime 15.2 atau lebih tinggi.
  • Mendaftarkan UDF Java dari JAR dengan spark.udf.registerJavaFunction memerlukan Databricks Runtime 18 LTS atau lebih tinggi. Lihat Mendaftarkan Java UDF dari JAR.

Important

Bangun file JAR Anda menggunakan versi Scala dan Apache Spark yang sama dengan yang digunakan oleh lingkungan komputasi yang menjalankannya. Ketidakcocokan dapat menyebabkan UDF gagal pada waktu pendaftaran atau panggilan.

  • Komputasi klasik: Cocokkan versi Scala dan Spark dari versi Databricks Runtime Anda. Lihat bagian Lingkungan sistem dari versi dan kompatibilitas catatan rilis Databricks Runtime untuk versi Anda. Misalnya, Databricks Runtime 18 LTS menggunakan Scala 2.13.16 dan Apache Spark 4.0.
  • Komputasi nirserver: Sesuaikan versi Scala dengan Versi Lingkungan Anda. Lihat versi Lingkungan.

Tandai dependensi Apache Spark sebagai provided agar tidak disertakan ke dalam JAR Anda. Hanya sertakan dependensi pihak ketiga yang digunakan UDF Anda.

Mendaftarkan fungsi sebagai UDF

Daftarkan fungsi Scala sebagai UDF menggunakan spark.udf.register:

val squared = (s: Long) => {
  s * s
}
spark.udf.register("square", squared)

Memanggil UDF di Spark SQL

Buat tampilan sementara, lalu panggil UDF dalam kueri SQL:

spark.range(1, 20).createOrReplaceTempView("test")
%sql select id, square(id) as id_squared from test

Menggunakan UDF dengan DataFrames

Anda juga dapat memanggil UDF menggunakan API DataFrame:

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

File-file dengan UDF

Important

Fitur ini ada di Beta. Admin ruang kerja dapat mengontrol akses ke fitur ini dari halaman Pratinjau . Lihat Kelola Pratinjau Azure Databricks.

Tipe Scala untuk sebuah file adalah FileRef. Gunakan sebagai parameter atau tipe return dalam UDF, baik sebagai tipe tingkat atas atau bersarang. Untuk tipe dan aturan bersarangnya, lihat FILE tipe.

Untuk membaca isi file di UDF, panggil salah satu dari berikut pada :FileRef

  • asLocalFile(): Mengembalikan java.io.File yang dapat diteruskan ke pustaka apa pun yang menerima jalur.
  • open(): Mengembalikan java.io.InputStream untuk membaca byte dari file. Penelepon menutupnya.

Untuk menghasilkan yang baru FileRef dalam UDF, panggil salah satu metode statis berikut:

  • FileRef.create(uri): Membuat referensi ke file di .uri
  • FileRef.fromBytes(bytes, destinationPath, contentType): Mengunggah bytes ke destinationPath jalur Volume sebagai file eksternal dan mengembalikan referensi.
  • FileRef.fromLocalFile(localFile, destinationPath, contentType): Mengunggah file lokal ke destinationPath jalur Volume sebagai file eksternal dan mengembalikan referensi.

Mengembalikan FileRef dari UDF yang menulis ke kolom FILE MANAGED tidak didukung.

Untuk contoh di Python, Scala, dan SQL, termasuk pemrosesan gambar, deteksi tipe file, dan ekstraksi bingkai video, lihat Proses file dengan UDF.

Daftarkan UDF Java dari JAR

Paketkan UDF sebagai JAR, tambahkan ke sesi Anda dengan spark.addArtifact, dan daftarkan kelas UDF dengan spark.udf.registerJavaFunction.

Catatan

Didukung pada mode akses standar dan komputasi tanpa server di Databricks Runtime 18 LTS atau yang lebih tinggi. Fungsi yang terdaftar berada dalam lingkup sesi dan tidak terdaftar di Unity Catalog.

Langkah-langkah berikut menjelaskan cara membuat proyek, menulis kelas UDF, membuat fat JAR, dan mendaftarkannya.

Langkah 1: Buat proyek Anda

Siapkan proyek di Scala atau Java.

Scala

Buat proyek Scala baru menggunakan sbt:

sbt new scala/scala-seed.g8

Ganti konten file Anda build.sbt dengan yang berikut ini. Atur scalaVersion dan versi spark-sql agar sesuai dengan komputasi Anda:

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

Aktifkan plugin sbt-assembly untuk membuat fat JAR. Buat atau edit project/assembly.sbt dan tambahkan:

addSbtPlugin("com.eed3si9n" % "sbt-assembly" % "2.0.0")

Java

Buat proyek Maven baru menggunakan arketipe mulai cepat:

mvn archetype:generate \
  -DgroupId=com.example \
  -DartifactId=my-udf \
  -DarchetypeArtifactId=maven-archetype-quickstart \
  -DinteractiveMode=false

Perintah ini membuat struktur proyek Maven standar dengan src/main/java direktori dan src/test/java .

Dalam pom.xml yang dihasilkan, di dalam tag <project></project>, tambahkan blok <properties> dan konfigurasikan maven-shade-plugin untuk membuat JAR fat:

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

Langkah 2: Tulis kelas UDF Anda

Kelas UDF Anda harus mengimplementasikan salah satu antarmuka org.apache.spark.sql.api.java.UDF(UDF1 hingga UDF22), dengan angka yang menunjukkan berapa banyak argumen input yang diterima UDF. Terapkan metode call() sesuai logika Anda.

Handler harus merupakan kelas Java. spark.udf.registerJavaFunction memuat kelas melalui refleksi, jadi kelas tersebut harus merupakan kelas publik tingkat atas (atau kelas publik bertingkat static) dengan konstruktor publik tanpa argumen. Scala class atau object tidak memenuhi persyaratan ini dan gagal pada waktu panggilan. Anda dapat membangun JAR dengan sbt, tetapi kelas UDF itu sendiri harus ditulis dalam Java.

Buat 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;
  }
}

Langkah 3: Buat fat JAR Anda

Kemas UDF Anda yang telah dikompilasi menjadi fat JAR.

Scala

Dari direktori akar proyek Anda, jalankan:

sbt clean assembly

fat JAR dibuat di target/scala-2.13/ dengan nama seperti my-udf-assembly-0.1.0-SNAPSHOT.jar.

Java

Dari direktori akar proyek Anda, jalankan:

mvn clean package

fat JAR dibuat di target/ dengan nama seperti my-udf-1.0-SNAPSHOT.jar.

Langkah 4: Unggah JAR Anda ke volume Katalog Unity

Unggah JAR ke volume Katalog Unity sehingga komputasi Anda dapat mengaksesnya. Jika Anda belum memiliki volume, buatlah sebuah volume:

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

Unggah file JAR Anda ke volume menggunakan Catalog Explorer:

  1. Di ruang kerja Azure Databricks Anda, klik Ikon data.Katalog untuk membuka Catalog Explorer.
  2. Pilih katalog, lalu pilih skema yang berisi volume Anda.
  3. Klik nama volume.
  4. Klik Unggah ke volume ini dan pilih file JAR Anda.
  5. Klik Unggah.
  6. Setelah pengunggahan selesai, klik nama file JAR Anda, lalu klik Salin jalur untuk menyalin jalur volume. Contohnya, /Volumes/my_catalog/my_schema/udf_jars/my-udf-assembly-0.1.0-SNAPSHOT.jar. Anda memerlukan jalur ini di langkah berikutnya.

Langkah 5: Daftarkan dan panggil UDF

Tambahkan JAR ke sesi Anda menggunakan jalur volumenya, daftarkan kelas UDF, dan panggil dari 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()

Pada komputasi dengan mode akses serverless dan standar, Anda harus menentukan tipe nilai kembalian yang eksplisit. Menghilangkan jenis pengembalian gagal dengan UC_COMMAND_NOT_SUPPORTED_IN_SHARED_ACCESS_MODE. Fungsi agregat yang ditentukan pengguna (UDAF) tidak didukung dengan registerJavaFunction.

Kueri mengembalikan output UDF, mengonfirmasi bahwa fungsi terdaftar dan dapat dipanggil:

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

Urutan evaluasi dan pengecekan null

Spark SQL (termasuk SQL dan API DataFrame dan Himpunan Data) tidak menjamin urutan evaluasi subekspresi. Spark tidak mengevaluasi input operator atau fungsi kiri-ke-kanan. Logika AND dan OR ekspresi tidak memiliki semantik sirkuit pendek kiri-ke-kanan.

Jangan bergantung pada efek samping atau urutan evaluasi ekspresi Boolean, maupun urutan klausa WHERE dan HAVING. Pengoptimal kueri dapat menyusun ulang ekspresi dan klausa ini. Jika UDF bergantung pada semantik arus pendek untuk pemeriksaan null, Spark tidak menjamin bahwa pemeriksaan null berjalan sebelum UDF. Contohnya:

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

Klausa WHERE ini tidak menjamin bahwa Spark memanggil strlen UDF setelah menyaring nilai null.

Untuk menangani pemeriksaan null, Databricks merekomendasikan salah satu hal berikut:

  • Buat UDF itu sendiri agar dapat menangani nilai null, dan lakukan pengecekan null di dalam UDF
  • Gunakan ekspresi IF atau CASE WHEN untuk melakukan pemeriksaan null dan memanggil UDF di cabang kondisional
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 Himpunan Data yang Ditik

Catatan

Fitur ini didukung pada kluster yang mendukung Unity Catalog dengan mode akses standar di Databricks Runtime 15.4 ke atas.

Gunakan API Himpunan Data yang ditik untuk menjalankan transformasi seperti peta, filter, dan agregasi pada Himpunan Data dengan fungsi yang ditentukan pengguna.

Contoh berikut menggunakan map() API untuk mengubah angka di kolom hasil menjadi string awalan:

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

Contoh ini menggunakan map(), tetapi pola yang sama berlaku untuk API Himpunan Data jenis lainnya seperti filter(), , mapPartitions(), foreach()foreachPartition(), reduce(), dan flatMap().

Fitur-fitur UDF Scala dan kompatibilitas dengan Databricks Runtime

Fitur berikut memerlukan versi minimum Databricks Runtime pada kluster dengan Unity Catalog yang diaktifkan dalam mode akses standar (bersama).

Fitur Versi Minimum Runtime Databricks
UDF skalar Databricks Runtime 14.2
Dataset.map, Dataset.mapPartitionsDataset.filter, Dataset.reduce,Dataset.flatMap Databricks Runtime 15.4
KeyValueGroupedDataset.flatMapGroups, KeyValueGroupedDataset.mapGroups Databricks Runtime 15.4
(Siaran streaming) foreachWriter Sink Databricks Runtime 15.4
(Siaran streaming) foreachBatch Databricks Runtime 16.1
(Siaran streaming) KeyValueGroupedDataset.flatMapGroupsWithState Databricks Runtime 16.2
spark.udf.registerJavaFunction(Java UDF dari JAR) Databricks Runtime 18 LTS