Felhasználó által definiált operátorok a Lakeflow Designerben

A Lakeflow Designer lehetővé teszi a felhasználó által definiált operátorok létrehozását, amelyek közvetlenül a vásznon jelennek meg a beépített operátorok mellett. Ezekkel kibővítheti a Lakeflow Designert saját üzleti logikájával, számításaival vagy integrációival.

A felhasználó által definiált operátoroknak három típusa van:

  • python-run-function: A munkaterületen tárolt, beágyazott Python-kódot tartalmazó önálló YAML-fájl. A Legjobb a DataFrame-szintű átalakításokhoz és külső integrációkhoz. Az engedélyek kezelése a munkaterület fájlszintjén lehetséges.
  • uc-udf: Körbefuttat egy Unity Catalog skaláris függvényt. Oszlopszintű átalakításokhoz a legjobb. A hozzáférésre a Unity Katalógus engedélyei vonatkoznak.
  • uc-udtf: Egy Unity Catalog táblaértékű függvény körbefuttatása. Leginkább táblaszintű átalakításokhoz ajánlott, például gépi tanulási klaszterezéshez és összesítéshez. A hozzáférésre a Unity Katalógus engedélyei vonatkoznak.
Funkció python-run-function uc-udf uc-udtf
Példa használati esetre DataFrame-átalakítások, API-integrációk, e-mail-értesítések Oszlopszintű számítások (BMI, kamatlábak) ML-klaszterezés, sorok közötti aggregálás
Bemenet Adatkeretek Önálló értékek Teljes táblázat sorról sorra
Output Adatkeretek Egyetlen érték Táblázat (több sor)
Unity Catalog-függvényt igényel No Igen Igen
Hozzáférés szabályozása Munkaterület fájlengedélyek Unity Catalog-engedélyek (EXECUTE, USE SCHEMA) Unity Catalog-engedélyek (EXECUTE, USE SCHEMA)
Támogatott nyelvek Csak Python SQL vagy Python SQL-burkolóban SQL vagy Python SQL-burkolóban

A felhasználó által definiált operátorok működése

A felhasználó által definiált operátorok a következőkből állnak:

  • Operátorlogika: Az operátor végrehajtásakor futó kód. Ez lehet beágyazott Python run() függvény (python-run-function esetén) vagy Unity Catalog függvény (uc-udf és uc-udtf esetén).
  • YAML-konfiguráció: Közli a Lakeflow Designerrel, hogyan jelenítse meg az operátort a felhasználói felületen, beleértve az operátor nevét, leírását, bemeneti paramétereit, felhasználói felületi vezérlőit és portjait. Minden operátortípus használja a sémát user-defined-operator-v0.1.0 .
  • Regisztrációs fájl: Olyan bejegyzés, amely .user_defined_operators.yaml lehetővé teszi, hogy a Lakeflow Designer felderítse az operátort.

Operátorlogika

Python felhasználó által definiált operátorlogika futtatása

Minden python-run-function operátornak definiálnia kell egy függvényt run() :

def run(config: Dict[str, Any], inputs: Dict[str, Any], spark) -> Dict[str, Any]:
  • config: Felhasználó által konfigurált értékek a felhasználói felületen, tulajdonságnév alapján kulcsolt.
  • inputs: Bemeneti DataFrame-ek, bemeneti port szerint kulcsolva name.
  • spark: Az aktív SparkSession.
  • Visszaadja: Egy szótár a kimeneti port name értékeit a DataFrame-ekre megfelelteti.

Az alábbi példa egy bemeneti DataFrame-ből szűri a sorokat:

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

Ha az operátor külső pipcsomagokat igényel, adja hozzá a environment mezőt a YAML-hez:

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

UDF és UDTF operátorlogika

Az UC-függvényeket SQL-ben vagy Python is megírhatja. Python függvények egy SQL-CREATE FUNCTION utasításba vannak csomagolva:

SQL-függvény:

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 függvény (SQL-be burkolva):

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)
$$;

Az UDF-ek egyszerre egyetlen értéket dolgoznak fel, és kiszámított értéket adnak vissza. Az UDTF-ek sorról sorra dolgozzák fel a táblákat, és az állapotot az összes sorban meg tudják tartani. Használja a uc-udf elemet oszlopszintű átalakításokhoz, a uc-udtf elemet pedig olyan műveletekhez, mint az ML-fürtözés vagy az aggregálás.

Az UDTF-ek emellett három fő módszert is meg kell határozniuk: __init__(), eval()és 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

Az UDTF visszatérési tábláknak rögzített, explicit típusokkal kell rendelkezniük. A visszatérési konfigurációban nem hivatkozhat bemeneti oszloptípusokra.

YAML-konfiguráció

