Spuštění a orchestrace poznámkového bloku NotebookUtils

Pomocí nástrojů poznámkového bloku můžete spustit poznámkový blok, spustit několik poznámkových bloků paralelně nebo ukončit poznámkový blok s hodnotou. Spuštěním následujícího příkazu získejte přehled dostupných metod:

notebookutils.notebook.help()

Následující tabulka uvádí dostupné metody spuštění a orchestrace poznámkového bloku:

Metoda Signature Description
run run(path: str, timeout_seconds: int = 90, arguments: dict = None, workspace: str = ""): str Spustí poznámkový blok a vrátí jeho výstupní hodnotu.
runMultiple runMultiple(dag: Any, config: dict = None): dict[str, dict[str, Any]] Spouští více poznámkových bloků současně s podporou relací závislostí.
validateDAG validateDAG(dag: Any): bool Ověří, jestli je definice DAG správně strukturovaná.
exit exit(value: str): None Ukončí aktuální poznámkový blok s hodnotou.

Operace CRUD poznámkového bloku (vytvoření, získání, aktualizace, odstranění, výpis) najdete v tématu Správa artefaktů poznámkového bloku.

Poznámka:

Parametr config in runMultiple() je k dispozici pouze v Pythonu. Scala a R tento parametr nepodporují.

Poznámka:

Nástroje poznámkového bloku se nevztahují na definice úloh Apache Sparku (SJD).

Odkaz na poznámkový blok

Metoda run() odkazuje na poznámkový blok a vrátí jeho výstupní hodnotu. Vnořená volání funkcí můžete v poznámkovém bloku spouštět interaktivně nebo v pipeline. Odkazovaný poznámkový blok běží ve fondu Sparku toho poznámkového bloku, který tuto funkci volá.

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

Například:

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

Návratová hodnota

Metoda run() vrátí přesný řetězec předaný do notebookutils.notebook.exit(value) v podřízeném poznámkovém bloku. Pokud není voláno exit() v podřízeném poznámkovém bloku, vrátí se prázdný řetězec ("").

Poznámkové bloky Fabric také podporují odkazování na poznámkové bloky napříč pracovními prostory zadáním ID pracovního prostoru.

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

Otevřete odkaz na snímek ve výstupu z buňky a zkontrolujte referenční spuštění. Snímek zaznamenává výsledky spuštění a pomáhá ladit odkazovaný poznámkový blok.

Snímek obrazovky s výsledkem referenčního spuštění

Snímek obrazovky ukázkového snímku.

Nastavte podřízené poznámkové bloky pro příjem parametrů

Když vytvoříte podřízený poznámkový blok, který je volán prostřednictvím run() nebo runMultiple(), nastavte buňku parametru tak, aby poznámkový blok mohl přijímat argumenty od nadřazeného poznámkového bloku:

  1. Vytvořte buňku kódu s výchozími hodnotami parametrů.
  2. Označte buňku jako parametrickou buňku tak, že v uživatelském rozhraní poznámkového bloku vyberete Označit buňku jako parametry.
  3. Během provádění se hodnoty buňky parametru nahradí argumenty předanými z nadřazeného objektu.
# This cell should be marked as "parameters" cell
# Default values are overridden when the notebook is called
date = "2024-01-01"
region = "US"

Návod

Výstupní hodnoty jsou vždy řetězce. Pokud potřebujete číselnou hodnotu v nadřazeném poznámkovém bloku, převeďte výsledek po načtení (například int(result)).

Úvahy

  • Poznámkový blok pro referencování mezi pracovními prostory je podporován verzí runtime 1.2 a vyšší.
  • Pokud používáte soubory v části Zdroj poznámkového bloku, použijte notebookutils.nbResPath v odkazovaném poznámkovém bloku, aby odkazoval na stejnou složku jako interaktivní spuštění.
  • Spuštění odkazu umožňuje spouštění podřízených poznámkových bloků pouze v případě, že používají stejný objekt lakehouse jako nadřazený objekt, dědí nadřazený objekt lakehouse nebo ani jedno nedefinuje. Spuštění se zablokuje, pokud podřízený notebook určuje jiný objekt lakehouse než rodičovský poznámkový blok. Pokud chcete tuto kontrolu obejít, nastavte useRootDefaultLakehouse: True je v argumentech.
  • Nevolejte notebookutils.notebook.exit(value) uvnitř try-catch bloku. Pokud je volání ukončení zabalené ve zpracování výjimek, neúčinkuje.

Paralelní spouštění více referenčních poznámkových bloků

