NotebookUtils-jegyzetfüzet futtatása és orchestrációja

A jegyzetfüzet segédprogramjaival jegyzetfüzetet futtathat, több jegyzetfüzetet futtathat párhuzamosan, vagy kiléphet egy értékekkel rendelkező jegyzetfüzetből. Futtassa a következő parancsot az elérhető módszerek áttekintéséhez:

notebookutils.notebook.help()

Az alábbi táblázat az elérhető jegyzetfüzet-futtatási és vezénylési módszereket sorolja fel:

Módszer Signature Leírás
run run(path: str, timeout_seconds: int = 90, arguments: dict = None, workspace: str = ""): str Futtat egy jegyzetfüzetet, és visszaadja a kilépési értékét.
runMultiple runMultiple(dag: Any, config: dict = None): dict[str, dict[str, Any]] Egyszerre több notebookot futtat, támogatva a függőségi kapcsolatokat.
validateDAG validateDAG(dag: Any): bool Ellenőrzi, hogy a DAG-definíció megfelelően van-e strukturálva.
exit exit(value: str): None Kilép az aktuális jegyzetfüzetből egy értékkel.

A jegyzetfüzet CRUD-műveleteiről (létrehozás, lekérés, frissítés, törlés, lista) lásd: Jegyzetfüzet-összetevők kezelése.

Megjegyzés:

A config paraméter runMultiple() csak Pythonban érhető el. A Scala és az R nem támogatja ezt a paramétert.

Megjegyzés:

A jegyzetfüzet-segédprogramok nem alkalmazhatók az Apache Spark-feladatdefiníciókra (SJD).

Hivatkozzon egy jegyzetfüzetre

A run() metódus egy jegyzetfüzetre hivatkozik, és visszaadja a kilépési értékét. A beágyazott függvényhívásokat interaktívan vagy folyamatban is futtathatja egy jegyzetfüzetben. A hivatkozott jegyzetfüzet a függvényt meghívó jegyzetfüzet Spark-készletében fut.

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

Például:

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

Visszaadott érték

A run() metódus a gyermekjegyzetfüzetben megadott pontos sztringet notebookutils.notebook.exit(value) adja vissza. Ha a gyermekjegyzetfüzetből nem hívják meg a exit()-t, akkor egy üres sztringet ("") ad vissza.

A hálójegyzetfüzetek a munkaterület azonosítójának megadásával támogatják a jegyzetfüzetek munkaterületek közötti hivatkozását is.

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

Nyissa meg a pillanatkép-hivatkozást a cella kimenetében a referenciafuttatás vizsgálatához. A pillanatkép rögzíti a futtatott eredményeket, és segít a hivatkozott jegyzetfüzet hibakeresésében.

Képernyőkép a referenciafuttatás eredményéről.

Képernyőkép egy pillanatkép-példáról.

Gyermekjegyzetfüzetek beállítása paraméterek fogadásához

Amikor létrehoz egy aljegyzetfüzetet, amelyet a run() vagy runMultiple() hív meg, állítson be egy paramétercellát, hogy a jegyzetfüzet argumentumokat fogadhasson a szülőtől.

  1. Hozzon létre egy kódcellát alapértelmezett paraméterértékekkel.
  2. Jelölje meg a cellát paramétercellaként, ha a Jegyzetfüzet felhasználói felületén kiválasztja a Cella megjelölése paraméterként lehetőséget.
  3. A végrehajtás során a paramétercellák értékeit a rendszer a szülőtől kapott argumentumokra cseréli.
# This cell should be marked as "parameters" cell
# Default values are overridden when the notebook is called
date = "2024-01-01"
region = "US"

Jótanács

A kilépési értékek mindig karakterláncok. Ha numerikus értékre van szüksége a szülőjegyzetfüzetben, konvertálja az eredményt a beolvasás után (például int(result)).

Megfontolások

  • A munkaterületek közötti referenciajegyzetfüzetet az 1.2-es és újabb futtatókörnyezet támogatja.
  • Ha a Notebook Resource alatt lévő fájlokat használja, a notebookutils.nbResPath hivatkozott jegyzetfüzetben győződjön meg arról, hogy ugyanarra a mappára mutat, amelyet az interaktív futtatás során használnak.
  • A referenciafuttatás lehetővé teszi, hogy a gyermekjegyzetfüzetek csak akkor fussanak, ha ugyanazt a lakehouse-t használják, mint a szülő, ha öröklik a szülő lakehouse-ját, vagy ha egyik sem definiál egy lakehouse-t. A végrehajtás le lesz tiltva, ha a gyermek a szülőjegyzetfüzettől eltérő tóházat ad meg. Az ellenőrzés megkerüléséhez állítsa be useRootDefaultLakehouse: True az argumentumokat.
  • Ne hívja meg a notebookutils.notebook.exit(value)-t a try-catch blokkon belül. A kilépési hívás nem lép érvénybe, ha a kivételkezelésbe van burkolva.

Hivatkozás több jegyzetfüzet párhuzamos futtatására

