Oturum kapsamındaki Scala ve Java UDF'leri

Important

Scala ve Java UDF'leri, idare, yeniden kullanım ve bulunabilirlik için Unity Kataloğu'nda kaydedilebilir. Bkz. Unity Kataloğu'nda Scala ve Java kullanıcı tanımlı işlevler (UDF' ler).

Bu sayfada, Azure Databricks oturum kapsamlı Scala ve Java UDF'lerinin nasıl oluşturulacağı açıklanır. Oturum kapsamındaki UDF'ler bir not defterinde veya görevde tanımlanır ve yalnızca geçerli SparkSession için geçerlidir. SQL dil başvurusu için bkz. Dış kullanıcı tanımlı skaler işlevler (UDF).

Yaklaşımınızı seçin

Scala veya Java UDF'yi aşağıdaki yollarla tanımlayabilirsiniz. Diller, yönetişim ve işlem kaynakları genelindeki tüm UDF türlerini karşılaştırmak için bkz: Unity Catalog tarafından yönetilen UDF'ler ile oturum kapsamlı UDF'ler.

Approach Açıklama
Satır içi Scala UDF Scala işlevini veya lambdayı kullanarak not defterinde UDF tanımlayın. Oturum kapsamında. Sunucusuz işlemde desteklenmez.
JAR'dan Java UDF kullanarak spark.udf.registerJavaFunctionJAR'dan önceden derlenmiş bir UDF sınıfını kaydedin. Oturum kapsamında. Sunucusuz işlemde desteklenir.
Unity Kataloğu ile yönetilen Scala veya Java UDF İdare, yeniden kullanım ve bulunabilirlik için Unity Kataloğu'nda bir UDF kaydedin. Sunucusuz işlemde desteklenir.

Gereksinimler

  • Standart erişim moduyla Unity Kataloğu özellikli işlemde Scala UDF'leri Databricks Runtime 14.2 veya üzerini gerektirir.
  • Unity Catalog etkinleştirilmiş kümelerde Scala UDF'ler için ARM instance desteği, Databricks Runtime 15.2 veya üzerini gerektirir.
  • Bir JAR'dan spark.udf.registerJavaFunction ile Java UDF kaydetmek için Databricks Runtime 18 LTS veya üstü gerekir. Bkz. JAR'dan Java UDF kaydetme.

Important

JAR dosyanızı, onu çalıştıran işlem ortamıyla aynı Scala ve Apache Spark sürümlerini kullanarak oluşturun. Bir uyuşmazlık olması, UDF'nin kayıt sırasında veya çağrı anında başarısız olmasına neden olabilir.

  • Klasik işlem: Databricks Runtime sürümünüzün Scala ve Spark sürümleriyle eşleşin. Databricks Runtime sürüm notlarının sürümleri ve sürümünüz için uyumluluk konusunun Sistem ortamı bölümüne bakın. Örneğin, Databricks Runtime 18 LTS Scala 2.13.16 ve Apache Spark 4.0 kullanır.
  • Sunucusuz işlem: Ortam Sürümünüzün Scala sürümüyle eşleşin. Çevre sürümlerine bakınız.

Apache Spark bağımlılığını, JAR dosyanıza dahil edilmemesi için provided olarak işaretleyin. Yalnızca UDF'nizin kullandığı üçüncü taraf bağımlılıklarını dahil edin.

İşlevi UDF olarak kaydetme

kullanarak spark.udf.registerScala işlevini UDF olarak kaydetme:

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

Spark SQL'de UDF'yi çağırma

Geçici görünüm oluşturun, ardından SQL sorgusunda UDF'yi çağırın:

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

DataFrame'lerle UDF kullanma

DataFrame API'sini kullanarak da bir UDF çağırabilirsiniz:

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

Important

Bu özellik Beta sürümündedir. Çalışma alanı yöneticileri Bu özelliğe erişimi Önizlemeler sayfasından denetleyebilir. Bkz. Azure Databricks önizlemelerini yönetme.

Bir dosyanın Scala tipi FileRef. Bunu bir UDF'de parametre veya dönüş tipi olarak, üst seviye tip olarak veya iç içe olarak kullanın. Tür ve iç içe yerleştirme kuralları için bkz. FILE tür.

Bir dosyanın içeriğini bir UDF'de okumak için, aşağıdaki FileRefadreslerden birini çağırın:

  • asLocalFile(): Bir yolu kabul eden herhangi bir kütüphaneye iletebileceğiniz bir java.io.File gönderi döndürür.
  • open(): Dosyanın baytlarını okumak için a java.io.InputStream döndürür. Arayan kapatıyor.