A YAML-konfiguráció tájékoztatja a Lakeflow Designert, hogyan jelenítse meg az operátort a felhasználói felületen. Meghatározza az operátor nevét, leírását, bemeneti paramétereit, felhasználói felületi vezérlőit és portjait. Minden konfigurációs mező egy típussal, címmel és opcionális x-ui vezérlőtippekkel rendelkező tulajdonság:

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

A YAML-sémával kapcsolatos részletes információkért, beleértve az összes widgettípust és konfigurációs beállítást, tekintse meg a felhasználó által definiált operátor YAML-referenciáját.

Kikötők

A portok határozzák meg az operátor bemeneteit és kimeneteit:

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

YAML Python függvényoperátorokhoz

A python-run-function operátorok esetében a YAML-fájl önálló, és tartalmaz egy run_function mezőt beágyazott Python kóddal:

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}

YAML a Unity Catalog-függvényekhez

UC-alapú operátorok esetén ágyazza be a YAML-konfigurációt megjegyzésként vagy docstringként a függvénybe.

SQL-ben (a /* ... */ megjegyzés használatával):

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

A Python (""" ... """ dokumentum használata):

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)
$$;

Az operátor regisztrálása és üzembe helyezése a Lakeflow Designerben

Ahhoz, hogy az operátor megjelenjen a Lakeflow Designerben, regisztrálja azt egy .user_defined_operators.yaml fájlban:

  • Munkaterület szintje: Helyezze a fájlt a munkaterület gyökerére, hogy az operátor látható legyen az összes felhasználó számára.
  • Felhasználói szint: Helyezze a fájlt a felhasználói kezdőlap mappájába (/Workspace/Users/<user-name>/.user_defined_operators.yaml), hogy az operátorok csak Ön számára legyenek láthatók.

A operators: szakasz támogatja a fájlelérési utakat, a Unity Catalog függvényhivatkozásait és a glob mintákat. A bejegyzéstípusok keverhetők:

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

Operátor frissítése vagy eltávolítása

Egy operátor kódjának módosításakor frissítse a felhasználó által definiált operátorokat a módosítás betöltéséhez. A menü Operátorok lapján kattintson a Frissítés ikonra.

  • Ha az operátor ugyanaz marad version, a frissítés betölti a frissített kódot.
  • Ha az operátornak van egy új version verziója, a vásznon található operátor a frissítés után felajánlja, hogy frissít rá (vagy megtartja az aktuális verziót).

Ha el szeretne távolítani egy operátort a Lakeflow Designerből, törölje a bejegyzését a helyről .user_defined_operators.yaml. A uc-udf és uc-udtf operátorok esetén a mögöttes Unity Catalog-függvényt a DROP FUNCTION paranccsal is eltávolíthatja, ha már nincs rá szüksége.

Speciális konfigurációk

Előnézeti mód

A Lakeflow Designer tervezési módban támogatja az előzetes verziókat. Külső API-kat hívó vagy külső rendszerekbe író operátorok esetén adjon hozzá egy is_preview konfigurációs tulajdonságot, hogy kihagyhassa a mellékhatásokat az előzetes verzióban. Ha az előnézeti mód engedélyezve van, a felhasználóknak kifejezetten a Futtatás gombra kell kattintanak az operátor mellékhatásokkal való végrehajtásához.

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

A Lakeflow Designer automatikusan true értékre állítja ezt az értéket az előnézet során. Ellenőrizze ezt a logikában a mellékhatások elkerüléséhez:

# 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 Catalog-kapcsolatok

Külső API-kat hívó UC-alapú SQL-operátorok esetén a Unity Catalog HTTP-kapcsolataival biztonságosan tárolhatja a hitelesítő adatokat:

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

Ezután használja a kapcsolatot az SQL UDF-ben a http_request() függvénnyel. További információ: Csatlakozás külső HTTP-szolgáltatásokhoz.

WorkspaceClient

A python-run-function operátorok esetében a Azure Databricks WorkspaceClient használatával érheti el a munkaterület erőforrásait és külső API-kat:

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

Teljes, felhasználó által definiált python-run-function operátor létrehozása

Az alábbi lépések végigvezetik Önt egy python-run-function operátor a semmiből történő létrehozásán.

1. lépés: A logika meghatározása

Írja meg a(z) run() függvényt egy jegyzetfüzetben:

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. lépés: A függvény tesztelése

A függvény interaktív tesztelése mintaadatokkal:

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. lépés: A YAML-konfiguráció létrehozása

Adja meg az operátor metaadatait, konfigurációs mezőit és portját egy YAML-fájlban:

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. lépés: A logika és a YAML kombinálása

Adja hozzá a run_function és a ports mezőket a teljes YAML-fájl létrehozásához. Mentse a munkaterületére, például /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. lépés: Az operátor regisztrálása

Adja hozzá a fájl elérési útját:.user_defined_operators.yaml

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

6. lépés: Az operátor használata a Lakeflow Designerben

