Lakeflow Designer'da kullanıcı tanımlı işleçler

Lakeflow Designer, yerleşik işleçlerin yanı sıra doğrudan tuvalde görünen kullanıcı tanımlı işleçler oluşturmanıza olanak tanır. Lakeflow Designer'ı kendi iş mantığınız, hesaplamalarınız veya tümleştirmelerinizle genişletmek için bunları kullanın.

Üç tür kullanıcı tanımlı işleç vardır:

  • python-run-function: Çalışma alanında depolanan satır içi Python içeren tek başına BIR YAML dosyası. DataFrame düzeyinde dönüştürmeler ve dış tümleştirmeler için en iyi yöntemdir. İzinler çalışma alanı dosya düzeyinde yönetilir.
  • uc-udf: Unity Catalog skaler fonksiyonunu sarar. Sütun düzeyindeki dönüşümler için en uygunudur. Access, Unity Kataloğu izinlerine tabidir.
  • uc-udtf: Unity Catalog tablo döndüren bir işlevi sarmalar. ML kümelemesi ve toplama gibi tablo düzeyinde dönüşümler için en iyi yöntemdir. Access, Unity Kataloğu izinlerine tabidir.
Özellik python-run-function uc-udf uc-udtf
Örnek kullanım örneği DataFrame dönüşümleri, API tümleştirmeleri, e-posta bildirimleri Sütun düzeyinde hesaplamalar (BMI, faiz oranları) ML kümelemesi, satırlar genelinde toplulaştırma
Giriş Veri Çerçeveleri Tek değerler Tablonun tamamı, satır satır
Çıkış Veri Çerçeveleri Tek değer Tablo (birden çok satır)
Unity Kataloğu işlevi gerektirir Hayır Yes Yes
Erişim idaresi Çalışma alanı dosyası izinleri Unity Kataloğu izinleri (EXECUTE, USE SCHEMA) Unity Kataloğu izinleri (EXECUTE, USE SCHEMA)
Desteklenen diller Yalnızca Python Bir SQL sarmalayıcısında SQL veya Python Bir SQL sarmalayıcısında SQL veya Python

Kullanıcı tanımlı işleçler nasıl çalışır?

Kullanıcı tanımlı işleç şunlardan oluşur:

  • İşleç mantığı: İşleç yürütürken çalışan kod. Bu satır içi bir Python run() işlevi (python-run-function için) veya Unity Catalog işlevi (uc-udf ve uc-udtf için) olabilir.
  • YAML yapılandırması: Lakeflow Designer'a, işlecin adı, açıklaması, giriş parametreleri, kullanıcı arabirimi pencere öğeleri ve bağlantı noktaları dahil olmak üzere kullanıcı arabiriminde işleci nasıl sunması gerektiğini bildirir. Tüm işleç türleri şemayı user-defined-operator-v0.1.0 kullanır.
  • Kayıt dosyası: Lakeflow Designer'ın .user_defined_operators.yaml işleci bulmasına olanak tanıyan bir girdi.

İşleç mantığı

Python işlev kullanıcı tanımlı işleç mantığını çalıştırma

Her python-run-function işleç bir run() işlev tanımlamalıdır:

def run(config: Dict[str, Any], inputs: Dict[str, Any], spark) -> Dict[str, Any]:
  • config: Kullanıcı arabiriminden kullanıcı tarafından yapılandırılan ve özellik adına göre anahtarlanan değerler.
  • inputs: Giriş bağlantı noktası nametarafından anahtarlanan Giriş Veri Çerçeveleri.
  • spark: Aktif SparkSession.
  • Döndürür: Çıkış bağlantı noktası name değerlerini DataFrame'lerle eşleyen bir sözlük.

Aşağıdaki örnek, bir giriş DataFrame'indeki satırları filtreler:

def run(config, inputs, spark):
    df = inputs["in"]
    filtered = df.filter(config["filter_expression"])
    return {"out": filtered}

Operatörünüz harici pip paketleri gerektiriyorsa alanı YAML'ye ekleyin environment :

environment:
  environment_version: '4'
  dependencies:
    - requests==2.31.0
    - beautifulsoup4==4.12.0

UDF ve UDTF işleç mantığı

UC fonksiyonlarını SQL veya Python'da yazabilirsiniz. Python işlevleri sql CREATE FUNCTION deyiminde sarmalanmıştır:

SQL işlevi:

