Notebook run a orchestrácia notebooku NotebookUtils

Použite nástroje zápisníka na spustenie zápisníka, spustenie viacerých zápisníkov paralelne alebo ukončenie zápisníka s hodnotou. Spustením nasledujúceho príkazu získate prehľad dostupných metód:

notebookutils.notebook.help()

Nasledujúca tabuľka uvádza dostupné metódy behu a orchestrácie zápisníka:

Method Podpis Description
run run(path: str, timeout_seconds: int = 90, arguments: dict = None, workspace: str = ""): str Spustí zápisník a vráti jeho výstupnú hodnotu.
runMultiple runMultiple(dag: Any, config: dict = None): dict[str, dict[str, Any]] Spúšťa viacero zápisníkov súčasne s podporou závislostných vzťahov.
validateDAG validateDAG(dag: Any): bool Overuje, či je DAG definícia správne štruktúrovaná.
exit exit(value: str): None Ukončí aktuálny zápisník s hodnotou.

Pre CRUD operácie zápisníka (vytváranie, získavanie, aktualizácia, mazanie, zoznam) pozri Spravovať artefakty zápisníka.

Poznámka

Parameter config v runMultiple() je dostupný iba v Pythone. Scala a R tento parameter nepodporujú.

Poznámka

Pomôcky poznámkového blokov nie sú použiteľné pre definície úloh Apache Spark (SJD).

Odkaz na poznámkový blok

Metóda run() odkazuje na zápisník a vracia jeho výstupnú hodnotu. Volania vnorených funkcií môžete spustiť interaktívne v poznámkovom bloke alebo v kanáli. Poznámkový blok, na ktorý sa odkazuje, je spustený vo fonde spark poznámkového bloku, ktorý túto funkciu volá.

notebookutils.notebook.run("notebook name", <timeout_seconds>, <arguments>, <workspace>)

Napríklad:

notebookutils.notebook.run("Sample1", 90, {"input": 20 })

Vrátená hodnota

Metóda run() vráti presný reťazec, ktorý bol odovzdaný do notebookutils.notebook.exit(value) detského zápisníka. Ak exit() sa v detskom zápisníku nezavolá, vráti sa prázdny reťazec ("").

Textilné zápisníky tiež podporujú referencovanie zápisníkov naprieč pracovnými priestormi zadaním ID pracovného priestoru.

notebookutils.notebook.run("Sample1", 90, {"input": 20 }, "fe0a6e2a-a909-4aa3-a698-0a651de790aa")

Otvorte snapshot link vo výstupe bunky, aby ste skontrolovali referenčný beh. Snapshot zachytáva výsledky behu a pomáha vám ladiť odkazovaný zápisník.

Snímka obrazovky s výsledkom referenčného spustenia.

Snímka obrazovky s príkladom snímky.

Nastavte detské zápisníky na prijímanie parametrov

Keď vytvoríte detský zápisník, ktorý sa volá cez run() alebo runMultiple(), nastavte parametrickú bunku tak, aby zápisník mohol prijímať argumenty od rodiča:

  1. Vytvorte kódovú bunku s predvolenými hodnotami parametrov.
  2. Označte bunku ako parameter cell výberom Mark cell ako parameter v rozhraní notebooku.
  3. Počas vykonávania sú hodnoty buniek parametrov nahradené argumentmi odovzdanými od rodiča.
# This cell should be marked as "parameters" cell
# Default values are overridden when the notebook is called
date = "2024-01-01"
region = "US"

Prepitné

Výstupné hodnoty sú vždy reťazce. Ak potrebujete číselnú hodnotu v rodičovskom zápisníku, preveďte výsledok po získaní (napríklad int(result)).

Zváženia

  • Referenčný poznámkový blok krížového pracovného priestoru je podporovaný verziou 1.2 a novšou verziou modulu runtime.
  • Ak používate súbory v časti Zdroj poznámkového bloku, použite notebookutils.nbResPath odkazovaný poznámkový blok, aby ste sa uistili, že odkazuje na rovnaký priečinok ako interaktívne spustenie.
  • Spustenie odkazu umožňuje spustenie podriadených poznámkových blokov iba v prípade, že používajú rovnaký jazerný dom ako nadradený objekt, dedia nadradený jazerársky dom alebo ho nedefinuje. Vykonávanie je zablokované, ak dieťa určí iný lakehouse než rodičovský zápisník. Aby ste túto kontrolu obišli, useRootDefaultLakehouse: True nastavte argumenty.
  • Nevolajte notebookutils.notebook.exit(value) v bloku try-catch . Výstupné volanie nenadobudne účinok, keď je zabalené v spracovaní výnimiek.

