เริ่มต้นใช้งาน Livy API สําหรับเซสชันการทํางานพร้อมกันสูง Fabric

นําไปใช้กับ:✅ Fabric วิศวกรรมข้อมูลและวิทยาศาสตร์ข้อมูล

เซสชันการทํางานพร้อมกันสูง (HC) ช่วยให้ผู้โทรหลายคนแชร์เซสชัน Spark เดียวโดยไม่รบกวนซึ่งกันและกัน แทนที่จะจัดเตรียมเซสชันแยกต่างหากสําหรับทุกปริมาณงาน คุณจะได้รับเซสชัน HC และ Fabric API จะกําหนด REPL ที่แยกได้ภายในเซสชันพื้นฐานที่ใช้ร่วมกัน

ในบทความนี้ คุณใช้ Fabric Livy API เพื่อรับเซสชัน HC ตรวจสอบการบรรจุเซสชัน เรียกใช้คําสั่งแบบขนาน และยืนยันการแยก REPL

ข้อกำหนดเบื้องต้น

แทนที่ตัวยึด {Entra_TenantID}, {Entra_ClientID}, , {Entra_ClientSecret}{Fabric_WorkspaceID}และ {Fabric_LakehouseID} ด้วยค่าของคุณเมื่อทําตามตัวอย่างในบทความนี้

เซสชันการทํางานพร้อมกันสูงคืออะไร?

เซสชันการทํางานพร้อมกันสูง (HC) ช่วยให้ผู้ใช้หรือกระบวนการหลายคนแชร์เซสชัน Spark เดียวได้ ผู้โทรแต่ละคนจะได้รับ REPL (Read-Eval-Print Loop) ที่แยกจากกันภายในเซสชันที่ใช้ร่วมกัน ข้อความจากผู้โทรที่แตกต่างกันจะไม่รบกวนซึ่งกันและกัน

การบรรจุเซสชัน

เมื่อคุณสร้างเซสชัน HC สองเซสชันที่มี sessionTag เดียวกัน API Fabric จะบรรจุเซสชันเหล่านั้นลงในเซสชัน Livy พื้นฐานเดียวกันเดียวกัน แต่ละเซสชัน HC จะได้รับ REPL ของตัวเอง ซึ่งให้:

  • ประสิทธิภาพของทรัพยากร: ผู้ใช้หลายคนแชร์เซสชัน Spark หนึ่งเซสชันแทนที่จะสร้างเซสชันของตนเอง
  • การแยก REPL: ตัวแปรและสถานะใน REPL หนึ่งจะไม่ปรากฏให้ผู้อื่นเห็น
  • การดําเนินการแบบขนาน: คําสั่งบน REPL ที่แตกต่างกันสามารถทํางานพร้อมกันได้

รหัสคีย์

ID ไม่ซ้ํากันต่อ ใช้สําหรับ
เซสชั่น HC id เซสชั่น HC สถานะโพล ลบเซสชัน
sessionId เซสชัน Livy (ใช้ร่วมกัน เมื่อบรรจุ) URL ของใบแจ้งยอด
replId REPL (บริบทที่แยกได้) URL ของใบแจ้งยอด

สำคัญ

และsessionIdreplIdจะใช้ได้เฉพาะเมื่อเซสชัน HC ถึงIdleสถานะ

เซสชัน HC แตกต่างจากเซสชัน Livy ปกติอย่างไร

แอสเพ็กต์ เซสชั่น Livy ปกติ เซสชั่น HC
จุดสิ้นสุด .../sessions .../highConcurrencySessions
คําชี้แจง ส่งไปยังเซสชันโดยตรง ส่งผ่าน REPL (/repls/{replId}/statements)
การเข้าซื้อกิจการ เซสชันจะกลายเป็น idle โดยตรง NotStartedจากนั้นAcquiringHighConcurrencySessionIdle
การบรรจุเซสชัน ไม่มีผลบังคับใช้ ตัวเลือก sessionTag ในการแชร์เซสชัน Spark พื้นฐาน

คําแนะนําทีละขั้นตอน

1. รับรองความถูกต้องด้วย Microsoft Entra

รับโทเค็นการเข้าถึงโดยใช้โฟลว์ข้อมูลประจําตัวไคลเอ็นต์ SPN แทนที่ค่าตัวยึดด้วยข้อมูลประจําตัวจริงของคุณ

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. สร้างเซสชัน HC สองเซสชันด้วยแท็กเซสชันเดียวกัน