Použijte notebookutils.notebook.runMultiple() ke spuštění více poznámkových bloků paralelně nebo v předdefinované topologické struktuře. Rozhraní API používá vícevláknovou implementaci v rámci relace Sparku, což znamená, že označované notebooky sdílejí výpočetní prostředky.

Pomocí notebookutils.notebook.runMultiple():

  • Spusťte několik poznámkových bloků současně, aniž byste museli čekat na dokončení každého z nich.

  • Pomocí jednoduchého formátu JSON určete závislosti a pořadí provádění poznámkových bloků.

  • Optimalizujte využití výpočetních prostředků Sparku a snižte náklady na vaše projekty Fabric.

  • Prohlédněte si snímky z každého záznamu běhu poznámkového bloku ve výstupu a pohodlně laďte a monitorujte úlohy poznámkového bloku.

  • Získejte výstupní hodnotu jednotlivých aktivit vedení a použijte je v podřízených úkolech.

Spusťte notebookutils.notebook.help("runMultiple") pro zobrazení dalších příkladů a podrobností o využití.

Proveďte jednoduchý seznam poznámkových bloků

Následující příklad spustí seznam poznámkových bloků paralelně:

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

Výsledek spuštění z kořenového poznámkového bloku je následující:

Snímek obrazovky zachycující seznam poznámkových bloků

Návratová hodnota

Metoda runMultiple() vrátí slovník, kde každý klíč je název aktivity a každá hodnota je slovník s následujícími klíči:

  • exitVal: Řetězec vrácený voláním funkce exit() v podřízeném poznámkovém bloku, nebo prázdný řetězec, pokud nebylo voláno exit().
  • exception: Objekt chyby, pokud aktivita selhala nebo None pokud byla úspěšná.

Spouštění poznámkových bloků se strukturou DAG

Následující příklad spouští poznámkové bloky ve struktuře DAG pomocí .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ýsledek spuštění z kořenového poznámkového bloku je následující:

Snímek obrazovky se zobrazením seznamu poznámkových bloků s parametry.

Referenční informace k parametrům DAG

Následující tabulka popisuje každé pole, které můžete použít v definici DAG:

Obor Úroveň Povinné Description
activities Kořen Ano Seznam objektů aktivit, které určují poznámkové bloky ke spuštění.
timeoutInSeconds Kořen Ne Maximální časový limit pro celý DAG. Výchozí hodnota je 43200 (12 hodin).
concurrency Kořen Ne Maximální počet notebooků, které lze spustit souběžně. Výchozí hodnota je 3krát dostupný počet jader procesoru. Tuto hodnotu nastavte explicitně, pokud potřebujete přísnější kontrolu nebo ji použijete 0 pro neomezenou souběžnost.
name Activity Ano Jedinečný název aktivity. Slouží k identifikaci výsledků a definování závislostí.
path Activity Ano Název položky poznámkového bloku nebo cesta ke spuštění.
timeoutPerCellInSeconds Activity Ne Maximální časový limit pro každou buňku v podřízeného poznámkovém bloku Výchozí hodnota je 90 sekund.
args Activity Ne Slovník parametrů, který se má předat do podřízeného poznámkového bloku.
workspace Activity Ne Název nebo ID pracovního prostoru, ve kterém se poznámkový blok nachází. Ve výchozím nastavení se podřízený poznámkový blok spouští ve stejném pracovním prostoru jako volající.
retry Activity Ne Počet pokusů o opakování v případě selhání aktivity Výchozí hodnota je 0.
retryIntervalInSeconds Activity Ne Doba čekání v sekundách mezi opakovanými pokusy. Výchozí hodnota je 0.
dependencies Activity Ne Seznam názvů aktivit, které se musí dokončit před zahájením této aktivity.

Odkazovat na výstupní hodnoty mezi aktivitami

Pomocí výrazu args můžete odkazovat na výstupní hodnotu aktivity závislosti v @activity() poli. Tento vzorec umožňuje předávat data mezi poznámkovými bloky 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)

Návod

Pomocí výrazu @activity('activity_name').exitValue()args v poli předejte výsledky z jedné aktivity do druhé v rámci DAG.

Vytvoření dynamického DAG

Struktury DAG můžete generovat programově pro scénáře, jako je zpracování ventilátorů napříč několika oddíly:

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)

Ověřit DAG

Slouží validateDAG() k ověření platnosti struktury DAG před spuštěním. Zachytává problémy, jako jsou duplicitní názvy aktivit, chybějící závislosti a cyklický odkaz.