Odkaz na paralelné spustenie viacerých poznámkových blokov

Použite notebookutils.notebook.runMultiple() na spúšťanie viacerých notebookov paralelne alebo v preddefinovanej topologickej štruktúre. API používa viacvláknovú implementáciu v rámci relácie Spark, čo znamená, že referencované behy zápisníka zdieľajú výpočtové zdroje.

S možnosťou notebookutils.notebook.runMultiple()môžete:

  • Vykonajte viacero poznámkových blokov súčasne bez čakania na dokončenie každého z nich.

  • Zadajte závislosti a poradie vykonávania pre poznámkové bloky pomocou jednoduchého formátu JSON.

  • Optimalizujte použitie výpočtových zdrojov služby Spark a znížte náklady na projekty služby Fabric.

  • Zobrazte snímky každého záznamu spúšťania poznámkového bloku vo výstupe a pohodlné ladenie/monitorovanie úloh poznámkového bloku.

  • Získajte hodnotu výstupu každej aktivity vedúceho pracovníka a použite ich v následných úlohách.

Bežte a pozrite notebookutils.notebook.help("runMultiple") si viac príkladov a detailov používania.

Spusti jednoduchý zoznam zápisníkov

Nasledujúci príklad zobrazuje zoznam zápisníkov paralelne:

notebookutils.notebook.runMultiple(["NotebookSimple", "NotebookSimple2"])

Výsledok spustenia koreňového poznámkového bloku je nasledovný:

Snímka obrazovky odkazovania na zoznam poznámkových blokov.

Vrátená hodnota

Metóda runMultiple() vracia slovník, kde každý kľúč je názov aktivity a každá hodnota je slovník s nasledujúcimi kľúčmi:

  • exitVal: Reťazec vrátený volaním detského zošita exit() , alebo prázdny reťazec, ak exit() nebol zavolaný.
  • exception: Chybový objekt, ak aktivita zlyhala, alebo None ak bola úspešná.

Spúšťajte notebooky so štruktúrou DAG

Nasledujúci príklad spúšťa zápisníky v DAG štruktúre pomocou .notebookutils.notebook.runMultiple()

# run multiple notebooks with parameters
DAG = {
    "activities": [
        {
            "name": "Process_1", # activity name, must be unique
            "path": "NotebookSimple", # notebook item name
            "timeoutPerCellInSeconds": 90, # max timeout for each cell, default to 90 seconds
            "args": {"p1": "changed value", "p2": 100}, # notebook parameters
            "workspace":"WorkspaceName" # both name and id are supported
        },
        {
            "name": "Process_2",
            "path": "NotebookSimple2",
            "timeoutPerCellInSeconds": 120,
            "args": {"p1": "changed value 2", "p2": 200},
            "workspace":"id" # both name and id are supported
        },
        {
            "name": "Process_1.1",
            "path": "NotebookSimple2",
            "timeoutPerCellInSeconds": 120,
            "args": {"p1": "changed value 3", "p2": 300},
            "retry": 1,
            "retryIntervalInSeconds": 10,
            "dependencies": ["Process_1"] # list of activity names that this activity depends on
        }
    ],
    "timeoutInSeconds": 43200, # max timeout for the entire DAG, default to 12 hours
    "concurrency": 12 # max number of notebooks to run concurrently, default to 3x CPU cores, 0 means unlimited
}
notebookutils.notebook.runMultiple(DAG, {"displayDAGViaGraphviz": False})

Výsledok spustenia koreňového poznámkového bloku je nasledovný:

Snímka obrazovky znázorňujúca odkaz na zoznam poznámkových blokov s parametrami.

Referencia parametra DAG

Nasledujúca tabuľka popisuje každé pole, ktoré môžete použiť v DAG definícii:

