Hinweis
Für den Zugriff auf diese Seite ist eine Autorisierung erforderlich. Sie können versuchen, sich anzumelden oder das Verzeichnis zu wechseln.
Für den Zugriff auf diese Seite ist eine Autorisierung erforderlich. Sie können versuchen, das Verzeichnis zu wechseln.
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
- Fabric Premium- oder Testkapazität mit einem Lakehouse.
- Ein Remoteclient wie Visual Studio Code mit PySpark und Python 3,10+.
- Ein Microsoft Entra Service Principal (SPN) mit Zugriff auf den Arbeitsbereich. Registern Sie eine Anwendung mit dem Microsoft Identity Platform.
- Ein Client-Geheimnis für den Serviceprinzipal. Hinzufügen und Verwalten von Anwendungsanmeldeinformationen
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_aundhc_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_aundsessionId_b) stimmen überein und bestätigen, dass beide HC-Sitzungen in derselben Livy-Sitzung verpackt wurden. - Die REPL-IDs (
replId_aundreplId_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
- Navigieren Sie in der linken Navigationsleiste zu Monitor .
- Wählen Sie den namen der letzten Aktivität aus, um Sitzungsdetails anzuzeigen.
- 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.