UDF'de yeni FileRef bir yöntem üretmek için aşağıdaki statik yöntemlerden birini çağırın:

  • FileRef.create(uri): uri konumundaki dosya için bir başvuru oluşturur.
  • FileRef.fromBytes(bytes, destinationPath, contentType): destinationPath ögesini bytes Volume yoluna harici bir dosya olarak yükler ve bir referans döndürür.
  • FileRef.fromLocalFile(localFile, destinationPath, contentType): Yerel bir dosyayı destinationPath Volume yoluna harici bir dosya olarak yükler ve bir referans döndürür.

FileRef sütununa yazan bir UDF'den FILE MANAGED döndürmek desteklenmez.

Python, Scala ve SQL'deki örnekler için, örneğin görüntü işleme, dosya tipi algılama ve video çerçeve çıkarma gibi konularda, bkz. UDF'lerle dosyaları işlemek.

JAR'dan Java UDF kaydetme

Bir UDF'yi JAR olarak paketleyip ile oturumunuza spark.addArtifactekleyin ve ile UDF sınıfını spark.udf.registerJavaFunctionkaydedin.

Not

Databricks Runtime 18 LTS veya üzerinde standart erişim modunda ve sunucusuz işlemde desteklenir. Kaydedilen işlev oturum kapsamındadır ve Unity Kataloğu'na kaydedilmez.

Aşağıdaki adımlar proje oluşturma, UDF sınıfı yazma, fat JAR oluşturma ve kaydetme adımlarını gösterir.

1. Adım: Projenizi oluşturma

Scala'da veya Java bir proje ayarlayın.

Scala

kullanarak sbtyeni bir Scala projesi oluşturun:

sbt new scala/scala-seed.g8

Dosyanızın build.sbt içeriğini aşağıdakilerle değiştirin. scalaVersion ve spark-sql sürümünü işlem ortamınıza uyacak şekilde ayarlayın:

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

Tüm bağımlılıkları içeren bir JAR oluşturmak için sbt-assembly eklentisini etkinleştirin. Oluşturun veya düzenleyin project/assembly.sbt ve ekleyin:

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

Java

Hızlı başlangıç arketipini kullanarak yeni bir Maven projesi oluşturun:

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

Bu komut, src/main/java ve src/test/java dizinleriyle standart Maven proje yapısını oluşturur.

Oluşturulan pom.xml içinde, <project></project> etiketleri arasına bir <properties> bloğu ekleyin ve fat JAR oluşturmak için maven-shade-plugin öğesini yapılandırın:

<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. Adım: UDF sınıfınızı yazma

UDF sınıfınızın arabirimlerden birini org.apache.spark.sql.api.java.UDF (UDF1 ile UDF22) uygulaması gerekir; burada sayı, UDF'nin kaç giriş bağımsız değişkeni aldığını gösterir. call() yöntemini mantığınızla uygulayın.

İşleyici bir Java sınıfı olmalıdır. spark.udf.registerJavaFunction sınıfını yansımaya göre yükler, bu nedenle genel no-arg oluşturuculu bir üst düzey (veya static iç içe) bir ortak sınıf olmalıdır. Scala class veya object bu gereksinimi karşılamaz ve arama zamanında başarısız olur. JAR'yi sbt ile oluşturabilirsiniz, ancak UDF sınıfının kendisi Java yazılmalıdır.

Oluştur 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. Adım: fat JAR'ınızı oluşturun

Derlenmiş UDF'nizi şişman bir JAR içinde paketleyin.

Scala

Proje kök dizininizden şunu çalıştırın:

sbt clean assembly

Fat JAR dosyası, target/scala-2.13/ içinde my-udf-assembly-0.1.0-SNAPSHOT.jar gibi bir adla oluşturulur.

Java

Proje kök dizininizden şunu çalıştırın:

mvn clean package

Fat JAR dosyası, target/ içinde my-udf-1.0-SNAPSHOT.jar gibi bir adla oluşturulur.

4. Adım: JAR'ınızı Unity Kataloğu birimine yükleme

İşlem kaynağınızın erişebilmesi için JAR dosyasını bir Unity Catalog birimine yükleyin. Henüz bir disk biriminiz yoksa bir disk birimi oluşturun:

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

