Erste Schritte mit der Livy-API für Fabric Hohen Parallelitätssitzungen

Gilt für:✅ Fabric Data Engineering and Data Science

Sitzungen mit hoher Parallelität (HC) ermöglichen es mehreren Benutzern, eine einzelne Spark-Sitzung gemeinsam zu nutzen, ohne sich gegenseitig zu stören. Statt für jede Workload eine separate Sitzung bereitzustellen, erwirbt man eine HC-Sitzung, und die Fabric-API weist dieser innerhalb einer freigegebenen zugrunde liegenden Sitzung eine isolierte REPL zu.

In diesem Artikel verwenden Sie die Fabric Livy-API, um HC-Sitzungen zu erwerben, Sitzungsverpackungen zu überprüfen, Anweisungen parallel auszuführen und REPL-Isolation zu bestätigen.

Voraussetzungen

Ersetzen Sie die Platzhalter {Entra_TenantID}, {Entra_ClientID}, {Entra_ClientSecret}, {Fabric_WorkspaceID} und {Fabric_LakehouseID} durch Ihre Werte, wenn Sie den Beispielen in diesem Artikel folgen.

Was sind Sitzungen mit hoher Gleichzeitigkeit?

Hohe Parallelitätssitzungen (HC) ermöglichen es mehreren Benutzern oder Prozessen, eine einzelne Spark-Sitzung gemeinsam zu nutzen. Jeder Anrufer erhält innerhalb der gemeinsamen Sitzung eine isolierte REPL (Read-Eval-Print Loop). Aussagen von verschiedenen Anrufern beeinträchtigen einander nicht.

Sitzungsverpackung

Wenn Sie zwei HC-Sitzungen mit demselben sessionTag erstellen, packt die Fabric-API sie in die same zugrunde liegende Livy-Sitzung. Jede HC-Sitzung erhält eine eigene REPL, die Folgendes bietet:

  • Ressourceneffizienz: Mehrere Benutzer teilen eine Spark-Sitzung, statt jede eigene Sitzung zu erstellen.
  • REPL-Isolation: Variablen und Status in einer REPL sind für andere Personen nicht sichtbar.
  • Parallele Ausführung: Anweisungen für verschiedene REPLs können gleichzeitig ausgeführt werden.

Tasten-IDs

ID Eindeutig pro Verwendung
HC-Sitzung id HC-Sitzung Umfragestatus, Sitzung löschen
sessionId Livy-Sitzung (freigegeben wenn verpackt) URLs der Aussage
replId REPL (isolierter Kontext) URLs der Aussage

Von Bedeutung

Die sessionId und replId sind nur verfügbar, wenn die HC-Sitzung den Zustand Idle erreicht.

Unterschiede zwischen HC-Sitzungen und regelmäßigen Livy-Sitzungen

Aspekt Regelmäßige Livy-Sitzung HC-Sitzung
Endpunkt .../sessions .../highConcurrencySessions
Anweisungen Direkt an die Sitzung übermittelt Übermittelt über eine REPL (/repls/{replId}/statements)
Erwerb Die Sitzung wird idle direkt NotStarteddann AcquiringHighConcurrencySessionIdle
Sitzungsverpackung Nicht anwendbar Optional sessionTag zum Teilen zugrunde liegender Spark-Sitzungen

Schritt-für-Schritt-Anleitung

1. Authentifizieren mit Microsoft Entra

Erhalten eines Zugriffstokens mithilfe des SPN-Clientanmeldeinformationsflusses. Ersetzen Sie die Platzhalterwerte durch Ihre tatsächlichen Anmeldeinformationen.

from msal import ConfidentialClientApplication

# Configuration — Replace with your actual values
tenant_id = "{Entra_TenantID}"       # Microsoft Entra tenant ID
client_id = "{Entra_ClientID}"       # Service principal application ID
client_secret = "{Entra_ClientSecret}"  # Service principal client secret

# OAuth settings
authority = f"https://login.microsoftonline.com/{tenant_id}"
scope = "https://analysis.windows.net/powerbi/api/.default"

app = ConfidentialClientApplication(
    client_id=client_id,
    authority=authority,
    client_credential=client_secret,
)

result = app.acquire_token_for_client(scopes=[scope])

if "access_token" in result:
    token = result["access_token"]
    print("Access token acquired successfully.")
else:
    raise RuntimeError(
        f"Failed to acquire token: {result.get('error_description', 'unknown error')}"
    )

2. Erstellen Sie zwei HC-Sitzungen mit demselben Sitzungstag

Erstellen Sie zwei HC-Sitzungen mithilfe von sessionTag: "demo-tag". Da sie dasselbe Tag teilen, packt die Fabric-API sie auf dieselbe zugrunde liegende Livy-Sitzung. Jede Sitzung erhält eine eigene isolierte REPL.

import json
import requests

# Fabric resource IDs — Replace with your actual values
workspace_id = "{Fabric_WorkspaceID}"
lakehouse_id = "{Fabric_LakehouseID}"