CREATE OR REPLACE FUNCTION my_catalog.my_schema.calculate_bmi(weight_kg DOUBLE, height_m DOUBLE)
RETURNS DOUBLE
LANGUAGE SQL
RETURN
  SELECT weight_kg / (height_m * height_m);

Python fonksiyonu (SQL ile sarmalanmış):

CREATE OR REPLACE FUNCTION my_catalog.my_schema.calculate_bmi(weight_kg DOUBLE, height_m DOUBLE)
RETURNS DOUBLE
LANGUAGE PYTHON
AS $$
  return weight_kg / (height_m ** 2)
$$;

UDF'ler tek seferde tek bir değeri işler ve hesaplanan bir değer döndürür. UDTF'ler tabloları satır satır işler ve tüm satırlar boyunca durumu koruyabilir. Sütun düzeyindeki dönüşümler için uc-udf ve ML kümeleme veya toplulaştırma gibi işlemler için uc-udtf kullanın.

Ek olarak, UDF'ler üç anahtar yöntemi tanımlamanızı gerektirir: __init__(), eval()ve terminate():

class MyOperator:
    def __init__(self):
        # Called before processing - initialize any values needed.

    def eval(self, row, id_column, columns, k):
        # Called one time per input row - accumulate data here.

    def terminate(self):
        # Called after all rows - perform final calculations and yield results.

Note

UDTF dönüş tabloları sabit ve açık türlere sahip olmalıdır. Dönüş yapılandırmasında giriş sütunu türlerine başvuramazsınız.

YAML yapılandırması

YAML yapılandırması, Lakeflow Designer'a kullanıcı arabiriminde işleci nasıl sunması gerektiğini bildirir. Operatörün adını, açıklamasını, giriş parametrelerini, kullanıcı arabirimi pencere öğelerini ve bağlantı noktalarını tanımlar. Her yapılandırma alanı tür, başlık ve isteğe bağlı x-ui pencere öğesi ipuçları içeren bir özelliktir:

config:
  type: object
  properties:
    my_param:
      type: string
      title: My Parameter
      x-ui:
        widget: input
    my_expression:
      type: string
      title: Column
      format: expression
      x-ui:
        widget: expression
        port: in
    my_number:
      type: number
      title: Count
      default: 10
      minimum: 0
      maximum: 100
  required:
    - my_param
    - my_expression

Tüm pencere öğesi türleri ve yapılandırma seçenekleri de dahil olmak üzere YAML şemasıyla ilgili tüm ayrıntılar için bkz. Kullanıcı tanımlı işleç YAML başvurusu.

Limanlar

Bağlantı noktaları, işleciniz için girişleri ve çıkışları tanımlar:

ports:
  input:
    - name: in
      title: Input Data
      mime: application/vnd.databricks.dataframe
      required: true
      allowMultiple: false
  output:
    - name: out
      title: Output Data

Python çalıştırma işlevi işleçleri için YAML

python-run-function işleçleri için YAML dosyası tek başınadır ve satır içi Python kodu içeren bir run_function alanı içerir:

schema: user-defined-operator-v0.1.0
type: python-run-function
name: Filter Rows
id: filter_rows
version: '1.0.0'
description: Filters rows based on a SQL expression.
config:
  type: object
  properties:
    filter_expression:
      type: string
      title: Filter Expression
      x-ui:
        widget: input
  required:
    - filter_expression
ports:
  input:
    - name: in
      title: Input
  output:
    - name: out
      title: Output
run_function:
  type: inline
  code: |
    def run(config, inputs, spark):
        df = inputs["in"]
        filtered = df.filter(config["filter_expression"])
        return {"out": filtered}

Unity Kataloğu işlevleri için YAML

UC tabanlı işleçler için YAML yapılandırmasını işlevinize açıklama veya docstring olarak ekleyin.

SQL'de (/* ... */ yorumunu kullanın):

RETURN(/*
  schema: user-defined-operator-v0.1.0
  type: uc-udf
  name: Calculate BMI
  id: calculate_bmi
  version: "1.0.0"
  description: Calculates BMI from weight and height.
  config:
    type: object
    properties:
      weight_kg:
        type: string
        title: Weight (in kg)
        format: expression
        x-ui:
          widget: expression
          port: in
      height_m:
        type: string
        title: Height (in meters)
        format: expression
        x-ui:
          widget: expression
          port: in
    required:
      - weight_kg
      - height_m
  ports:
    input:
      - name: in
        title: Input Data
    output:
      - name: out
        title: Output
    */
  SELECT weight_kg / (height_m * height_m)
);

