Introduzione all'API Livy delle sessioni di concorrenza elevata di Fabric

Si applica a:✅ Fabric Data Engineering and Data Science

Le sessioni di concorrenza elevata (HC) consentono a più chiamanti di condividere una singola sessione Spark senza interferire tra loro. Invece di effettuare il provisioning di una sessione separata per ogni carico di lavoro, si acquisisce una sessione HC e l'API Fabric assegna un REPL isolato all'interno di una sessione sottostante condivisa.

In questo articolo si usa l'API Fabric Livy per acquisire sessioni HC, verificare la compressione della sessione, eseguire istruzioni in parallelo e confermare l'isolamento REPL.

Prerequisiti

Sostituire i segnaposto {Entra_TenantID}, {Entra_ClientID}, {Entra_ClientSecret}, {Fabric_WorkspaceID}e {Fabric_LakehouseID} con i valori quando si seguono gli esempi in questo articolo.

Che cosa sono le sessioni di concorrenza elevata?

Le sessioni di concorrenza elevata (HC) consentono a più utenti o processi di condividere una singola sessione Spark. Ogni chiamante ottiene un REPL isolato (Read-Eval-Print Loop) all'interno della sessione condivisa. Le dichiarazioni di chiamanti diversi non interferiscono tra loro.

Compressione della sessione

Quando si creano due sessioni HC con la stessa sessionTag, l'API Fabric li include nella sessione same sottostante Livy. Ogni sessione HC ottiene il proprio REPL, che fornisce:

  • Efficienza delle risorse: più utenti condividono una sessione Spark invece di crearne una personalizzata.
  • Isolamento REPL: le variabili e lo stato in un REPL non sono visibili ad altri utenti.
  • Esecuzione parallela: le istruzioni in repls diverse possono essere eseguite simultaneamente.

ID chiave

Documento d'identità Univoco per Usato per
Sessione HC id Sessione HC Stato del sondaggio, elimina sessione
sessionId Sessione livy (condivisa al momento della compilazione) URL delle dichiarazioni
replId REPL (contesto isolato) URL delle dichiarazioni

Importante

sessionId e replId sono disponibili solo dopo che la sessione HC raggiunge lo Idle stato.

Differenze tra le sessioni HC e le sessioni Livy regolari

Aspetto Sessione regolare di Livy Sessione HC
Punto finale .../sessions .../highConcurrencySessions
Dichiarazioni Inviato direttamente alla sessione Inviato tramite REPL (/repls/{replId}/statements)
Acquisizione La sessione diventa idle direttamente NotStartedpoi poi AcquiringHighConcurrencySessionIdle
Compressione della sessione Non applicabile Opzionale sessionTag per condividere sessioni Spark di base

Guida passo-passo

1. Eseguire l'autenticazione con Microsoft Entra

Acquisire un token di accesso usando il flusso delle credenziali client SPN. Sostituire i valori segnaposto con le credenziali effettive.

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. Creare due sessioni HC con lo stesso tag di sessione

Creare due sessioni HC utilizzando sessionTag: "demo-tag". Poiché condividono lo stesso tag, l'API Fabric li comprime nella sessione same sottostante Livy. Ogni sessione ottiene il proprio REPL isolato.

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. Eseguire il polling di entrambe le sessioni fino a quando non sono pronte e verificare la compressione delle sessioni

Ogni sessione passa attraverso questi stati: NotStarted, AcquiringHighConcurrencySessione quindi Idle.

Quando entrambe le sessioni sono Idle, l'output conferma i dettagli seguenti sulla compressione della sessione:

  • I due ID sessione HC (hc_id_a e hc_id_b) sono diversi, confermando che ogni chiamata "acquire" ha restituito una sessione HC distinta.
  • La corrispondenza degli ID sessione Livy sottostanti (sessionId_a e sessionId_b) conferma che entrambe le sessioni HC sono state compresse nella stessa sessione Livy.
  • Gli ID REPL (replId_a e replId_b) sono diversi, confermando che ogni sessione HC ha un proprio contesto di esecuzione isolato.

Il codice seguente esegue il polling di entrambe le sessioni fino a che sono pronte e stampa l'output di verifica.

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. Inviare istruzioni a entrambi i REPL in parallelo

Inviare due richieste POST (una per ciascun REPL) prima di interrogare ciascuno per ottenere i risultati. Poiché i REPL condividono la stessa sessione Spark, entrambe le istruzioni possono essere eseguite contemporaneamente. Questo codice definisce anche la poll_statement funzione helper usata nei passaggi rimanenti.

# 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. Verificare l'isolamento REPL

Impostare una variabile x = 42 in REPL A, quindi provare ad accedervi da REPL B. Anche se entrambi i REPL condividono la stessa sessione Spark, le relative variabili sono isolate.

# 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. Pulire entrambe le sessioni HC

Eliminare entrambe le sessioni HC per rilasciare le risorse. Usare la sessione idHC, non l'oggetto sottostante 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}")

Visualizza le attività nell'hub di monitoraggio

  1. Passare a Monitor nel menu di navigazione a sinistra.
  2. Selezionare il nome dell'attività più recente per visualizzare i dettagli della sessione.
  3. Si noti che entrambe le sessioni HC condividono la stessa sessione Spark sottostante, che conferma la compressione della sessione.

Informazioni di riferimento sugli endpoint API

Operation metodo Punto finale
Creare una sessione HC POST /v1/workspaces/{workspaceId}/lakehouses/{lakehouseId}/livyapi/versions/2023-12-01/highConcurrencySessions
Recuperare la sessione HC GET .../highConcurrencySessions/{highConcurrencySessionId}
Eliminare la sessione HC DELETE .../highConcurrencySessions/{highConcurrencySessionId}
Invia dichiarazione POST .../highConcurrencySessions/{sessionId}/repls/{replId}/statements
Istruzione Get GET .../highConcurrencySessions/{sessionId}/repls/{replId}/statements/{statementId}
Istruzione Annulla POST .../highConcurrencySessions/{sessionId}/repls/{replId}/statements/{statementId}/cancel

Annotazioni

Le operazioni di creazione, recupero ed eliminazione usano la sessione HC id. Le operazioni delle dichiarazioni utilizzano l'oggetto Livy sessionId sottostante.