Több jegyzetfüzet párhuzamos vagy előre definiált topológiai struktúrában való futtatására használható notebookutils.notebook.runMultiple() . Az API egy Spark-munkameneten belül többszálú implementációt használ, ami azt jelenti, hogy a hivatkozott jegyzetfüzetek számítási erőforrásokat osztanak meg.

A(z) notebookutils.notebook.runMultiple() segítségével a következőket teheted:

  • Több jegyzetfüzetet hajthatsz végre egyszerre, anélkül, hogy meg kellene várni mindegyik befejezését.

  • Egyszerű JSON-formátum használatával adja meg a jegyzetfüzetek függőségeit és végrehajtási sorrendjét.

  • Optimalizálja a Spark számítási erőforrások használatát, és csökkentse a Fabric-projektek költségeit.

  • A kimenetben megtekintheti az egyes jegyzetfüzet-futtatási rekordok pillanatképeit, és kényelmesen hibakeresést/monitorozást végezhet a jegyzetfüzet-feladatokban.

  • Szerezze be az egyes vezetői tevékenységek kilépési értékét, és használja őket az alsóbb rétegbeli feladatokban.

Futtassa a notebookutils.notebook.help("runMultiple") parancsot további példák és a használati részletek megtekintéséhez.

Jegyzetfüzetek egyszerű listájának futtatása

Az alábbi példa párhuzamosan futtatja a jegyzetfüzetek listáját:

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

A gyökérjegyzetfüzet végrehajtási eredménye a következő:

Képernyőkép a jegyzetfüzetek listájáról.

Visszaadott érték

A runMultiple() metódus egy szótárat ad vissza, amelyben minden kulcs a tevékenység neve, és minden érték egy szótár, amely a következő kulcsokkal rendelkezik:

  • exitVal: A gyermekjegyzetfüzet exit() hívása által visszaadott karakterlánc, vagy üres karakterlánc, ha a exit() nem hívták meg.
  • exception: Hibaobjektum, ha a tevékenység sikertelen volt, vagy None sikeres volt.

Jegyzetfüzetek futtatása DAG-struktúrával

Az alábbi példa egy DAG-struktúrában futtat jegyzetfüzeteket a használatával 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})

A gyökérjegyzetfüzet végrehajtási eredménye a következő:

Képernyőkép a paraméterekkel rendelkező jegyzetfüzetek listájáról.

DAG paraméterhivatkozás

Az alábbi táblázat a DAG-definícióban használható mezőket ismerteti:

szakterület szint Szükséges Leírás
activities Gyökér Igen A futtatandó jegyzetfüzeteket meghatározó tevékenységobjektumok listája.
timeoutInSeconds Gyökér No A teljes DAG maximális időtúllépése. Az alapértelmezett érték 43200 (12 óra).
concurrency Gyökér No Az egyidejűleg futtatandó jegyzetfüzetek maximális száma. Az alapértelmezett érték a processzormagok számának háromszorosa. Ezt az értéket explicit módon állíthatja be, ha szigorúbb vezérlésre van szüksége, vagy korlátlan egyidejűséghez használja 0 .
name Activity Igen A tevékenység egyedi neve. Az eredmények azonosítására és a függőségek meghatározására szolgál.
path Activity Igen A végrehajtandó jegyzetfüzetelem neve vagy elérési útja.
timeoutPerCellInSeconds Activity No Maximális időtúllépés a gyermekjegyzetfüzet minden egyes cellájában. Az alapértelmezett érték 90 másodperc.
args Activity No A gyermekjegyzetfüzetnek átadni kívánt paraméterek szótára.
workspace Activity No A munkaterület neve vagy azonosítója, ahol a jegyzetfüzet található. Alapértelmezés szerint a gyermekjegyzetfüzet ugyanabban a munkaterületen fut, mint a hívó.
retry Activity No Újrapróbálkozási kísérletek száma, ha a tevékenység meghiúsul. Az alapértelmezett érték 0.
retryIntervalInSeconds Activity No Várakozási idő másodpercben az újrapróbálkozások között. Az alapértelmezett érték 0.
dependencies Activity No Azoknak a tevékenységneveknek a listája, amelyeket a tevékenység megkezdése előtt végre kell hajtani.

Hivatkozás a tevékenységek közötti kilépési értékekre

Használhatja a args kifejezést egy függőségi tevékenység kilépési értékére való hivatkozáshoz a @activity() mezőben. Ez a minta lehetővé teszi az adatok továbbítását a jegyzetfüzetek között a DAG-ban.

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)

Jótanács

@activity('activity_name').exitValue() A mezőben található args kifejezéssel az egyik tevékenység eredményeit továbbíthatja egy másiknak egy DAG-on belül.

Dinamikus DAG létrehozása

Programmatikai szempontból hozhat létre DAG-struktúrákat olyan forgatókönyvekhez, mint például a szétosztási feldolgozás több partíción keresztül.

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)

DAG ellenőrzése

Annak ellenőrzésére használható validateDAG() , hogy a DAG-struktúra érvényes-e a végrehajtás előtt. Olyan problémákat fog fel, mint az ismétlődő tevékenységnevek, a hiányzó függőségek és a körkörös hivatkozások.

notebookutils.notebook.validateDAG(DAG)