Python (""" ... """ docstring kullanın):

AS $$
  """
  schema: user-defined-operator-v0.1.0
  type: uc-udf
  name: Calculate BMI
  id: calculate_bmi
  version: "1.0.0"
  description: Calculates BMI from weight and height.
  config:
    type: object
    properties:
      weight_kg:
        type: string
        title: Weight (in kg)
        format: expression
        x-ui:
          widget: expression
          port: in
      height_m:
        type: string
        title: Height (in meters)
        format: expression
        x-ui:
          widget: expression
          port: in
    required:
      - weight_kg
      - height_m
  ports:
    input:
      - name: in
        title: Input Data
    output:
      - name: out
        title: Output
  """

  return weight_kg / (height_m ** 2)
$$;

Operatörünüzü kaydetme ve Lakeflow Designer'a dağıtma

Operatörünüzün Lakeflow Designer'da görünmesi için bunu bir .user_defined_operators.yaml dosyaya kaydedin:

  • Çalışma alanı düzeyi: İşleci tüm kullanıcılara görünür hale getirmek için dosyayı çalışma alanınızın köküne yerleştirin.
  • Kullanıcı düzeyi: İşleçleri yalnızca sizin için görünür hale getirmek için dosyayı kullanıcı giriş klasörünüze (/Workspace/Users/<user-name>/.user_defined_operators.yaml) yerleştirin.

operators: bölümü dosya yollarını, Unity Kataloğu işlev başvurularını ve glob desenlerini destekler. Giriş türlerini karıştırabilirsiniz:

operators:
  # File path (python-run-function operators)
  - /Workspace/Users/me/udos/my_operator.yaml
  # Glob pattern (registers all matching files)
  - /Workspace/Users/me/udos/transforms/*.yaml
  # UC function reference (uc-udf and uc-udtf operators)
  - catalog: my_catalog
    schema: my_schema
    functionName: my_function

Bir işleci güncelleştirme veya kaldırma

Bir işlecin kodunu değiştirdiğinizde, değişikliği yüklemek için kullanıcı tanımlı işleçlerinizi yenileyin. Menünün İşleçler sekmesinde Yenile simgesine tıklayın..

  • Operatör aynı version kullanırsa, yenileme güncellenmiş kodu yükler.
  • İşlecin yeni bir version sürümü varsa, yenilediğinizde tuvaldeki işleç sizi buna yükseltmeniz (veya geçerli sürümü korumanız) için yönlendirir.

Lakeflow Designer'dan bir operatörü kaldırmak için .user_defined_operators.yaml içindeki ilgili girdiyi silin. uc-udf ve uc-udtf işleçleri için, artık ihtiyacınız yoksa, altyapıdaki Unity Catalog işlevini DROP FUNCTION ile de kaldırabilirsiniz.

Gelişmiş yapılandırmalar

Önizleme modu

Lakeflow Designer, tasarım modundayken önizlemeleri destekler. Dış API'leri çağıran veya dış sistemlere yazan işleçler için, önizleme sırasında yan etkileri atlayabileceğiniz bir is_preview yapılandırma özelliği ekleyin. Önizleme modu etkinleştirildiğinde, kullanıcıların işleci yan efektlerle yürütmek için Çalıştır'a açıkça tıklaması gerekir.

config:
  type: object
  properties:
    is_preview:
      type: boolean
      format: is_preview
      default: false

Lakeflow Designer, önizleme sırasında bu değeri otomatik olarak olarak true olarak ayarlar. Yan etkileri atlamak için bunu mantığınızda kontrol edin:

# In a python-run-function
if config.get("is_preview"):
    return {"out": inputs["in"]}

# In a UC function (SQL)
CASE WHEN is_preview THEN 'preview' ELSE /* actual work */ END

Unity Kataloğu bağlantıları

Dış API'leri çağıran UC tabanlı SQL işleçleri için, kimlik bilgilerini güvenli bir şekilde depolamak için Unity Kataloğu HTTP bağlantılarını kullanın:

CREATE CONNECTION my_api_connection TYPE HTTP OPTIONS (
  host 'https://api.example.com',
  port '443',
  base_path '/v1/',
  bearer_token 'your-token-here'
);

Ardından, SQL UDF'nizde http_request() işleviyle bağlantıyı kullanın. Ayrıntılar için bkz. Dış HTTP hizmetlerine bağlanma.