สร้างเซสชัน HC สองเซสชันโดยใช้sessionTag: "demo-tag" เนื่องจากใช้แท็กเดียวกัน API Fabric จึงบรรจุลงในเซสชัน Livy พื้นฐานเดียวกันเดียวกัน แต่ละเซสชันจะได้รับ 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. สํารวจทั้งสองเซสชันจนกว่าจะพร้อมและตรวจสอบการบรรจุเซสชัน

แต่ละเซสชันจะเปลี่ยนผ่านสถานะเหล่านี้: NotStarted, , AcquiringHighConcurrencySessionและ Idleจากนั้น

เมื่อทั้งสองเซสชันเป็น Idleเอาต์พุตจะยืนยันรายละเอียดต่อไปนี้เกี่ยวกับการบรรจุเซสชัน:

  • รหัสเซสชัน HC สองรายการ (hc_id_a และ hc_id_b) แตกต่างกัน โดยยืนยันว่าการเรียก "รับ" แต่ละครั้งส่งคืนเซสชัน HC ที่แตกต่างกัน
  • รหัสเซสชัน Livy พื้นฐาน (sessionId_a และ sessionId_b) ตรงกัน ซึ่งยืนยันว่าเซสชัน HC ทั้งสองถูกบรรจุไว้ในเซสชัน Livy เดียวกัน
  • รหัส REPL (replId_a และ replId_b) แตกต่างกัน โดยยืนยันว่าแต่ละเซสชัน HC มีบริบทการดําเนินการที่แยกจากกันของตัวเอง

โค้ดต่อไปนี้จะสํารวจทั้งสองเซสชันจนกว่าจะพร้อมและพิมพ์ผลลัพธ์การยืนยัน:

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. ส่งใบแจ้งยอดไปยัง REPL ทั้งสองพร้อมกัน

ส่งคําขอ POST สองรายการ (หนึ่งรายการต่อ REPL) ก่อนสํารวจผลลัพธ์ เนื่องจาก REPL ใช้เซสชัน Spark เดียวกัน คําสั่งทั้งสองจึงสามารถทํางานพร้อมกันได้ รหัสนี้ยังกําหนด poll_statement ฟังก์ชันตัวช่วยที่ใช้ในขั้นตอนที่เหลือ

# 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. ตรวจสอบการแยก REPL

ตั้งค่าตัวแปร x = 42 ใน REPL A จากนั้นลองเข้าถึงตัวแปรจาก REPL B แม้ว่า REPL ทั้งสองจะใช้เซสชัน Spark เดียวกัน แต่ ตัวแปรจะถูกแยกออกจากกัน

# 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. ล้างเซสชัน HC ทั้งสอง

ลบเซสชัน HC ทั้งสองเพื่อนําทรัพยากรออกใช้ ใช้เซสชัน idHC ไม่ใช่ พื้นฐาน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}")

ดูงานของคุณในฮับการตรวจสอบ

  1. ไปที่ ตรวจสอบ ในการนําทางด้านซ้าย
  2. เลือกชื่อกิจกรรมล่าสุดเพื่อดูรายละเอียดเซสชัน
  3. โปรดสังเกตว่าเซสชัน HC ทั้งสองใช้เซสชัน Spark พื้นฐานเดียวกัน ซึ่งจะยืนยันการบรรจุเซสชัน

การอ้างอิงปลายทาง API

การดำเนินการ เมธอด ปลายทาง
สร้างเซสชัน HC POST /v1/workspaces/{workspaceId}/lakehouses/{lakehouseId}/livyapi/versions/2023-12-01/highConcurrencySessions
รับเซสชัน HC GET .../highConcurrencySessions/{highConcurrencySessionId}
ลบเซสชัน HC DELETE .../highConcurrencySessions/{highConcurrencySessionId}
ส่งใบแจ้งยอด POST .../highConcurrencySessions/{sessionId}/repls/{replId}/statements
รับใบแจ้งยอด GET .../highConcurrencySessions/{sessionId}/repls/{replId}/statements/{statementId}
ยกเลิกใบแจ้งยอด POST .../highConcurrencySessions/{sessionId}/repls/{replId}/statements/{statementId}/cancel

Note

การดําเนินการสร้าง รับ และลบใช้เซสชัน idHC การดําเนินการใบแจ้งยอดใช้ Livy sessionIdพื้นฐาน