# Construct the HC session endpoint URL
livy_base_url = (
    f"https://api.fabric.microsoft.com/v1"
    f"/workspaces/{workspace_id}"
    f"/lakehouses/{lakehouse_id}"
    f"/livyapi/versions/2023-12-01"
    f"/highConcurrencySessions"
)

headers = {"Authorization": f"Bearer {token}"}
session_tag = "demo-tag"

print(f"HC session endpoint: {livy_base_url}")
print(f"Session tag: {session_tag}")
print()

# Create HC Session A
print("Creating HC Session A...")
resp_a = requests.post(livy_base_url, headers=headers, json={"sessionTag": session_tag})
assert resp_a.status_code == 202, f"Failed: {resp_a.status_code} — {resp_a.text}"
session_a = resp_a.json()
hc_id_a = session_a["id"]
print(f"  HC session A id: {hc_id_a}  state: {session_a['state']}")

# Create HC Session B
print("Creating HC Session B...")
resp_b = requests.post(livy_base_url, headers=headers, json={"sessionTag": session_tag})
assert resp_b.status_code == 202, f"Failed: {resp_b.status_code} — {resp_b.text}"
session_b = resp_b.json()
hc_id_b = session_b["id"]
print(f"  HC session B id: {hc_id_b}  state: {session_b['state']}")

session_url_a = f"{livy_base_url}/{hc_id_a}"
session_url_b = f"{livy_base_url}/{hc_id_b}"

3. Beide Sitzungen so lange abfragen, bis sie bereit sind, und dann die Sitzungskonfiguration überprüfen.

Jede Sitzung wechselt durch diese Zustände: NotStarted, AcquiringHighConcurrencySession, und dann Idle.

Sobald beide Sitzungen abgeschlossen sind Idle, bestätigt die Ausgabe die folgenden Details zum Session-Pack-Prozess:

  • Die beiden HC-Sitzungs-IDs (hc_id_a und hc_id_b) unterscheiden sich, wobei bestätigt wird, dass jeder "Acquire"-Aufruf eine eindeutige HC-Sitzung zurückgegeben hat.
  • Die zugrunde liegenden Livy-Sitzungs-IDs (sessionId_a und sessionId_b) stimmen überein und bestätigen, dass beide HC-Sitzungen in derselben Livy-Sitzung verpackt wurden.
  • Die REPL-IDs (replId_a und replId_b) unterscheiden sich, wobei bestätigt wird, dass jede HC-Sitzung über einen eigenen isolierten Ausführungskontext verfügt.

Der folgende Code fragt beide Sitzungen ab, bis sie bereit sind, und druckt die Überprüfungsausgabe:

import time

ACQUIRING_STATES = {"NotStarted", "starting", "AcquiringHighConcurrencySession"}
POLL_INTERVAL = 5

def poll_until_ready(url, label):
    """Poll an HC session until it leaves the acquisition states."""
    print(f"[{label}] Polling...")
    while True:
        resp = requests.get(url, headers=headers, timeout=30)
        resp.raise_for_status()
        data = resp.json()
        state = data.get("state", "unknown")
        print(f"  [{label}] state={state}  sessionId={data.get('sessionId', 'N/A')}  replId={data.get('replId', 'N/A')}")
        if state in ("Dead", "Killed", "Failed"):
            raise RuntimeError(f"[{label}] Session failed: {state}")
        if state not in ACQUIRING_STATES:
            return data
        time.sleep(POLL_INTERVAL)

ready_a = poll_until_ready(session_url_a, "A")
ready_b = poll_until_ready(session_url_b, "B")

livy_session_id_a = ready_a["sessionId"]
livy_session_id_b = ready_b["sessionId"]
repl_id_a = ready_a["replId"]
repl_id_b = ready_b["replId"]

print()
print("=" * 50)
print("SESSION PACKING VERIFICATION")
print("=" * 50)
print(f"HC session A id:    {hc_id_a}")
print(f"HC session B id:    {hc_id_b}")
print(f"HC IDs differ:      {hc_id_a != hc_id_b}")
print()
print(f"Livy sessionId A:   {livy_session_id_a}")
print(f"Livy sessionId B:   {livy_session_id_b}")
print(f"Same Livy session:  {livy_session_id_a == livy_session_id_b}")
print()
print(f"REPL A:             {repl_id_a}")
print(f"REPL B:             {repl_id_b}")
print(f"REPLs differ:       {repl_id_a != repl_id_b}")

4. Anweisungen parallel an beide REPLs übermitteln

Übermitteln Sie zwei POST-Anfragen (eine pro REPL), bevor Sie eine der beiden auf Ergebnisse abfragen. Da die REPLs dieselbe Spark-Sitzung verwenden, können beide Anweisungen gleichzeitig ausgeführt werden. Dieser Code definiert auch die poll_statement Hilfsfunktion, die in den verbleibenden Schritten verwendet wird.

# Build statement URLs for each REPL
stmts_url_a = f"{livy_base_url}/{livy_session_id_a}/repls/{repl_id_a}/statements"
stmts_url_b = f"{livy_base_url}/{livy_session_id_b}/repls/{repl_id_b}/statements"