notebookutils.notebook.validateDAG(DAG)

Návratová hodnota

Metoda validateDAG() vrátí True , pokud je struktura DAG platná nebo vyvolá výjimku, pokud ověření selže.

Návod

V produkčních pracovních postupech vždy volejte validateDAG() před runMultiple(), abyste včas zachytili strukturální chyby.

Řešit selhání runMultiple

Metoda runMultiple() vrátí slovník, kde každý klíč je název aktivity a každá hodnota obsahuje exitVal (řetězec) a exception (objekt chyby nebo None). Částečné výsledky můžete zkontrolovat i v případě, že některé aktivity selžou:

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']}")

Úvahy

  • Stupeň paralelismu při spuštění více poznámkových bloků je omezen celkovými dostupnými výpočetními prostředky relace Sparku.
  • Výchozí počet souběžných notebooků je 3x dostupný počet jader CPU. Tuto hodnotu můžete přizpůsobit, ale nadměrný paralelismus může vést k problémům se stabilitou a výkonem kvůli vysokému využití výpočetních prostředků. Pokud dojde k problémům, zvažte rozdělení notebooků na několik volání runMultiple nebo snížení úrovně souběžnosti úpravou pole souběžnosti v parametru DAG.
  • Výchozí časový limit pro celý DAG je 12 hodin a výchozí časový limit pro každou buňku v podřízeném poznámkovém bloku je 90 sekund. Časový limit můžete změnit nastavením polí timeoutInSeconds a timeoutPerCellInSeconds v parametru DAG.
  • Konfigurace retry a retryIntervalInSeconds pro aktivity, které můžou selhat kvůli přechodným problémům, jako jsou vypršení časového limitu sítě nebo nedostupnost dočasné služby
  • Paralelní poznámkové bloky sdílejí výpočetní prostředky v rámci jedné relace Spark. Monitorujte využití prostředků, abyste se vyhnuli kolizím zatížení paměti a procesoru.

Opustit poznámkový blok

Metoda exit() ukončí poznámkový blok s hodnotou. Vnořená volání funkcí můžete v poznámkovém bloku spouštět interaktivně nebo v pipeline.

  • Při interaktivním volání exit() funkce z notebooku vyvolá Fabric notebook výjimku, přeskočí spuštění následujících buněk a udrží Spark relaci aktivní.

  • Když orchestrujete notebook v pipeline, která volá exit() funkci, aktivita notebooku vrátí výstupní hodnotu. Tím se dokončí spuštění kanálu a zastaví se relace Sparku.

  • Při volání funkce exit() v poznámkovém bloku, který je právě odkazován, Fabric Spark zastaví další spuštění tohoto odkazovaného poznámkového bloku a bude pokračovat ve spouštění dalších buněk v hlavním poznámkovém bloku, který volá funkci run(). Například: Notebook1 obsahuje tři buňky a volá exit() funkci ve druhé buňce. Poznámkový blok 2 má pět buněk a vyvolává run(notebook1) ve třetí buňce. Při spuštění Notebook2 se Notebook1 zastaví na druhé buňce, když dosáhne funkce exit(). Notebook2 nadále spouští čtvrtou a pátou buňku.

notebookutils.notebook.exit("value string")

Návratové chování

Metoda exit() nevrací hodnotu. Ukončí aktuální poznámkový blok a předá zadaný řetězec volajícímu poznámkovému bloku nebo kanálu.

Poznámka:

Funkce exit() přepíše aktuální výstup buňky. Pokud se chcete vyhnout ztrátě výstupu jiných příkazů kódu, zavolejte notebookutils.notebook.exit() v samostatné buňce.

Důležité

Nevolejte notebookutils.notebook.exit() uvnitř try-catch bloku. Ukončení nebude mít účinek, pokud je zabaleno do zpracování výjimek. Aby exit() volání fungovalo správně, musí být na nejvyšší úrovni kódu.

Například:

Poznámkový blok Sample1 obsahuje následující dvě buňky:

  • Buňka 1 definuje vstupní parametr s výchozí hodnotou nastavenou na 10.

  • Buňka 2 opouští poznámkový blok s input jako výstupní hodnotou.

Snímek obrazovky s ukázkovým poznámkovým blokem výstupní funkce

Ukázku 1 můžete spustit v jiném poznámkovém bloku s výchozími hodnotami:

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

Výstup:

10

Ukázku 1 můžete spustit v jiném poznámkovém bloku a nastavit vstupní hodnotu na 20:

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

Výstup:

20