Nyissa meg a Lakeflow Designert, és ellenőrizze, hogy az operátor megjelenik-e az operátorpalettán. Húzza a vászonra, csatlakoztassa a bemenetet, konfigurálja az oszlop nevét, és futtasson egy előnézetet.

Teljes felhasználó által definiált UC-operátor létrehozása

Az alábbi lépések végigvezetik az UC-alapú uc-udf operátorok létrehozásának lépésein.

1. lépés: A logika meghatározása

A függvénylogika írása és tesztelése jegyzetfüzetben:

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

2. lépés: A YAML-konfiguráció létrehozása

Adja meg az operátor metaadatait, a konfigurációs mezőket és a portokat:

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. lépés: A logika és a YAML kombinálása

Hozd létre a Unity Catalog-függvényt úgy, hogy a YAML docstringként legyen beágyazva. A példa létrehozza a függvényt -ben main.example_output. Először hozd létre a sémát, ha nem létezik:

CREATE SCHEMA IF NOT EXISTS main.example_output
CREATE OR REPLACE FUNCTION main.example_output.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. lépés: A függvény tesztelése

Tesztelje a függvényt. Visszaadja a 10 értéket.

SELECT main.example_output.double_value(5) AS result;

5. lépés: Az operátor regisztrálása

Adja hozzá a Unity Catalog függvényhivatkozást a .user_defined_operators.yaml fájlhoz:

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

6. lépés: Az operátor használata a Lakeflow Designerben

Nyissa meg a Lakeflow Designert, és ellenőrizze, hogy az operátor megjelenik-e az operátorpalettán. Húzza a vászonra, csatlakoztassa a bemenetet, és futtasson egy előnézetet.

Hibaelhárítás

Issue Megoldás
Az operátor nem jelenik meg a Lakeflow Designerben. Ellenőrizze, hogy a(z) .user_defined_operators.yaml létezik-e, és felsorolja-e a függvényt vagy a fájl elérési útját. Operátorok esetén python-run-function ellenőrizze a fájl elérési útját, és hogy a YAML-fájl elérhető-e.
A séma érvényesítése sikertelen. Ellenőrizze a YAML-t a hivatalos séma alapján itt: https://your-workspace.cloud.databricks.com/static/schemas/user-defined-operator-v0.1.0.json.
Engedély megtagadva. UC-alapú operátorok esetén ellenőrizze, hogy a felhasználók rendelkeznek-e EXECUTE jogosultsággal a függvényre és USE SCHEMA jogosultsággal a sémára. Operátorok esetén python-run-function ellenőrizze, hogy a felhasználók olvasási hozzáféréssel rendelkeznek-e a YAML-fájlhoz.
python-run-function operátor futásidőben hibát jelez. Ellenőrizze, hogy a run() függvény aláírása megegyezik-e a(z) def run(config, inputs, spark) értékkel. Ellenőrizze, hogy a kódban szereplő portnevek megfelelnek-e a YAML-nek, és hogy a visszaadott szótárkulcsok megfelelnek-e a kimeneti port name értékeinek.
Az UDTF helytelen típusokat ad vissza. Az UDTF visszatérési típusainak explicitnek kell lenniük; a bemeneti oszloptípusokra nem hivatkozhat.

Permissions

Engedély Alkalmazás célja
Olvasási hozzáférés a következőhöz .user_defined_operators.yaml: . Fedezze fel az operátort.
Olvasási hozzáférés a YAML-fájlhoz (python-run-function csak). Töltse be az operátordefiníciót.
EXECUTE a Unity Catalog függvényen (csak UC-alapú operátorok). Indítsa el az operátort.
USE SCHEMA a sémán (csak UC-alapú operátorok esetén). Hozzáférés a függvény létrehozásához szükséges sémához.
Egyéb engedélyek Az operátortól függően előfordulhat, hogy a felhasználók más engedélyeket igényelnek. Például USE CONNECTION egy Unity Catalog-kapcsolaton HTTP API-hívásokhoz.

További erőforrások

Ismerkedjen meg a következő oktatóanyagokkal:

Example Típus Description
Gmail-üzenet feladója python-run-function DataFrame-adatok küldése CSV e-mail mellékletként a Gmailen keresztül.
Kamatkalkulátor uc-udf A jövőbeli befektetési értékek kiszámítása az összetett kamat képlet használatával.
K-közép klaszterezés uc-udtf Adatok csoportosítása fürtökbe a scikit-learn használatával.
Slack-üzenet küldése uc-udf Értesítések küldése a Slack-csatornákra API-val.
Minden felhasználói felületi widget uc-udf Referencia operátor, amely az összes elérhető felhasználói felületi widgetet megjeleníti.

A YAML-sémára való teljes hivatkozásért tekintse meg a felhasználó által definiált OPERÁTOR YAML-referenciáját.