WorkspaceClient

python-run-function işleçleri için çalışma alanı kaynaklarına ve dış API'lere erişmek için Azure Databricks WorkspaceClient kullanabilirsiniz:

def run(config, inputs, spark):
    from databricks.sdk import WorkspaceClient
    w = WorkspaceClient()
    # Use w to access workspace resources

Tam bir python-run-function kullanıcı tanımlı operatörü oluşturun

Aşağıdaki adımlar, sıfırdan bir python-run-function operatör oluşturma adımlarını göstermektedir.

1. Adım: Mantığı tanımlama

İşlevinizi run() not defterine yazın:

from typing import Dict, Any

def run(config: Dict[str, Any], inputs: Dict[str, Any], spark) -> Dict[str, Any]:
    from pyspark.sql import functions as F
    df = inputs["in"]
    result = df.withColumn(config["column_name"], F.current_timestamp())
    return {"out": result}

2. Adım: İşlevi test edin

İşlevi örnek verilerle etkileşimli olarak test edin:

test_df = spark.createDataFrame(
    [("Alice", 100), ("Bob", 200)],
    ["name", "amount"]
)

result = run(
    config={"column_name": "processed_at"},
    inputs={"in": test_df},
    spark=spark
)

result["out"].show()

3. Adım: YAML yapılandırmasını oluşturma

YAML dosyasında işleç meta verilerini, yapılandırma alanlarını ve bağlantı noktalarını tanımlayın:

schema: user-defined-operator-v0.1.0
type: python-run-function
name: Add Timestamp
id: transforms.add_timestamp
version: '1.0.0'
description: Adds a timestamp column to the input DataFrame.
config:
  type: object
  properties:
    column_name:
      type: string
      title: Column Name
      default: processed_at
      x-ui:
        widget: input
  required:
    - column_name

4. Adım: Mantığı ve YAML'yi birleştirme

YAML dosyasının tamamını oluşturmak için run_function ve ports alanlarını ekleyin. Çalışma alanınıza kaydedin, örneğin /Workspace/Users/<user-name>/udos/add_timestamp.yaml:

schema: user-defined-operator-v0.1.0
type: python-run-function
name: Add Timestamp
id: transforms.add_timestamp
version: '1.0.0'
description: Adds a timestamp column to the input DataFrame.
config:
  type: object
  properties:
    column_name:
      type: string
      title: Column Name
      default: processed_at
      x-ui:
        widget: input
  required:
    - column_name
ports:
  input:
    - name: in
      title: Input
  output:
    - name: out
      title: Output
run_function:
  type: inline
  code: |
    from typing import Dict, Any

    def run(config: Dict[str, Any], inputs: Dict[str, Any], spark) -> Dict[str, Any]:
        from pyspark.sql import functions as F
        df = inputs["in"]
        result = df.withColumn(config["column_name"], F.current_timestamp())
        return {"out": result}

5. Adım: İşleci kaydetme

Dosyanıza .user_defined_operators.yaml dosya yolunu ekleyin:

operators:
  - /Workspace/Users/<user-name>/udos/add_timestamp.yaml

6. Adım: Lakeflow Designer'da operatörü kullanma

Lakeflow Designer'ı açın ve işlecin işleç paletinde göründüğünü doğrulayın. Tuvale sürükleyin, bir giriş bağlayın, sütun adını yapılandırın ve bir önizleme çalıştırın.

Tam bir UC kullanıcı tanımlı operatörü oluşturun

Aşağıdaki adımlar, UC tabanlı uc-udf bir işleç oluşturma adımlarını gösterir.

1. Adım: Mantığı tanımlama

İşlev mantığınızı not defterine yazın ve test edin:

def double_value(input_value: float) -> float:
    if input_value is None:
        return None
    return input_value * 2

2. Adım: YAML yapılandırmasını oluşturma

İşleç meta verilerini, yapılandırma alanlarını ve bağlantı noktalarını tanımlayın:

schema: user-defined-operator-v0.1.0
type: uc-udf
name: Double Value
id: math.double_value
version: '1.0.0'
description: Doubles the input value
config:
  type: object
  properties:
    input_value:
      type: string
      title: Input Value
      format: expression
      x-ui:
        widget: expression
        port: input_data
  required:
    - input_value
ports:
  input:
    - name: input_data
      title: Input
  output:
    - name: out
      title: Output

3. Adım: Mantığı ve YAML'yi birleştirme