Pole Úroveň Požaduje sa Description
activities Koreň Áno Zoznam objektov aktivít, ktoré definujú, ktoré zápisníky majú bežať.
timeoutInSeconds Koreň Nie Maximálny časový limit pre celý DAG. Predvolená hodnota je 43200 (12 hodín).
concurrency Koreň Nie Maximálny počet zápisníkov, ktoré musia bežať súčasne. Predvolený je 3-násobok dostupného počtu jadier CPU. Nastavte túto hodnotu explicitne, ak potrebujete presnejšiu kontrolu, alebo použite 0 na neobmedzenú súbežnosť.
name Aktivita Áno Jedinečný názov pre túto aktivitu. Používa sa na identifikáciu výsledkov a definovanie závislostí.
path Aktivita Áno Názov položky zápisníka alebo cesta na vykonanie.
timeoutPerCellInSeconds Aktivita Nie Maximálny časový limit pre každú bunku v detskom zošite. Predvolená je 90 sekúnd.
args Aktivita Nie Slovník parametrov, ktoré sa posielajú do detského zápisníka.
workspace Aktivita Nie Názov pracovného priestoru alebo ID, kde sa zápisník nachádza. Predvolene detský zápisník beží v rovnakom pracovnom priestore ako volajúci.
retry Aktivita Nie Počet pokusov o opätovné pokusy, ak aktivita zlyhá. Predvolená hodnota je 0.
retryIntervalInSeconds Aktivita Nie Čas čakania medzi opakovanými pokusmi je v sekundách. Predvolená hodnota je 0.
dependencies Aktivita Nie Zoznam názvov aktivít, ktoré musia byť dokončené pred začiatkom tejto aktivity.

Referenčné výstupné hodnoty medzi aktivitami

Výstupnú hodnotu závislosti môžete v args poli použiť pomocou výrazu @activity() . Tento vzor vám umožňuje prenášať dáta medzi zápisníkmi v DAG.

DAG = {
    "activities": [
        {
            "name": "Extract",
            "path": "ExtractData",
            "timeoutPerCellInSeconds": 120,
            "args": {"source": "prod_db"}
        },
        {
            "name": "Transform",
            "path": "TransformData",
            "timeoutPerCellInSeconds": 180,
            "args": {
                "data_path": "@activity('Extract').exitValue()"
            },
            "dependencies": ["Extract"]
        }
    ]
}

results = notebookutils.notebook.runMultiple(DAG)

Prepitné

Použite @activity('activity_name').exitValue() výraz v args poli na prenos výsledkov z jednej aktivity do druhej v rámci DAG.

Postavte dynamický DAG

DAG štruktúry môžete programovo generovať pre situácie ako fan-out spracovanie cez viaceré partície:

def create_fan_out_dag(partitions):
    activities = []

    for partition in partitions:
        activities.append({
            "name": f"Process_{partition}",
            "path": "ProcessPartition",
            "timeoutPerCellInSeconds": 180,
            "args": {"partition": partition}
        })

    activities.append({
        "name": "Aggregate",
        "path": "AggregateResults",
        "timeoutPerCellInSeconds": 120,
        "dependencies": [f"Process_{p}" for p in partitions]
    })

    return {"activities": activities, "concurrency": 25}

partitions = ["2024-01", "2024-02", "2024-03", "2024-04"]
dag = create_fan_out_dag(partitions)

results = notebookutils.notebook.runMultiple(dag)

Validujte DAG

Použite na validateDAG() overenie, či je vaša DAG štruktúra platná pred spustením. Zachytáva problémy ako duplicitné názvy aktivít, chýbajúce závislosti a kruhové odkazy.

notebookutils.notebook.validateDAG(DAG)

Vrátená hodnota

Metóda vráti, validateDAG()True ak je DAG štruktúra platná, alebo vytvorí výnimku, ak validácia zlyhá.

Prepitné

Vždy volajte validateDAG() vopred runMultiple() v produkčných pracovných postupoch, aby ste včas odhalili štrukturálne chyby.

Riešenie behuViacnásobné zlyhania

Metóda runMultiple() vracia slovník, kde každý kľúč je názov aktivity a každá hodnota obsahuje exitVal (reťazec) a exception (chybový objekt alebo None). Čiastočné výsledky môžete skontrolovať aj vtedy, keď niektoré aktivity zlyhajú:

from notebookutils.common.exceptions import RunMultipleFailedException

try:
    results = notebookutils.notebook.runMultiple(DAG)
except RunMultipleFailedException as ex:
    results = ex.result

