Not
Bu sayfaya erişim yetkilendirme gerektiriyor. Oturum açmayı veya dizinleri değiştirmeyi deneyebilirsiniz.
Bu sayfaya erişim yetkilendirme gerektiriyor. Dizinleri değiştirmeyi deneyebilirsiniz.
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:
- Unity Kataloğu: Unity Kataloğu'nda UDTF'yi yönetilen nesne olarak kaydedin.
- Oturum kapsamlı: Geçerli not defterine veya işe izole edilmiş yerel
SparkSessionöğesine kaydolun. Bkz. Python kullanıcı tanımlı tablo işlevleri (UDF).
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 TABLEparametresi 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. -
evalyö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ğiIpMatchertarafından gerçekleştirilir. -
HANDLERyan 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.
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);Görüntü açıklamalı alt yazıları oluşturmak için Unity Kataloğu Python UDTF'sini oluşturun:
- 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.
-
evalyö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. - 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.
-
terminateyö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:
- Unity Kataloğu hizmeti için kimlik bilgileri desteklenmez.
- Özel bağımlılıklar desteklenmez.