Unity Kataloğu'nda Python kullanıcı tanımlı tablo işlevleri (UDF)

Unity Kataloğu kullanıcı tanımlı tablo işlevi (UDTF), skaler değerler yerine tam tabloları döndüren işlevleri kaydeder. Her çağrıdan tek bir sonuç değeri döndüren skaler işlevlerin aksine, UDF'ler sql deyiminin FROM yan tümcesinde çağrılır ve birden çok satır ve sütun döndürebilir.

UDF'ler özellikle şunlar için kullanışlıdır:

  • Dizileri veya karmaşık veri yapılarını birden çok satıra dönüştürme
  • Dış API'leri veya hizmetleri SQL iş akışlarıyla tümleştirme
  • Özel veri oluşturma veya zenginleştirme mantığı uygulama
  • Satırlar arasında durum bilgisi olan işlemleri gerektiren verileri işleme

Her bir UDTF çağrısı sıfır veya daha fazla parametre kabul eder. Bu bağımsız değişkenler, giriş tablolarının tamamını temsil eden skaler ifadeler veya tablo bağımsız değişkenleri olabilir.

UDTF'ler iki şekilde kaydedilebilir:

Gereksinimler

Unity Kataloğu Python UDF'leri aşağıdaki işlem türlerinde desteklenir:

  • Sunucusuz not defterleri ve işler
  • Standart erişim moduyla klasik işlem (Databricks Runtime 17.1 ve üzeri)
  • SQL ambarı (sunucusuz veya profesyonel)

Unity Kataloğu'nda UDTF oluşturma

Unity Kataloğu'nda yönetilen bir UDTF oluşturmak için SQL DDL kullanın. UDTF'ler bir SQL deyiminin FROM ifadesi kullanılarak çağrılır.

CREATE OR REPLACE FUNCTION square_numbers(start INT, end INT)
RETURNS TABLE (num INT, squared INT)
LANGUAGE PYTHON
HANDLER 'SquareNumbers'
DETERMINISTIC
AS $$
class SquareNumbers:
    """
    Basic UDTF that computes a sequence of integers
    and includes the square of each number in the range.
    """
    def eval(self, start: int, end: int):
        for num in range(start, end + 1):
            yield (num, num * num)
$$;

SELECT * FROM square_numbers(1, 5);

+-----+---------+
| num | squared |
+-----+---------+
| 1   | 1       |
| 2   | 4       |
| 3   | 9       |
| 4   | 16      |
| 5   | 25      |
+-----+---------+

Azure Databricks, Python UDF'lerini Python sınıfları olarak uygular ve çıkış satırları veren zorunlu eval bir yöntem kullanır.

Tablo bağımsız değişkenleri

Uyarı

Databricks Runtime 17.2 ve sonrasında TABLE bağımsız değişkenleri desteklenir.

UDTF'ler, tüm tabloları giriş bağımsız değişkeni olarak kabul ederek duruma dayalı karmaşık dönüşümleri ve toplama işlemlerini etkinleştirebilir.

eval() ve terminate() yaşam döngüsü yöntemleri

UDTF'lerdeki tabloya ait bağımsız değişkenler, her satırı işlemek için aşağıdaki işlevleri kullanır:

  • eval(): Giriş tablosundaki her satır için bir kez çağrılır. Bu, ana işleme yöntemidir ve gereklidir.
  • terminate(): eval() tarafından tüm satırlar işlendiğinde her bölümün sonunda bir kez çağrılır. Son toplu sonuçları vermek veya temizleme işlemleri gerçekleştirmek için bu yöntemi kullanın. Bu yöntem isteğe bağlıdır ancak toplamalar, sayım veya toplu işlem gibi durum bilgisi olan işlemler için gereklidir.

ve eval() yöntemleri hakkında terminate() daha fazla bilgi için bkz. Apache Spark belgeleri: Python UDTF.

Satır erişim desenleri

eval(), TABLE bağımsız değişkenlerinden satırları pyspark.sql.Row nesneleri olarak alır. Değerlere sütun adına (row['id'], row['name']) veya dizine (row[0], row[1]) göre erişebilirsiniz.

  • Şema esnekliği: TABLE Bağımsız değişkenleri şema tanımları olmadan bildirme (örneğin, , data TABLEt TABLE). işlevi herhangi bir tablo yapısını kabul eder, bu nedenle kodunuz gerekli sütunların mevcut olduğunu doğrulamalıdır.

Bkz Örnek: IP adreslerini CIDR ağ bloklarıyla eşleştirme ve Örnek: Azure Databricks görüntü uç noktalarını kullanarak toplu resim altyazısı oluşturma.

