หมายเหตุ
การเข้าถึงหน้านี้ต้องได้รับการอนุญาต คุณสามารถลอง ลงชื่อเข้าใช้หรือเปลี่ยนไดเรกทอรีได้
การเข้าถึงหน้านี้ต้องได้รับการอนุญาต คุณสามารถลองเปลี่ยนไดเรกทอรีได้
นําไปใช้กับ:✅ Fabric วิศวกรรมข้อมูลและวิทยาศาสตร์ข้อมูล
เซสชันการทํางานพร้อมกันสูง (HC) ช่วยให้ผู้โทรหลายคนแชร์เซสชัน Spark เดียวโดยไม่รบกวนซึ่งกันและกัน แทนที่จะจัดเตรียมเซสชันแยกต่างหากสําหรับทุกปริมาณงาน คุณจะได้รับเซสชัน HC และ Fabric API จะกําหนด REPL ที่แยกได้ภายในเซสชันพื้นฐานที่ใช้ร่วมกัน
ในบทความนี้ คุณใช้ Fabric Livy API เพื่อรับเซสชัน HC ตรวจสอบการบรรจุเซสชัน เรียกใช้คําสั่งแบบขนาน และยืนยันการแยก REPL
ข้อกำหนดเบื้องต้น
- Fabric ความจุพรีเมียมหรือทดลองใช้ พร้อมเลคเฮาส์
- ไคลเอนต์ระยะไกล เช่น Visual Studio Code ที่มี PySpark และ Python 3.10+
- บริการหลัก (SPN) ของ Microsoft Entra ที่มีการเข้าถึงพื้นที่ทํางาน ลงทะเบียนใบสมัครกับ แพลตฟอร์มข้อมูลประจำตัวของ Microsoft
- ข้อมูลลับของไคลเอ็นต์สําหรับบริการหลัก เพิ่มและจัดการข้อมูลประจําตัวของแอปพลิเคชัน
แทนที่ตัวยึด {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}")
ดูงานของคุณในฮับการตรวจสอบ
- ไปที่ ตรวจสอบ ในการนําทางด้านซ้าย
- เลือกชื่อกิจกรรมล่าสุดเพื่อดูรายละเอียดเซสชัน
- โปรดสังเกตว่าเซสชัน 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พื้นฐาน