for activity_name, result in results.items():
    if result["exception"]:
        print(f"{activity_name} failed: {result['exception']}")
    else:
        print(f"{activity_name} succeeded: {result['exitVal']}")

Zváženia

  • Stupeň paralelného spustenia viacerých poznámkových blokov je obmedzený na celkový dostupný výpočtový prostriedok relácie služby Spark.
  • Predvolený počet súbežných notebookov je trojnásobok dostupného počtu jadier CPU. Túto hodnotu si môžete prispôsobiť, ale nadmerný paralelizmus môže viesť k problémom so stabilitou a výkonom kvôli vysokej spotrebe výpočtových zdrojov. Ak sa vyskytnú problémy, zvážte rozdelenie poznámkových blokov do viacerých runMultiple volaní alebo zníženie súbežnosti úpravou poľa súbežnosti v parametri DAG.
  • Predvolený časový limit pre celý DAG je 12 hodín a predvolený časový limit pre každú bunku v detskom zápisníku je 90 sekúnd. Časový limit môžete zmeniť nastavením časového limituInSeconds a časového limitu políPerCellInSeconds v parametri DAG.
  • Nakonfigurujte a pre retry aktivity retryIntervalInSeconds , ktoré by mohli zlyhať kvôli prechodným problémom, ako sú časové limity siete alebo dočasná nedostupnosť služieb.
  • Paralelné notebooky zdieľajú výpočtové zdroje v rámci jednej Spark relácie. Monitorujte využitie zdrojov, aby ste sa vyhli tlaku pamäte a konkurencii CPU.

Ukončenie poznámkového bloku

Metóda exit() ukončí zápisník s hodnotou. Volania vnorených funkcií môžete spustiť interaktívne v poznámkovom bloke alebo v kanáli.

  • Keď interaktívne voláte funkciu exit() z notebooku, Fabric notebook vyhodí výnimku, preskočí spúšťanie ďalších buniek a udrží Spark session nažive.

  • Keď orchestrujete zápisník v pipeline, ktorý volá funkciu exit() , aktivita zápisníka vráti výstupnú hodnotu. Tým sa dokončí pipeline run a Spark session sa zastaví.

  • Keď zavoláte funkciu exit() v zápisníku, na ktorý sa odkazuje, Fabric Spark zastaví ďalšie vykonávanie odkazovaného zápisníka a pokračuje v spustení ďalších buniek v hlavnom zápisníku, ktoré volajú funkciu run() . Napríklad: Notebook1 má tri bunky a volá funkciu exit() v druhej bunke. Notebook2 má päť buniek a volá run(notebook1) v tretej bunke. Keď spustíte Notebook2, Notebook1 sa zastaví na druhej bunke pri kliknutí na funkciu exit() . Notebook2 naďalej spúšťa svoju štvrtú bunku a piatu bunku.

notebookutils.notebook.exit("value string")

Správanie návratu

Metóda exit() nevracia hodnotu. Ukončí aktuálny zápisník a odovzdá poskytnutý reťazec volajúcemu zápisníku alebo pipeline.

Poznámka

Funkcia exit() prepíše aktuálny výstup bunky. Ak sa chcete vyhnúť strate výstupu iných príkazov kódu, zavolajte notebookutils.notebook.exit() v samostatnej bunke.

Dôležité

Nevolajte notebookutils.notebook.exit() v bloku try-catch . Výstup sa neprejaví, keď je zabalený v spracovaní výnimiek. Hovor exit() musí byť na najvyššej úrovni kódu, aby fungoval správne.

Napríklad:

Zápisník Sample1 má nasledujúce dve bunky:

  • Bunka 1 definuje vstupný parameter s predvolenou hodnotou nastavenou na 10.

  • Bunka 2 ukončí poznámkový blok so vstupom ako výstupnou hodnotou.

Snímka obrazovky zobrazujúca ukážkový poznámkový blok funkcie exit.

Ukážku1 môžete spustiť v inom poznámkovom bloku s predvolenými hodnotami:

exitVal = notebookutils.notebook.run("Sample1")
print (exitVal)

Output:

10

Ukážku1 môžete spustiť v inom poznámkovom bloku a nastaviť vstupnú hodnotu na 20:

exitVal = notebookutils.notebook.run("Sample1", 90, {"input": 20 })
print (exitVal)

Output:

20