Katalog Gezgini'ni kullanarak JAR dosyanızı birime yükleyin:

  1. Azure Databricks çalışma alanınızda katalog gezginini açmak için Data icon.Catalog öğesine tıklayın.
  2. Kataloğu seçin ve ardından biriminizi içeren şemayı seçin.
  3. Birim adına tıklayın.
  4. Bu birime yükle'ye tıklayın ve JAR dosyanızı seçin.
  5. Yükle'ye tıklayın.
  6. Yükleme tamamlandıktan sonra JAR dosyanızın adına tıklayın, ardından birim yolunu kopyalamak için Yolu kopyala seçeneğine tıklayın. Örneğin, /Volumes/my_catalog/my_schema/udf_jars/my-udf-assembly-0.1.0-SNAPSHOT.jar. Sonraki adımda bu yola ihtiyacınız vardır.

5. Adım: UDF'yi kaydetme ve çağırma

Jar dosyasını birim yolunu kullanarak oturumunuza ekleyin, UDF sınıfını kaydedin ve Spark SQL'den çağırın:

# 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()

Sunucusuz ve standart erişim modundaki işlem kaynaklarında, açık bir dönüş türü belirtmeniz gerekir. Dönüş türü belirtilmezse UC_COMMAND_NOT_SUPPORTED_IN_SHARED_ACCESS_MODE hatasına neden olur. registerJavaFunction, kullanıcı tanımlı toplama işlevlerini (UDAF) desteklemez.

Sorgu UDF çıkışını döndürerek işlevin kayıtlı ve çağrılabilir olduğunu onaylar:

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

Değerlendirme sırası ve null denetimi

Spark SQL (SQL ve DataFrame ile Veri Kümesi API'leri dahil) alt ifade değerlendirmesinin sırasını garanti etmez. Spark, bir işlecin veya işlevin girişlerini soldan sağa değerlendirmez. Mantıksal AND ve OR ifadelerde soldan sağa kısa devre semantiği yoktur.

Boolean ifadelerin yan etkilerine, değerlendirilme sırasına veya WHERE ve HAVING tümcelerinin sırasına güvenmeyin. Sorgu iyileştiricisi bu ifadeleri ve yan tümceleri yeniden sıralayabilir. Bir UDF, null denetimi için kısa devre değerlendirme semantiğine dayanıyorsa Spark, null denetiminin UDF'den önce çalışacağını garanti etmez. Örneğin:

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

Bu WHERE ifade, Spark’ın null değerleri ayıkladıktan sonra strlen UDF’yi çağıracağını garanti etmez.

Databricks, null denetimini işlemek için aşağıdakilerden birini önerir:

  • UDF'nin kendisini null-aware yapın ve UDF içinde null denetimi yapın
  • Null denetimi yapmak ve koşullu dalda UDF'yi çağırmak için IF veya CASE WHEN ifadelerini kullanın.
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

Yazılan Veri Kümesi API'leri

Not

Bu özellik, Databricks Runtime 15.4 ve üzerinde standart erişim moduna sahip Unity Kataloğu etkin kümelerde desteklenir.

Kullanıcı tanımlı bir işlevle veri kümelerinde eşleme, filtre ve toplama gibi dönüştürmeleri çalıştırmak için yazılan Veri Kümesi API'lerini kullanın.

Aşağıdaki örnek, bir sonuç sütunundaki map() bir sayıyı ön ekli dize olarak değiştirmek için API'yi kullanır:

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

Bu örnekte map() kullanılır, ancak aynı kalıp filter(), mapPartitions(), foreach(), foreachPartition(), reduce() ve flatMap() gibi diğer türlendirilmiş Veri Kümesi API'leri için de geçerlidir.

Scala UDF özellikleri ve Databricks Runtime uyumluluğu

Aşağıdaki özellikler, standart (paylaşılan) erişim modunda Unity Kataloğu özellikli kümelerde en düşük Databricks Runtime sürümlerini gerektirir.

Özellik En Düşük Databricks Runtime sürümü
Skaler UDF'ler Databricks Runtime 14.2
Dataset.map, Dataset.mapPartitions, Dataset.filter, Dataset.reduce, , Dataset.flatMap Databricks Runtime 15.4
KeyValueGroupedDataset.flatMapGroups, KeyValueGroupedDataset.mapGroups Databricks Runtime 15.4
(Akış) foreachWriter Sink Databricks Runtime 15.4
(Akış) foreachBatch Databricks Çalışma Zamanı 16.1
(Akış) KeyValueGroupedDataset.flatMapGroupsWithState Databricks Runtime 16.2
spark.udf.registerJavaFunction (bir JAR dosyasından Java UDF) Databricks Runtime 18 LTS