# Fire both statement POSTs before polling
print("Submitting to REPL A: print('Hello from REPL A')")
resp_a = requests.post(stmts_url_a, headers=headers, json={"code": "print('Hello from REPL A')", "kind": "pyspark"})
assert resp_a.status_code in (200, 201), f"Failed: {resp_a.text}"
stmt_a = resp_a.json()
stmt_url_a = f"{stmts_url_a}/{stmt_a['id']}"

print("Submitting to REPL B: print('Hello from REPL B')")
resp_b = requests.post(stmts_url_b, headers=headers, json={"code": "print('Hello from REPL B')", "kind": "pyspark"})
assert resp_b.status_code in (200, 201), f"Failed: {resp_b.text}"
stmt_b = resp_b.json()
stmt_url_b = f"{stmts_url_b}/{stmt_b['id']}"

print("Both statements submitted. Polling for results...")

# Poll both statements
def poll_statement(url, label):
    while True:
        resp = requests.get(url, headers=headers, timeout=30)
        resp.raise_for_status()
        data = resp.json()
        if data.get("state") not in ("waiting", "running"):
            return data
        time.sleep(5)

result_a = poll_statement(stmt_url_a, "A")
result_b = poll_statement(stmt_url_b, "B")

output_a = result_a.get("output", {}).get("data", {}).get("text/plain", "")
output_b = result_b.get("output", {}).get("data", {}).get("text/plain", "")

print()
print("=" * 50)
print("PARALLEL EXECUTION RESULTS")
print("=" * 50)
print(f"REPL A output: {output_a}")
print(f"REPL B output: {output_b}")

5. Überprüfen der REPL-Isolation

Legen Sie eine Variable x = 42 in REPL A fest, und versuchen Sie dann, von REPL B darauf zuzugreifen. Obwohl beide REPLs dieselbe Spark-Sitzung gemeinsam nutzen, sind ihre Variablen isoliert.

# Set x = 42 in REPL A
print("[A] Setting x = 42...")
resp = requests.post(stmts_url_a, headers=headers, json={"code": "x = 42; print(x)", "kind": "pyspark"})
stmt_url = f"{stmts_url_a}/{resp.json()['id']}"
result_a = poll_statement(stmt_url, "A")
output_a = result_a.get("output", {}).get("data", {}).get("text/plain", "")
print(f"[A] Output: {output_a}")

# Try to read x from REPL B — should get NameError
print("\n[B] Trying to read x (expect NameError)...")
code_b = "try:\n    print(x)\nexcept NameError as e:\n    print(f'NameError: {e}')"
resp = requests.post(stmts_url_b, headers=headers, json={"code": code_b, "kind": "pyspark"})
stmt_url = f"{stmts_url_b}/{resp.json()['id']}"
result_b = poll_statement(stmt_url, "B")
output_b = result_b.get("output", {}).get("data", {}).get("text/plain", "")
print(f"[B] Output: {output_b}")

print()
print("=" * 50)
print("REPL ISOLATION RESULTS")
print("=" * 50)
print(f"REPL A (x = 42): {output_a}")
print(f"REPL B (print(x)): {output_b}")

6. Bereinigen beider HC-Sitzungen

Löschen Sie beide HC-Sitzungen, um Ressourcen freizugeben. Verwenden Sie die HC-Sitzung id, nicht die darunterliegende sessionId.

for label, url in [("A", session_url_a), ("B", session_url_b)]:
    print(f"[{label}] Deleting HC session...")
    resp = requests.delete(url, headers=headers)
    if resp.status_code in (200, 204):
        print(f"[{label}] Deleted successfully.")
    elif resp.status_code == 404:
        print(f"[{label}] Already deleted.")
    else:
        print(f"[{label}] Unexpected response: {resp.status_code} — {resp.text}")

Ansehen Ihrer Jobs im Monitoring-Hub

  1. Navigieren Sie in der linken Navigationsleiste zu Monitor .
  2. Wählen Sie den namen der letzten Aktivität aus, um Sitzungsdetails anzuzeigen.
  3. Beachten Sie, dass beide HC-Sitzungen dieselbe zugrunde liegende Spark-Sitzung verwenden, die das Packen von Sitzungen bestätigt.

Referenz der API-Schnittstellen

Vorgang Methode Endpunkt
HC-Sitzung erstellen POST /v1/workspaces/{workspaceId}/lakehouses/{lakehouseId}/livyapi/versions/2023-12-01/highConcurrencySessions
HC-Sitzung abrufen GET .../highConcurrencySessions/{highConcurrencySessionId}
HC-Sitzung löschen DELETE .../highConcurrencySessions/{highConcurrencySessionId}
Submit-Anweisung POST .../highConcurrencySessions/{sessionId}/repls/{replId}/statements
Get-Anweisung GET .../highConcurrencySessions/{sessionId}/repls/{replId}/statements/{statementId}
Cancel-Anweisung POST .../highConcurrencySessions/{sessionId}/repls/{replId}/statements/{statementId}/cancel

Hinweis

Erstellen-, Get- und Delete-Vorgänge verwenden die HC-Sitzung id. Anweisungsoperationen verwenden die zugrunde liegende Livy sessionId.