Visszaadott érték

A validateDAG() metódus akkor ad vissza, True ha a DAG-struktúra érvényes, vagy kivételt eredményez, ha az ellenőrzés sikertelen.

Jótanács

Mindig hívja meg a validateDAG()-t az runMultiple() előtt a produkciós munkafolyamatokban, hogy korán elkapja a strukturális hibákat.

RunMultiple hibák kezelése

A runMultiple() metódus egy szótárt ad vissza, amelyben minden kulcs a tevékenység neve, és minden érték tartalmaz egy (sztringet exitVal ) és egy exception (hibaobjektumot vagy None). A részleges eredményeket akkor is megvizsgálhatja, ha egyes tevékenységek meghiúsulnak:

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

Megfontolások

  • A több jegyzetfüzet-futtatás párhuzamossági foka a Spark-munkamenet teljes rendelkezésre álló számítási erőforrására korlátozódik.
  • Az egyidejű jegyzetfüzetek alapértelmezett száma a rendelkezésre álló processzormagok számának háromszorosa. Ezt az értéket testre szabhatja, de a túlzott párhuzamosság stabilitási és teljesítménybeli problémákhoz vezethet a magas számítási erőforrás-használat miatt. Ha problémák merülnek fel, érdemes lehet a jegyzetfüzeteket több runMultiple hívásra különválasztani, vagy csökkenteni az egyidejűséget a DAG paraméter egyidejűségi mezőjének módosításával.
  • A teljes DAG alapértelmezett időtúllépése 12 óra, a gyermekjegyzetfüzet minden cellájának alapértelmezett időtúllépése pedig 90 másodperc. Az időtúllépést a DAG paraméter timeoutInSeconds és timeoutPerCellInSeconds mezőinek beállításával módosíthatja.
  • Konfigurálja a retry és retryIntervalInSeconds elemeket azon tevékenységekhez, amelyek átmeneti problémák, például hálózati időtúllépés vagy ideiglenes szolgáltatás elérhetetlensége miatt meghiúsulhatnak.
  • A párhuzamos jegyzetfüzetek egyetlen Spark-munkameneten belül osztják meg a számítási erőforrásokat. Az erőforrás-kihasználtság monitorozása a memóriaterhelés és a cpu-versengés elkerülése érdekében.

Kilépés a jegyzetfüzetből

A exit() metódus egy értékkel rendelkező jegyzetfüzetből kilép. A beágyazott függvényhívásokat interaktívan vagy folyamatban is futtathatja egy jegyzetfüzetben.

  • Amikor egy függvényt egy jegyzetfüzetből interaktívan hív meg, a Fabric jegyzetfüzet kivételt dob, kihagyja a további cellák futtatását, és életben tartja a Spark munkamenetet.

  • Ha egy függvényt meghívó exit() folyamatban lévő jegyzetfüzetet vezényel, a jegyzetfüzet-tevékenység kilépési értékkel tér vissza. Ezzel befejezi a folyamatfuttatást, és leállítja a Spark-munkamenetet.

  • Amikor meghív egy függvényt egy exit() hivatkozott jegyzetfüzetben, a Fabric Spark leállítja a hivatkozott jegyzetfüzet további végrehajtását, és továbbra is futtatja a függvényt hívó fő jegyzetfüzet következő celláit run() . Például: A Jegyzetfüzet1 három cellával rendelkezik, és meghív egy függvényt exit() a második cellában. A Jegyzetfüzet2 öt cellával és hívásokkal run(notebook1) rendelkezik a harmadik cellában. Amikor futtatja a Notebook2-t, a Notebook1 a második cellánál megáll, amikor eléri a exit() függvényt. A Notebook2 továbbra is futtatja a negyedik és az ötödik celláit.

notebookutils.notebook.exit("value string")

Visszatérési viselkedés

A exit() metódus nem ad vissza értéket. Leállítja az aktuális jegyzetfüzetet, és átadja a megadott sztringet a hívó jegyzetfüzetnek vagy folyamatnak.

Megjegyzés:

A exit() függvény felülírja az aktuális cellakimenetet. Az egyéb kódkivonatok kimenetének elvesztésének elkerülése érdekében hívja meg a notebookutils.notebook.exit() egy külön cellában.

Fontos

Ne hívja meg a notebookutils.notebook.exit()-t a try-catch blokkon belül. A kilépés nem lép érvénybe, ha kivételkezelés alá kerül. A exit() hívásnak a kód legfelső szintjén kell lennie a megfelelő működéshez.

Például:

A Minta1 jegyzetfüzet a következő két cellával rendelkezik:

  • Az 1. cella egy bemeneti paramétert határoz meg, amelynek alapértelmezett értéke 10.

  • A 2. cella az input kilépési értékként lép ki a jegyzetfüzetből.

Képernyőkép a kilépési függvény mintajegyzetfüzetéről.

A Minta1 egy másik, alapértelmezett értékekkel rendelkező jegyzetfüzetben is futtatható:

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

kimenet:

10

A Minta1 egy másik jegyzetfüzetben is futtatható, és a bemeneti érték 20 lehet:

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

kimenet:

20