Docstring olarak katıştırılmış YAML ile Unity Kataloğu işlevini oluşturun:

CREATE OR REPLACE FUNCTION main.my_schema.double_value(input_value DOUBLE)
RETURNS DOUBLE
LANGUAGE PYTHON
AS $$
  """
  schema: user-defined-operator-v0.1.0
  type: uc-udf
  name: Double Value
  id: math.double_value
  version: "1.0.0"
  description: Doubles the input value
  config:
    type: object
    properties:
      input_value:
        type: string
        title: Input Value
        format: expression
        x-ui:
          widget: expression
          port: input_data
    required:
      - input_value
  ports:
    input:
      - name: input_data
        title: Input
    output:
      - name: out
        title: Output
  """

  def double_value(input_value: float) -> float:
      if input_value is None:
          return None
      return input_value * 2

  return double_value(input_value)
$$

4. Adım: İşlevi test edin

SELECT main.my_schema.double_value(5) AS result;
-- Should return: 10

5. Adım: İşleci kaydetme

Unity Catalog işlev referansını .user_defined_operators.yaml dosyanıza ekleyin:

operators:
  - catalog: main
    schema: my_schema
    functionName: double_value

6. Adım: Lakeflow Designer'da operatörü kullanma

Lakeflow Designer'ı açın ve işlecin işleç paletinde göründüğünü doğrulayın. Tuvale sürükleyin, bir giriş bağlayın ve bir önizleme çalıştırın.

Troubleshooting

Issue Çözüm
Operatör Lakeflow Designer'da görünmüyor. Mevcut olup olmadığını .user_defined_operators.yaml denetleyin ve işlevinizi veya dosya yolunuzu listeler. İşleçler için python-run-function dosya yolunu ve YAML dosyasının erişilebilir olduğunu doğrulayın.
Şema doğrulaması başarısız oluyor. YAML'nizi konumundaki https://your-workspace.cloud.databricks.com/static/schemas/user-defined-operator-v0.1.0.jsonresmi şemaya göre doğrulayın.
İzin reddedildi. UC tabanlı operatörler için, kullanıcıların fonksiyon üzerinde EXECUTE ve şema üzerinde USE SCHEMA iznine sahip olduğunu doğrulayın. python-run-function işleçleri için, kullanıcıların YAML dosyasına okuma erişimine sahip olduğunu doğrulayın.
python-run-function işleci çalışma zamanında başarısız oluyor. İşlev imzasının run() ile eşleşir def run(config, inputs, spark)olup olmadığını denetleyin. Koddaki bağlantı noktası adlarının YAML ile eşleştiklerini ve dönüş sözlük anahtarlarının çıkış bağlantı noktası name değerleriyle eşleştiklerini doğrulayın.
UDTF yanlış türler döndürüyor. UDTF dönüş türleri açık olmalıdır; giriş sütun türlerine başvuramazsınız.

Permissions

İzin Purpose
. İşleci bulun.
YAML dosyasına okuma izni (python-run-functionyalnızca). İşleç tanımını yükleyin.
Unity Kataloğu işlevinde EXECUTE (yalnızca UC tabanlı işleçler). Operatörü çalıştırın.
USE SCHEMA şemada (yalnızca UC tabanlı işleçler). İşlevin oluşturulduğu şemaya erişin.
Diğer izinler Operatörünüze bağlı olarak, kullanıcılar başka izinler gerektirebilir. Örneğin, USE CONNECTION HTTP API çağrıları için Unity Kataloğu bağlantısında.

Ek kaynaklar

Aşağıdaki eğitimleri inceleyin:

Example Türü Description
Gmail e-posta göndereni python-run-function DataFrame verilerini Gmail aracılığıyla CSV e-posta eki olarak gönderin.
Bileşik faiz hesaplayıcısı uc-udf Bileşik faiz formülünü kullanarak gelecekteki yatırım değerlerini hesaplayın.
K ortalamaları kümeleme uc-udtf scikit-learn kullanarak verileri kümeler halinde segmentlere ayırma.
Slack iletisi gönderme uc-udf API aracılığıyla Slack kanallarına bildirim gönderin.
Tüm kullanıcı arabirimi pencere öğeleri uc-udf Mevcut tüm kullanıcı arabirimi pencere öğelerini gösteren referans işleci.

YAML şemasına tam bir başvuru için bkz. Kullanıcı tanımlı işleç YAML başvurusu.