Dinamik çıkış şemasını hesaplama (çok biçimli UDF'ler)

Uyarı

Çok biçimli UC UDF'ler Databricks Runtime 18.1 ve üzerini gerektirir.

Çok biçimli UDTF, çıkış şemasını sorgu zamanında sütunları önceden bildirmek yerine statik analyze() bir yöntem kullanarak dinamik olarak belirler. Bir tane oluşturmak için sütun tanımları olmadan kullanın RETURNS TABLE ve işleyici sınıfında bir analyze() yöntem tanımlayın.

Aşağıdaki örnek, çağıran tarafından belirtilen alanları bir JSON dizesinden ayıklar ve bağımsız değişkene fields bağlı olarak farklı sütunlar döndürür:

CREATE OR REPLACE FUNCTION extract_fields(json_str STRING, fields STRING)
RETURNS TABLE
LANGUAGE PYTHON
HANDLER 'ExtractFields'
AS $$
class ExtractFields:
    @staticmethod
    def analyze(json_str, fields):

        # Build the output schema from the requested field names
        from pyspark.sql.types import StructType, StructField, StringType
        from pyspark.sql.udtf import AnalyzeResult
        col_names = [f.strip() for f in fields.value.split(",")]
        return AnalyzeResult(
            StructType([StructField(name, StringType()) for name in col_names])
        )

    def eval(self, json_str: str, fields: str):
        # Parse the JSON and yield only the requested fields
        import json
        data = json.loads(json_str)
        col_names = [f.strip() for f in fields.split(",")]
        yield tuple(data.get(name) for name in col_names)
$$;

-- Extract the name and city
SELECT * FROM extract_fields(
  '{"name": "Alice", "age": 30, "city": "Seattle"}',
  'name, city'
);
+-------+---------+
| name  | city    |
+-------+---------+
| Alice | Seattle |
+-------+---------+

analyze Yöntemini tanımlama

İşleyici sınıfı, UDTF ile aynı bağımsız değişkenleri kabul eden ve çıkış şemasını açıklayan @staticmethod döndüren analyze adlı bir AnalyzeResult yöntemi içermelidir. Azure Databricks, işlevi yürütmeden önce şemayı çözümlemek için sorgu planlama zamanında analyze() çağırır.

her parametresi analyze sınıfının bir örneğidir AnalyzeArgument :

Alan Description
dataType Giriş bağımsız değişkeninin türü DataTypeolarak belirtilmiştir. Giriş tablosu bağımsız değişkenleri için bu, tablonun sütunlarını temsil eden bir StructType tablodur.
value Giriş bağımsız değişkeninin Optional[Any]olarak değeri. Bu, tablo bağımsız değişkenleri veya sabit olmayan ifadeler içindir None.
isTable Giriş bağımsız değişkeninin bir tablo bağımsız değişkeni olup olmadığını BooleanType olarak kontrol edin.
isConstantExpression Giriş bağımsız değişkeninin BooleanTypeolarak sabit katlanabilir bir ifade olup olmadığı.

analyze yöntemi, AnalyzeResult sınıfının bir örneğini döndürür.

Alan Description
schema sonuç tablosunun StructTypeolarak şeması.
withSinglePartition ise True, tüm giriş satırlarını aynı UDTF sınıf örneğine gönderir.
partitionBy Eğer boş değilse, girdi satırlarını belirtilen ifadelere göre böler, böylece her benzersiz birleşim ayrı bir UDTF örneği tarafından işlenir.
orderBy Boş değilse, her bölüm içindeki satırların sırasını belirtir.
select Boş değilse, UDTF'nin giriş TABLE bağımsız değişkeninden hangi sütunları alacağını belirtir.

Uyarı

Unity Kataloğu'ndaki çok biçimli UDTF'ler için tüm içeri aktarmaları yöntem gövdesine analyze() yerleştirmeniz gerekir. Üst düzey importlar Unity Kataloğu güvenli alan ortamında kullanılamaz.

Durumu analyze'dan eval'ye ilet

analyze Yöntemi sorgu planlama zamanında bir kez çalıştırıldığından, sabit bağımsız değişkenleri önceden işlemek, yapılandırmaları ayrıştırmak veya arama derlemek için kullanabilirsiniz. Bu sonuçları eval adresine iletmek için, özel alanlara sahip bir @dataclass alt sınıfı olan AnalyzeResult oluşturun, ardından bunu analyze'den döndürün ve __init__ yönteminde kabul edin. Bu, her satır için pahalı çalışmaların yinelenmesinden kaçınıyor.

Aşağıdaki örnek, bir dil kodunu bir kez analyze tam dil adına çözümler ve bu kodu iletir; bu nedenle eval aramayı yinelemeden her satırı etiketleyebilir:

CREATE OR REPLACE FUNCTION tag_language(t TABLE, lang_code STRING)
RETURNS TABLE
LANGUAGE PYTHON
HANDLER 'TagLanguage'
AS $$
class TagLanguage:
    @staticmethod
    def analyze(t, lang_code):
        from dataclasses import dataclass
        from pyspark.sql.types import StructType, StructField, StringType
        from pyspark.sql.udtf import AnalyzeResult

        @dataclass
        class LangResult(AnalyzeResult):
            language: str = ""

        # Resolve the language code to a full name once during planning
        languages = {"en": "English", "es": "Spanish", "fr": "French", "de": "German"}
        return LangResult(
            schema=StructType([
                StructField("text", StringType()),
                StructField("language", StringType())
            ]),
            language=languages.get(lang_code.value, "Unknown")
        )

    def __init__(self, result):
        self._language = result.language

    def eval(self, row, lang_code: str):
        # Tag each row with the pre-resolved language name
        yield (row['text'], self._language)
$$;

SELECT * FROM tag_language(
  TABLE(VALUES ('Hola mundo'), ('Buenos días') t(text)),
  'es'
);
+-------------+----------+
| text        | language |
+-------------+----------+
| Hola mundo  | Spanish  |
| Buenos días | Spanish  |
+-------------+----------+

İletme durumuyla ilgili daha fazla şema ve ayrıntı için bkz. eval.

analyze yönteminden bölümlemeyi belirtin

Çok biçimli bir UDTF, bir tablo bağımsız değişkenini kabul ettiğinde, analyze yöntemi partitionBy, orderBy, withSinglePartition ve select ayarlarını yaparak giriş satırlarının UDTF örnekleri arasında nasıl dağıtıldığını AnalyzeResult üzerinde denetleyebilir. Bu, çağıranların SQL'de PARTITION BY veya ORDER BY belirtme gereksinimini ortadan kaldırır.

Tam bölümlendirme API'si ve örnekleri için bkz. Giriş satırlarının yöntemle bölümlendirilmesini analyze belirtme.

Ortam yalıtımı

Uyarı

Paylaşılan yalıtım ortamları Databricks Runtime 17.2 ve üzerini gerektirir. Önceki sürümlerde, tüm Unity Kataloğu Python UDF'leri katı yalıtım modunda çalışır.

Aynı sahip ve oturuma sahip Unity Kataloğu Python UDF'leri varsayılan olarak bir yalıtım ortamını paylaşabilir. Bu, başlatılması gereken ayrı ortamların sayısını azaltarak performansı artırır ve bellek kullanımını azaltır.

Katı yalıtım

UDTF'nin her zaman kendi, tamamen yalıtılmış ortamında çalıştığından emin olmak için karakteristik yan tümcesini STRICT ISOLATION ekleyin.

Çoğu UDF'nin katı yalıtıma ihtiyacı yoktur. Standart veri işleme UDF'leri varsayılan paylaşılan yalıtım ortamından yararlanır ve daha düşük bellek tüketimiyle daha hızlı çalışır.

STRICT ISOLATION Karakteristik yan tümcesini UDTF'lere ekleyin:

  • eval(), exec() veya benzer işlevleri kullanarak girdiyi kod olarak çalıştırın.
  • Dosyaları yerel dosya sistemine yazın.
  • Genel değişkenleri veya sistem durumunu değiştirin.
  • Ortam değişkenlerine erişme veya değişkenleri değiştirme.

Aşağıdaki UDTF örneği özel bir ortam değişkeni ayarlar, değişkeni geri okur ve değişkeni kullanarak bir sayı kümesini çarpar. UDTF, işlem ortamını değiştirdiği için STRICT ISOLATION içinde çalıştırın. Aksi takdirde, aynı ortamdaki diğer UDF'ler/UDTF'ler için ortam değişkenlerini sızdırabilir veya geçersiz kılabilir ve bu da yanlış davranışa neden olabilir.

CREATE OR REPLACE TEMPORARY FUNCTION multiply_numbers(factor STRING)
RETURNS TABLE (original INT, scaled INT)
LANGUAGE PYTHON
STRICT ISOLATION
HANDLER 'Multiplier'
AS $$
import os

class Multiplier:
    def eval(self, factor: str):
        # Save the factor as an environment variable
        os.environ["FACTOR"] = factor

        # Read it back and convert it to a number
        scale = int(os.getenv("FACTOR", "1"))

        # Multiply 0 through 4 by the factor
        for i in range(5):
            yield (i, i * scale)
$$;

SELECT * FROM multiply_numbers("3");

İşlevinizin tutarlı sonuçlar üretip üretmediğini ayarlama DETERMINISTIC

Aynı girişler için aynı çıkışları oluşturuyorsa işlev tanımınıza ekleyin DETERMINISTIC . Bu, sorgu iyileştirmelerinin performansı geliştirmesine olanak tanır.

Varsayılan olarak, Batch Unity Kataloğu Python UDF'lerinin açıkça bildirilmediği sürece belirleyici olmadığı varsayılır. Belirlenemeyen işlevlere örnek olarak şunlar verilebilir: rastgele değerler oluşturma, geçerli saatlere veya tarihlere erişme veya dış API çağrıları yapma.

Bkz.CREATE FUNCTION (SQL, Python, Scala ve Java).

Pratik örnekler

Aşağıdaki örneklerde, basit veri dönüştürmelerinden karmaşık dış tümleştirmelere kadar ilerleyen Unity Kataloğu Python UDF'leri için gerçek dünya kullanım örnekleri gösterilmektedir.

Örnek: Yeniden uygulama explode

Spark yerleşik explode işlevini sağlarken, kendi versiyonunuzu oluşturmak, tek bir girdi alıp birden fazla çıktı satırı üretmenin temel UDTF modelini gösterir.

CREATE OR REPLACE FUNCTION my_explode(arr ARRAY<STRING>)
RETURNS TABLE (element STRING)
LANGUAGE PYTHON
HANDLER 'MyExplode'
DETERMINISTIC
AS $$
class MyExplode:
    def eval(self, arr):
        if arr is None:
            return
        for element in arr:
            yield (element,)
$$;

İşlevi doğrudan bir SQL sorgusunda kullanın:

SELECT element FROM my_explode(array('apple', 'banana', 'cherry'));
+---------+
| element |
+---------+
| apple   |
| banana  |
| cherry  |
+---------+

Veya birleştirme ile LATERALmevcut tablo verilerine uygulayın:

SELECT s.*, e.element
FROM my_items AS s,
LATERAL my_explode(s.items) AS e;

Örnek: REST API aracılığıyla IP adresi coğrafi konumu

Bu örnek, UDF'lerin dış API'leri doğrudan SQL iş akışınızla nasıl tümleştirebileceğini gösterir. Analistler, ayrı ETL işlemleri gerektirmeden tanıdık SQL söz dizimlerini kullanarak verileri gerçek zamanlı API çağrılarıyla zenginleştirebilir.

CREATE OR REPLACE FUNCTION ip_to_location(ip_address STRING)
RETURNS TABLE (city STRING, country STRING)
LANGUAGE PYTHON
HANDLER 'IPToLocationAPI'
AS $$
class IPToLocationAPI:
    def eval(self, ip_address):
        import requests
        api_url = f"https://api.ip-lookup.example.com/{ip_address}"
        try:
            response = requests.get(api_url)
            response.raise_for_status()
            data = response.json()
            yield (data.get('city'), data.get('country'))
        except requests.exceptions.RequestException as e:
            # Return nothing if the API request fails
            return
$$;

Uyarı

Python UDTF'leri, 80, 443 ve 53 numaralı bağlantı noktaları üzerinden TCP/UDP ağ trafiğine, sunucusuz hesaplama veya standart erişim modunda yapılandırılmış hesaplama kullanılırken izin verir.

Web günlüğü verilerini coğrafi bilgilerle zenginleştirmek için işlevini kullanın:

SELECT
  l.timestamp,
  l.request_path,
  geo.city,
  geo.country
FROM web_logs AS l,
LATERAL ip_to_location(l.ip_address) AS geo;

Bu yaklaşım, önceden işlenmiş arama tabloları veya ayrı veri işlem hatları gerektirmeden gerçek zamanlı coğrafi analiz sağlar. UDTF HTTP isteklerini, JSON ayrıştırma ve hata işlemeyi işleyerek dış veri kaynaklarını standart SQL sorguları aracılığıyla erişilebilir hale getirir.

Örnek: IP adreslerini CIDR ağ bloklarıyla eşleştirme

Bu örnek, karmaşık SQL mantığı gerektiren yaygın bir veri mühendisliği görevi olan CIDR ağ bloklarıyla eşleşen IP adreslerini gösterir.

İlk olarak, hem IPv4 hem de IPv6 adresleriyle örnek veriler oluşturun:

-- An example IP logs with both IPv4 and IPv6 addresses
CREATE OR REPLACE TEMPORARY VIEW ip_logs AS
VALUES
  ('log1', '192.168.1.100'),
  ('log2', '10.0.0.5'),
  ('log3', '172.16.0.10'),
  ('log4', '8.8.8.8'),
  ('log5', '2001:db8::1'),
  ('log6', '2001:db8:85a3::8a2e:370:7334'),
  ('log7', 'fe80::1'),
  ('log8', '::1'),
  ('log9', '2001:db8:1234:5678::1')
t(log_id, ip_address);

Ardından UDTF'yi tanımlayın ve kaydedin. Python sınıf yapısına dikkat edin:

  • t TABLE parametresi herhangi bir şemaya sahip bir giriş tablosu kabul eder. UDTF, sağlanan sütunları işlemek için otomatik olarak uyarlanır. Bu esneklik, işlev imzasını değiştirmeden farklı tablolarda aynı işlevi kullanabileceğiniz anlamına gelir. Ancak, uyumluluğu sağlamak için satırların şemasını dikkatle denetlemeniz gerekir.
  • __init__ yöntemi, büyük ağ listesini yükleme gibi tek seferlik ağır kurulumlar için kullanılır. Bu çalışma, giriş tablosunun bölümü başına bir kez gerçekleşir.
  • eval yöntemi her satırı işler ve çekirdek eşleştirme mantığını içerir. Bu yöntem giriş bölümündeki her satır için tam olarak bir kez yürütülür ve her yürütme ilgili bölüm için UDTF sınıfının karşılık gelen örneği IpMatcher tarafından gerçekleştirilir.
  • HANDLER yan tümcesi, UDTF mantığını uygulayan Python sınıfının adını belirtir.
CREATE OR REPLACE TEMPORARY FUNCTION ip_cidr_matcher(t TABLE)
RETURNS TABLE(log_id STRING, ip_address STRING, network STRING, ip_version INT)
LANGUAGE PYTHON
HANDLER 'IpMatcher'
COMMENT 'Match IP addresses against a list of network CIDR blocks'
AS $$
class IpMatcher:
    def __init__(self):
        import ipaddress
        # Heavy initialization - load networks once per partition
        self.nets = []
        cidrs = ['192.168.0.0/16', '10.0.0.0/8', '172.16.0.0/12',
                 '2001:db8::/32', 'fe80::/10', '::1/128']
        for cidr in cidrs:
            self.nets.append(ipaddress.ip_network(cidr))

    def eval(self, row):
        import ipaddress
	    # Validate that required fields exist
        required_fields = ['log_id', 'ip_address']
        for field in required_fields:
            if field not in row:
                raise ValueError(f"Missing required field: {field}")
        try:
            ip = ipaddress.ip_address(row['ip_address'])
            for net in self.nets:
                if ip in net:
                    yield (row['log_id'], row['ip_address'], str(net), ip.version)
                    return
            yield (row['log_id'], row['ip_address'], None, ip.version)
        except ValueError:
            yield (row['log_id'], row['ip_address'], 'Invalid', None)
$$;

Unity Kataloğu'na kaydedildiğinden ip_cidr_matcher, TABLE() söz dizimini kullanarak SQL'den doğrudan çağırabilirsiniz.

-- Process all IP addresses
SELECT
  *
FROM
  ip_cidr_matcher(t => TABLE(ip_logs))
ORDER BY
  log_id;
+--------+-------------------------------+-----------------+-------------+
| log_id | ip_address                    | network         | ip_version  |
+--------+-------------------------------+-----------------+-------------+
| log1   | 192.168.1.100                 | 192.168.0.0/16  | 4           |
| log2   | 10.0.0.5                      | 10.0.0.0/8      | 4           |
| log3   | 172.16.0.10                   | 172.16.0.0/12   | 4           |
| log4   | 8.8.8.8                       | null            | 4           |
| log5   | 2001:db8::1                   | 2001:db8::/32   | 6           |
| log6   | 2001:db8:85a3::8a2e:370:7334  | 2001:db8::/32   | 6           |
| log7   | fe80::1                       | fe80::/10       | 6           |
| log8   | ::1                           | ::1/128         | 6           |
| log9   | 2001:db8:1234:5678::1         | 2001:db8::/32   | 6           |
+--------+-------------------------------+-----------------+-------------+

Örnek: Azure Databricks görsel uç noktalarını kullanarak toplu görüntü altyazılama

Bu örnekte uç noktaya hizmet veren bir Azure Databricks görüntü modeli kullanılarak toplu görüntü alt yazıları gösterilmektedir. terminate() kullanılarak toplu işleme ve bölüm tabanlı yürütme sergilenir.

  1. Genel görüntü URL'leriyle tablo oluşturma:

    CREATE OR REPLACE TEMPORARY VIEW sample_images AS
    VALUES
        ('https://upload.wikimedia.org/wikipedia/commons/thumb/d/dd/Gfp-wisconsin-madison-the-nature-boardwalk.jpg/2560px-Gfp-wisconsin-madison-the-nature-boardwalk.jpg', 'scenery'),
        ('https://upload.wikimedia.org/wikipedia/commons/thumb/a/a7/Camponotus_flavomarginatus_ant.jpg/1024px-Camponotus_flavomarginatus_ant.jpg', 'animals'),
        ('https://upload.wikimedia.org/wikipedia/commons/thumb/1/15/Cat_August_2010-4.jpg/1200px-Cat_August_2010-4.jpg', 'animals'),
        ('https://upload.wikimedia.org/wikipedia/commons/thumb/c/c5/M101_hires_STScI-PRC2006-10a.jpg/1024px-M101_hires_STScI-PRC2006-10a.jpg', 'scenery')
    images(image_url, category);
    
  2. Görüntü açıklamalı alt yazıları oluşturmak için Unity Kataloğu Python UDTF'sini oluşturun:

    1. Toplu iş boyutu, Azure Databricks API belirteci, görüntü modeli uç noktası ve çalışma alanı URL'si gibi yapılandırmayla UDTF'yi başlatın.
    2. eval yöntemine göre, görüntü URL'lerini bir arabelleğe toplayın. Arabellek toplu iş boyutuna ulaştığında toplu işlem tetikler. Bu, görüntü başına tek tek çağrılar yerine tek bir API çağrısında birden çok görüntünün birlikte işlenmesini sağlar.
    3. Toplu işleme yönteminde tüm arabelleğe alınan görüntüleri elde edin, base64 olarak kodlayın ve Databricks VisionModel API'sine tek bir API isteği olarak gönderin. Model, tüm görüntüleri aynı anda işler ve toplu işlemin tamamı için açıklamalı alt yazılar döndürür.
    4. terminate yöntemi, her bölümün sonunda tam olarak bir kez yürütülür. terminate yönteminde arabellekteki kalan görüntüleri işleyin ve toplanan tüm açıklamalı alt yazıları sonuç olarak verin.

Uyarı

<workspace-url> ile ()https://your-workspace.cloud.databricks.com değerini, gerçek Azure Databricks çalışma alanı URL'nizle değiştirin.

CREATE OR REPLACE TEMPORARY FUNCTION batch_inference_image_caption(data TABLE, api_token STRING)
RETURNS TABLE (caption STRING)
LANGUAGE PYTHON
HANDLER 'BatchInferenceImageCaption'
COMMENT 'batch image captioning by sending groups of image URLs to a Databricks vision endpoint and returning concise captions for each image.'
AS $$
class BatchInferenceImageCaption:
    def __init__(self):
        self.batch_size = 3
        self.vision_endpoint = "databricks-claude-sonnet-4-5"
        self.workspace_url = "<workspace-url>"
        self.image_buffer = []
        self.results = []

    def eval(self, row, api_token):
        self.image_buffer.append((str(row[0]), api_token))
        if len(self.image_buffer) >= self.batch_size:
            self._process_batch()

    def terminate(self):
        if self.image_buffer:
            self._process_batch()
        for caption in self.results:
            yield (caption,)

    def _process_batch(self):
        batch_data = self.image_buffer.copy()
        self.image_buffer.clear()

        import base64
        import httpx
        import requests

        # API request timeout in seconds
        api_timeout = 60
        # Maximum tokens for vision model response
        max_response_tokens = 300
        # Temperature controls randomness (lower = more deterministic)
        model_temperature = 0.3

        # create a batch for the images
        batch_images = []
        api_token = batch_data[0][1] if batch_data else None

        for image_url, _ in batch_data:
            image_response = httpx.get(image_url, timeout=15)
            image_data = base64.standard_b64encode(image_response.content).decode("utf-8")
            batch_images.append(image_data)

        content_items = [{
            "type": "text",
            "text": "Provide brief captions for these images, one per line."
        }]
        for img_data in batch_images:
            content_items.append({
                "type": "image_url",
                "image_url": {
                    "url": "data:image/jpeg;base64," + img_data
                }
            })

        payload = {
            "messages": [{
                "role": "user",
                "content": content_items
            }],
            "max_tokens": max_response_tokens,
            "temperature": model_temperature
        }

        response = requests.post(
            self.workspace_url + "/serving-endpoints/" +
            self.vision_endpoint + "/invocations",
            headers={
                'Authorization': 'Bearer ' + api_token,
                'Content-Type': 'application/json'
            },
            json=payload,
            timeout=api_timeout
        )

        result = response.json()
        batch_response = result['choices'][0]['message']['content'].strip()

        lines = batch_response.split('\n')
        captions = [line.strip() for line in lines if line.strip()]

        while len(captions) < len(batch_data):
            captions.append(batch_response)

        self.results.extend(captions[:len(batch_data)])
$$;

Toplu resim yazısı UDTF'sini kullanmak için örnek görüntüler tablosunu kullanarak çağırın:

Uyarı

Databricks API belirtecinin your_secret_scope gerçek gizli kapsamı ve anahtar adını api_token ile değiştirin.

SELECT
  caption
FROM
  batch_inference_image_caption(
    data => TABLE(sample_images),
    api_token => secret('your_secret_scope', 'api_token')
  )
+---------------------------------------------------------------------------------------------------------------+
| caption                                                                                                       |
+---------------------------------------------------------------------------------------------------------------+
| Wooden boardwalk cutting through vibrant wetland grasses under blue skies                                     |
| Black ant in detailed macro photography standing on a textured surface                                        |
| Tabby cat lounging comfortably on a white ledge against a white wall                                          |
| Stunning spiral galaxy with bright central core and sweeping blue-white arms against the black void of space. |
+---------------------------------------------------------------------------------------------------------------+

Kategoriye göre resim yazısı kategorisi de oluşturabilirsiniz:

SELECT
  *
FROM
  batch_inference_image_caption(
    TABLE(sample_images)
    PARTITION BY category ORDER BY (category),
    secret('your_secret_scope', 'api_token')
  )
+------------------------------------------------------------------------------------------------------+
| caption                                                                                              |
+------------------------------------------------------------------------------------------------------+
| Black ant in detailed macro photography standing on a textured surface                               |
| Stunning spiral galaxy with bright center and sweeping blue-tinged arms against the black of space.  |
| Tabby cat lounging comfortably on white ledge against white wall                                     |
| Wooden boardwalk cutting through lush wetland grasses under blue skies                               |
+------------------------------------------------------------------------------------------------------+

Örnek: ML modeli değerlendirmesi için ROC eğrisi ve AUC hesaplaması

Bu örnekte, scikit-learn kullanılarak ikili sınıflandırma modeli değerlendirmesi için bilgi işlem alıcısı çalışma özelliği (ROC) eğrileri ve eğrinin altındaki alan (AUC) puanları gösterilmektedir.

Bu örnekte çeşitli önemli desenler yer alınıyor:

  • Dış kitaplık kullanımı: "ROC eğrisi" hesaplaması için scikit-learn'i tümleştirir
  • Durum Korumalı Toplama: Ölçümleri hesaplamadan önce tüm satırlardaki tahminleri bir araya getirir
  • terminate() yöntem kullanımı: Tam veri kümesini işler ve sonuçları yalnızca tüm satırlar değerlendirildikten sonra verir
  • Hata işleme: Giriş tablosunda gerekli sütunların mevcut olduğunu doğrular

UDTF, eval() yöntemini kullanarak tüm tahminleri bellekte biriktirir ve ardından terminate() yönteminde tam ROC eğrisini hesaplayarak verir. Bu düzen, hesaplama için tam veri kümesini gerektiren ölçümler için kullanışlıdır.

CREATE OR REPLACE TEMPORARY FUNCTION compute_roc_curve(t TABLE)
RETURNS TABLE (threshold DOUBLE, true_positive_rate DOUBLE, false_positive_rate DOUBLE, auc DOUBLE)
LANGUAGE PYTHON
HANDLER 'ROCCalculator'
COMMENT 'Compute ROC curve and AUC using scikit-learn'
AS $$
class ROCCalculator:
    def __init__(self):
        from sklearn import metrics
        self._roc_curve = metrics.roc_curve
        self._roc_auc_score = metrics.roc_auc_score

        self._true_labels = []
        self._predicted_scores = []

    def eval(self, row):
        if 'y_true' not in row or 'y_score' not in row:
            raise KeyError("Required columns 'y_true' and 'y_score' not found")

        true_label = row['y_true']
        predicted_score = row['y_score']

        label = float(true_label)
        self._true_labels.append(label)
        self._predicted_scores.append(float(predicted_score))

    def terminate(self):
        false_pos_rate, true_pos_rate, thresholds = self._roc_curve(
            self._true_labels,
            self._predicted_scores,
            drop_intermediate=False
        )

        auc_score = float(self._roc_auc_score(self._true_labels, self._predicted_scores))

        for threshold, tpr, fpr in zip(thresholds, true_pos_rate, false_pos_rate):
            yield float(threshold), float(tpr), float(fpr), auc_score
$$;

Tahminlerle örnek ikili sınıflandırma verileri oluşturun:

CREATE OR REPLACE TEMPORARY VIEW binary_classification_data AS
SELECT *
FROM VALUES
  ( 1, 1.0, 0.95, 'high_confidence_positive'),
  ( 2, 1.0, 0.87, 'high_confidence_positive'),
  ( 3, 1.0, 0.82, 'medium_confidence_positive'),
  ( 4, 0.0, 0.78, 'false_positive'),
  ( 5, 1.0, 0.71, 'medium_confidence_positive'),
  ( 6, 0.0, 0.65, 'false_positive'),
  ( 7, 0.0, 0.58, 'true_negative'),
  ( 8, 1.0, 0.52, 'low_confidence_positive'),
  ( 9, 0.0, 0.45, 'true_negative'),
  (10, 0.0, 0.38, 'true_negative'),
  (11, 1.0, 0.31, 'low_confidence_positive'),
  (12, 0.0, 0.15, 'true_negative'),
  (13, 0.0, 0.08, 'high_confidence_negative'),
  (14, 0.0, 0.03, 'high_confidence_negative')
AS data(sample_id, y_true, y_score, prediction_type);

ROC eğrisini ve AUC'yi hesaplama:

SELECT
    threshold,
    true_positive_rate,
    false_positive_rate,
    auc
FROM compute_roc_curve(
  TABLE(
    SELECT y_true, y_score
    FROM binary_classification_data
    WHERE y_true IS NOT NULL AND y_score IS NOT NULL
    ORDER BY sample_id
  )
)
ORDER BY threshold DESC;
+-----------+---------------------+----------------------+-------+
| threshold | true_positive_rate  | false_positive_rate  | auc   |
+-----------+---------------------+----------------------+-------+
| 1.95      | 0.0                 | 0.0                  | 0.786 |
| 0.95      | 0.167               | 0.0                  | 0.786 |
| 0.87      | 0.333               | 0.0                  | 0.786 |
| 0.82      | 0.5                 | 0.0                  | 0.786 |
| 0.78      | 0.5                 | 0.125                | 0.786 |
| 0.71      | 0.667               | 0.125                | 0.786 |
| 0.65      | 0.667               | 0.25                 | 0.786 |
| 0.58      | 0.667               | 0.375                | 0.786 |
| 0.52      | 0.833               | 0.375                | 0.786 |
| 0.45      | 0.833               | 0.5                  | 0.786 |
| 0.38      | 0.833               | 0.625                | 0.786 |
| 0.31      | 1.0                 | 0.625                | 0.786 |
| 0.15      | 1.0                 | 0.75                 | 0.786 |
| 0.08      | 1.0                 | 0.875                | 0.786 |
| 0.03      | 1.0                 | 1.0                  | 0.786 |
+-----------+---------------------+----------------------+-------+

Örnek: Tablo bağımsız değişkeni üzerinden dinamik sütun projeksiyonu işlemi

Bu örnek, çok biçimli UDF'leri tablo bağımsız değişkenleriyle birleştirir. UDTF bir tabloyu ve virgülle ayrılmış bir sütun adı listesini kabul eder, ardından girişten yalnızca bu sütunları projeler. yöntemi, analyze giriş tablosunun şemasını inceler ve yalnızca istenen sütunları içeren bir çıkış şeması oluşturur.

CREATE OR REPLACE FUNCTION project_columns(t TABLE, columns STRING)
RETURNS TABLE
LANGUAGE PYTHON
HANDLER 'ProjectColumns'
AS $$
class ProjectColumns:
    @staticmethod
    def analyze(t, columns):
        from pyspark.sql.types import StructType
        from pyspark.sql.udtf import AnalyzeResult

        requested = [c.strip() for c in columns.value.split(",")]
        input_schema = t.dataType
        output_fields = []
        for field in input_schema.fields:
            if field.name in requested:
                output_fields.append(field)
        if not output_fields:
            raise ValueError(
                f"None of the requested columns {requested} "
                f"exist in the input table"
            )
        return AnalyzeResult(schema=StructType(output_fields))

    def eval(self, row, columns: str):
        requested = [c.strip() for c in columns.split(",")]
        yield tuple(row[col] for col in requested if col in row)
$$;

Tablodan belirli sütunları seçmek için işlevini kullanın:

SELECT * FROM project_columns(
  TABLE(SELECT * FROM samples.nyctaxi.trips LIMIT 5),
  'pickup_zip, dropoff_zip, fare_amount'
);

Sınırlamalar

Unity Kataloğu Python UDF'leri için aşağıdaki sınırlamalar geçerlidir:

Ek kaynaklar