mssql-python 驅動程式使用同步輸入輸出,且不提供原生 async/await 支援。 原生非同步支援已列入驅動程式路線圖。 在此之前,你可以透過以下變通模式將 mssql-python 整合到非同步應用程式:
- 執行緒池執行器用於卸載阻塞呼叫。
- 同步作業的非同步包裝器。
- 與 FastAPI 等非同步框架整合。
Note
本文中使用 ThreadPoolExecutor 的模式在背景執行緒中執行同步的 mssql-python 呼叫。 此方法相較原生非同步驅動程式增加了執行緒開銷。 對於 I/O 綁定的資料庫工作負載,開銷通常是可接受的。
何時使用非同步模式
線程池方法在以下情況下運作良好:
- 你的應用程式已經在使用
asyncio(例如 FastAPI、aiohttp 或 Discord 機器人),你需要整合資料庫呼叫,同時不阻擋事件迴圈。 - 資料庫查詢是受 I/O 限制,而非 CPU 限制。 執行緒池讓事件迴圈在等待 Microsoft SQL 的同時處理其他請求。
- 你的並發性是中等程度(數十筆同時查詢,而非數千筆)。
純同步應用則可跳過這些模式。 直接使用驅動程式搭配同步程式碼,實現直接且較低的開銷執行。
執行緒池執行器模式
以下範例展示了如何將同步的 mssql-python 呼叫包裝在 a ThreadPoolExecutor 中,以便使用 asyncio。
基本的非同步包裝器
建立一個簡單的輔助函式,在執行緒池中執行同步的 mssql-python 操作並等待結果。
import asyncio
from concurrent.futures import ThreadPoolExecutor
import mssql_python
from functools import partial
from typing import Any, Callable
# Create a dedicated thread pool for database operations
db_executor = ThreadPoolExecutor(max_workers=10, thread_name_prefix="db_")
async def run_in_executor(func: Callable, *args, **kwargs) -> Any:
"""Run a synchronous function in the thread pool."""
loop = asyncio.get_running_loop()
if kwargs:
func = partial(func, **kwargs)
return await loop.run_in_executor(db_executor, func, *args)
# Database functions
def _execute_query(connection_string: str, query: str, params: dict = None) -> list:
"""Synchronous query execution."""
conn = mssql_python.connect(connection_string)
cursor = conn.cursor()
try:
cursor.execute(query, params or {})
if cursor.description:
columns = [col[0] for col in cursor.description]
return [dict(zip(columns, row)) for row in cursor.fetchall()]
return []
finally:
cursor.close()
conn.close()
def _execute_scalar(connection_string: str, query: str, params: dict = None) -> Any:
"""Synchronous scalar query."""
conn = mssql_python.connect(connection_string)
cursor = conn.cursor()
try:
cursor.execute(query, params or {})
return cursor.fetchval()
finally:
cursor.close()
conn.close()
# Async interfaces
async def async_query(connection_string: str, query: str, params: dict = None) -> list:
"""Execute query asynchronously."""
return await run_in_executor(_execute_query, connection_string, query, params)
async def async_scalar(connection_string: str, query: str, params: dict = None) -> Any:
"""Execute scalar query asynchronously."""
return await run_in_executor(_execute_scalar, connection_string, query, params)
# Usage
async def main():
conn_str = "Server=<server>.database.windows.net;Database=<database>;Authentication=ActiveDirectoryDefault;Encrypt=yes"
# Execute query asynchronously
products = await async_query(conn_str, "SELECT * FROM Production.Product WHERE ProductSubcategoryID = %(cat)s", {"cat": 5})
print(f"Found {len(products)} products")
# Execute scalar asynchronously
count = await async_scalar(conn_str, "SELECT COUNT(*) FROM Production.Product")
print(f"Total products: {count}")
asyncio.run(main())
Note
此範例會依序等待兩個查詢,因此它們會依序執行。 關鍵字 await 釋放事件迴圈,讓每個查詢在等待時執行其他任務,但不會重疊這兩個查詢。 若要同時執行獨立查詢,請將它們排程為 ,如asyncio.gather非同步連線池區段所示:
products, count = await asyncio.gather(
async_query(conn_str, "SELECT * FROM Production.Product WHERE ProductSubcategoryID = %(cat)s", {"cat": 5}),
async_scalar(conn_str, "SELECT COUNT(*) FROM Production.Product"),
)
非同步連線池
本節說明如何針對 mssql-python 的內建連線池建立適合非同步使用的包裝器,供 asyncio 應用程式使用。
Note
mssql-python 驅動程式內建連線池功能。 此處所示的非同步連線池會以非同步內容管理器包裝同步的集區連線,以供在 asyncio 應用程式中使用。 如果你只從執行緒池執行器呼叫 mssql-python,就不需要管理自訂池。
池化非同步資料庫類別
建立一個可重複使用的非同步連線池類別,管理 mssql-python 連線佇列並提供非同步查詢執行方法。
import asyncio
from concurrent.futures import ThreadPoolExecutor
from contextlib import asynccontextmanager
from typing import Any, Optional
import mssql_python
from dataclasses import dataclass
from queue import Queue, Empty
import threading
@dataclass
class PooledConnection:
"""Wrapper for pooled connection."""
connection: Any
cursor: Any
in_use: bool = False
class AsyncDatabasePool:
"""Async-friendly connection pool for mssql-python."""
def __init__(self, connection_string: str, pool_size: int = 10):
self.connection_string = connection_string
self.pool_size = pool_size
self._pool: Queue[PooledConnection] = Queue(maxsize=pool_size)
self._executor = ThreadPoolExecutor(max_workers=pool_size, thread_name_prefix="dbpool_")
self._lock = threading.Lock()
self._initialized = False
async def initialize(self):
"""Initialize the connection pool."""
if self._initialized:
return
loop = asyncio.get_running_loop()
async def create_connection():
def _create():
conn = mssql_python.connect(self.connection_string)
cursor = conn.cursor()
return PooledConnection(connection=conn, cursor=cursor)
return await loop.run_in_executor(self._executor, _create)
# Create initial connections
tasks = [create_connection() for _ in range(self.pool_size)]
connections = await asyncio.gather(*tasks)
for conn in connections:
self._pool.put(conn)
self._initialized = True
async def acquire(self, timeout: float = 30.0) -> PooledConnection:
"""Acquire a connection from the pool."""
loop = asyncio.get_running_loop()
def _acquire():
try:
conn = self._pool.get(timeout=timeout)
conn.in_use = True
return conn
except Empty:
raise TimeoutError("Could not acquire connection from pool")
return await loop.run_in_executor(self._executor, _acquire)
def release(self, conn: PooledConnection):
"""Release a connection back to the pool."""
conn.in_use = False
try:
conn.connection.commit()
except Exception:
conn.connection.rollback()
self._pool.put(conn)
@asynccontextmanager
async def connection(self):
"""Async context manager for getting a connection."""
conn = await self.acquire()
try:
yield conn
except Exception:
conn.connection.rollback()
raise
else:
conn.connection.commit()
finally:
self.release(conn)
async def execute(self, query: str, params: dict = None) -> list:
"""Execute query and return results."""
async with self.connection() as conn:
loop = asyncio.get_running_loop()
def _execute():
conn.cursor.execute(query, params or {})
if conn.cursor.description:
columns = [col[0] for col in conn.cursor.description]
return [dict(zip(columns, row)) for row in conn.cursor.fetchall()]
return []
return await loop.run_in_executor(self._executor, _execute)
async def execute_scalar(self, query: str, params: dict = None) -> Any:
"""Execute query and return single value."""
async with self.connection() as conn:
loop = asyncio.get_running_loop()
def _execute():
conn.cursor.execute(query, params or {})
return conn.cursor.fetchval()
return await loop.run_in_executor(self._executor, _execute)
async def close(self):
"""Close all connections in the pool."""
while not self._pool.empty():
try:
conn = self._pool.get_nowait()
conn.cursor.close()
conn.connection.close()
except Empty:
break
self._executor.shutdown(wait=True)
# Usage
async def main():
pool = AsyncDatabasePool("Server=<server>.database.windows.net;Database=<database>;Authentication=ActiveDirectoryDefault;Encrypt=yes", pool_size=5)
await pool.initialize()
try:
# Execute queries concurrently
tasks = [
pool.execute("SELECT * FROM Production.Product WHERE ProductSubcategoryID = %(cat)s", {"cat": i})
for i in range(1, 6)
]
results = await asyncio.gather(*tasks)
for i, products in enumerate(results, 1):
print(f"Category {i}: {len(products)} products")
# Single scalar query
total = await pool.execute_scalar("SELECT COUNT(*) FROM Production.Product")
print(f"Total: {total}")
finally:
await pool.close()
# Only run the demo when this file is executed directly, not when imported.
if __name__ == "__main__":
asyncio.run(main())
與 FastAPI 的整合
FastAPI 的 lifespan 情境管理器會自動處理池初始化與清理。
非同步 FastAPI 搭配 mssql-python
這個範例是建立在 非同步連線池 的部分之上。 將該區段的程式碼存到一個名為 db.py的檔案中,然後在旁邊的 main.py 檔案中建立下一個 FastAPI 應用程式。 pool 範例使用 if __name__ == "__main__": 保護其示範程式,因此匯入 db.py 時不會執行該示範程式。 這個應用程式會使用 lifespan 上下文管理器,在啟動時初始化資源池,並在關閉時將其清理。
from fastapi import FastAPI, Depends, HTTPException
from contextlib import asynccontextmanager
from typing import Optional
import asyncio
from db import AsyncDatabasePool # The pool class from the previous section
# Initialize pool on startup
pool: Optional[AsyncDatabasePool] = None
@asynccontextmanager
async def lifespan(app: FastAPI):
"""Manage database pool lifecycle."""
global pool
pool = AsyncDatabasePool(
"Server=<server>.database.windows.net;Database=<database>;Authentication=ActiveDirectoryDefault;Encrypt=yes",
pool_size=10
)
await pool.initialize()
yield
await pool.close()
app = FastAPI(lifespan=lifespan)
async def get_db():
"""Dependency for database access."""
return pool
@app.get("/products")
async def list_products(db: AsyncDatabasePool = Depends(get_db)):
products = await db.execute("SELECT ProductID, Name, ListPrice FROM Production.Product")
return {"products": products}
@app.get("/products/{product_id}")
async def get_product(product_id: int, db: AsyncDatabasePool = Depends(get_db)):
products = await db.execute(
"SELECT ProductID, Name, ListPrice FROM Production.Product WHERE ProductID = %(id)s",
{"id": product_id}
)
if not products:
raise HTTPException(status_code=404, detail="Product not found")
return products[0]
@app.get("/stats")
async def get_stats(db: AsyncDatabasePool = Depends(get_db)):
# Execute multiple queries concurrently
product_count, subcategory_count, total_value = await asyncio.gather(
db.execute_scalar("SELECT COUNT(*) FROM Production.Product"),
db.execute_scalar("SELECT COUNT(*) FROM Production.ProductSubcategory"),
db.execute_scalar("SELECT SUM(ListPrice) FROM Production.Product"),
)
return {
"products": product_count,
"subcategories": subcategory_count,
"total_value": float(total_value) if total_value else 0
}
安裝相依套件,並用像 Uvicorn 這樣的 ASGI 伺服器來執行應用程式。 從包含 main.py 和 db.py的資料夾執行此指令:
pip install fastapi uvicorn mssql-python
uvicorn main:app --reload
伺服器運行中,開啟 http://127.0.0.1:8000/products、 http://127.0.0.1:8000/products/1,或 http://127.0.0.1:8000/stats 呼叫每個端點。
背景任務
在不阻塞應用程式事件迴圈的情況下,定期執行排程的資料庫操作。
非同步背景工作者
實作一個任務執行器,能在指定間隔執行註冊的資料庫操作,避免重複並行執行。 這個範例是建立在 非同步連線池 的部分,所以請將該區段的程式碼儲存為 db.py。 然後,將以下程式碼存檔放在 worker.py 旁邊。 它會設定日誌,讓每次執行都回報結果。
import asyncio
from typing import Callable, Any
from dataclasses import dataclass
from datetime import datetime
import logging
from db import AsyncDatabasePool # The pool class from the Async connection pool section
logger = logging.getLogger(__name__)
@dataclass
class Task:
"""Background task definition."""
name: str
func: Callable
interval: float # seconds
last_run: datetime = None
running: bool = False
class AsyncTaskRunner:
"""Run database tasks in the background."""
def __init__(self, pool: AsyncDatabasePool):
self.pool = pool
self.tasks: dict[str, Task] = {}
self._running = False
def register(self, name: str, func: Callable, interval: float):
"""Register a periodic task."""
self.tasks[name] = Task(name=name, func=func, interval=interval)
async def _run_task(self, task: Task):
"""Execute a single task."""
if task.running:
return
task.running = True
try:
await task.func(self.pool)
task.last_run = datetime.now()
logger.info(f"Task {task.name} completed")
except Exception as e:
logger.error(f"Task {task.name} failed: {e}")
finally:
task.running = False
async def start(self):
"""Start the task runner."""
self._running = True
while self._running:
now = datetime.now()
for task in self.tasks.values():
if task.last_run is None or \
(now - task.last_run).total_seconds() >= task.interval:
asyncio.create_task(self._run_task(task))
await asyncio.sleep(1) # Check every second
def stop(self):
"""Stop the task runner."""
self._running = False
# Example tasks
async def count_products(pool: AsyncDatabasePool):
"""Read-only task: report the current product count."""
total = await pool.execute_scalar("SELECT COUNT(*) FROM Production.Product")
logger.info("count_products: %s products", total)
async def check_low_inventory(pool: AsyncDatabasePool):
"""Report how many products are below an inventory threshold."""
rows = await pool.execute(
"""
SELECT ProductID, LocationID, Quantity
FROM Production.ProductInventory
WHERE Quantity < %(threshold)s
""",
{"threshold": 100},
)
logger.info("check_low_inventory: %s rows below threshold", len(rows))
# Usage
async def main():
logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s %(message)s")
pool = AsyncDatabasePool(
"Server=<server>.database.windows.net;Database=<database>;Authentication=ActiveDirectoryDefault;Encrypt=yes",
pool_size=3,
)
await pool.initialize()
runner = AsyncTaskRunner(pool)
runner.register("count_products", count_products, interval=2)
runner.register("low_inventory", check_low_inventory, interval=3)
# Run the runner in the background, let it cycle a few times, then stop.
# In a real app, run the runner for the application's lifetime instead,
# for example from a FastAPI lifespan handler.
runner_task = asyncio.create_task(runner.start())
await asyncio.sleep(7)
runner.stop()
await runner_task
await pool.close()
if __name__ == "__main__":
asyncio.run(main())
啟動工作程序:
python worker.py
每個註冊的任務執行時都會記錄,所以你會看到每隔幾秒就重複一次的輸出:
2026-07-17 15:07:25 INFO count_products: 504 products
2026-07-17 15:07:25 INFO Task count_products completed
2026-07-17 15:07:26 INFO check_low_inventory: 179 rows below threshold
2026-07-17 15:07:26 INFO Task low_inventory completed
並行查詢執行
使用 asyncio.gather 與信號量來限制同時執行的查詢數量。 本節的範例建立在 非同步連線池(Async connection pool )部分,請將該區的程式碼存為 , db.py 並在旁邊的獨立檔案中執行每個範例。
帶有信號量的平行查詢
同時執行多個查詢,同時使用信號量限制同時操作次數,避免執行緒池飽和。
import asyncio
from db import AsyncDatabasePool # The pool class from the Async connection pool section
async def parallel_queries(pool: AsyncDatabasePool, queries: list[tuple[str, dict]],
max_concurrent: int = 5) -> list:
"""Execute multiple queries with concurrency limit."""
semaphore = asyncio.Semaphore(max_concurrent)
async def run_query(query: str, params: dict):
async with semaphore:
return await pool.execute(query, params)
tasks = [run_query(q, p) for q, p in queries]
return await asyncio.gather(*tasks)
# Usage
async def main():
pool = AsyncDatabasePool(
"Server=<server>.database.windows.net;Database=<database>;Authentication=ActiveDirectoryDefault;Encrypt=yes",
pool_size=5,
)
await pool.initialize()
try:
queries = [
("SELECT * FROM Production.Product WHERE ProductSubcategoryID = %(cat)s", {"cat": 1}),
("SELECT * FROM Production.Product WHERE ProductSubcategoryID = %(cat)s", {"cat": 2}),
("SELECT * FROM Sales.SalesOrderHeader WHERE Status = %(status)s", {"status": 5}),
("SELECT * FROM Sales.Customer WHERE TerritoryID = %(territory)s", {"territory": 1}),
]
results = await parallel_queries(pool, queries, max_concurrent=3)
for i, rows in enumerate(results, 1):
print(f"Query {i}: {len(rows)} rows")
finally:
await pool.close()
if __name__ == "__main__":
asyncio.run(main())
每個查詢同時執行,輸出會回報每筆回傳的列數:
Query 1: 32 rows
Query 2: 43 rows
Query 3: 31465 rows
Query 4: 3520 rows
要載入大量資料列時,不要採用並行處理。 如 Bulk copy 所述,請改為減少來回次數。
串流傳送大量結果
利用 OFFSET/FETCH 分頁從大型結果集產生資料列,以保持記憶體使用量的限制。 這個範例是建立在 非同步連線池 的部分,所以請將該段的程式碼存為 , db.py 並在旁邊執行這個範例。
用於大型資料集的非同步產生器
實作一個非同步產生器函式,能隨時擷取結果頁面,讓呼叫者能在不將所有資料載入記憶體的情況下,遍歷大型資料集。
import asyncio
from db import AsyncDatabasePool # The pool class from the Async connection pool section
async def stream_results(pool: AsyncDatabasePool, query: str,
params: dict = None, chunk_size: int = 1000):
"""Stream query results as async generator."""
offset = 0
while True:
paged_query = f"""
{query}
ORDER BY (SELECT NULL)
OFFSET %(offset)s ROWS
FETCH NEXT %(limit)s ROWS ONLY
"""
chunk_params = {**(params or {}), "offset": offset, "limit": chunk_size}
results = await pool.execute(paged_query, chunk_params)
if not results:
break
for row in results:
yield row
offset += chunk_size
# Allow event loop to process other tasks
await asyncio.sleep(0)
# Usage
async def main():
pool = AsyncDatabasePool(
"Server=<server>.database.windows.net;Database=<database>;Authentication=ActiveDirectoryDefault;Encrypt=yes",
pool_size=5,
)
await pool.initialize()
try:
processed = 0
async for order in stream_results(
pool, "SELECT SalesOrderID FROM Sales.SalesOrderHeader", chunk_size=500
):
processed += 1
print(f"Streamed {processed} orders in chunks of 500.")
finally:
await pool.close()
if __name__ == "__main__":
asyncio.run(main())
產生器一次只取得一頁,因此無論結果集多大,記憶體都會保持有界:
Streamed 31465 orders in chunks of 500.
最佳做法
運用這些指引,保持非同步模式的安全與效率。
適當的執行器規模設定
執行緒池的大小應該與 I/O 受限的工作負載相匹配,而非僅是 CPU 數量。
import os
# Rule of thumb: 2-4x CPU cores for I/O-bound database work
cpu_count = os.cpu_count() or 4
pool_size = cpu_count * 2
db_executor = ThreadPoolExecutor(max_workers=pool_size)
優雅的關機
停止任務執行程式,排空進行中的工作,然後依序關閉資源池和執行器。
async def graceful_shutdown(pool: AsyncDatabasePool, runner: AsyncTaskRunner):
"""Gracefully shut down all async components."""
# Stop accepting new tasks
runner.stop()
# Wait for running tasks to complete
await asyncio.sleep(2)
# Close database pool
await pool.close()
# Shutdown executor
db_executor.shutdown(wait=True)
錯誤處理
僅在暫時性故障時重試,並採用具延遲上限與隨機抖動的退避策略。 重用is_transient_error重試邏輯與連線韌性中的分類器,讓錯誤的認證資訊或語法錯誤等永久性失敗立即失敗,而不是重試。
import asyncio
import random
import mssql_python
# Reuse is_transient_error() from the Retry logic article.
async def resilient_query(pool: AsyncDatabasePool, query: str,
params: dict = None, retries: int = 3,
base_delay: float = 1.0, max_delay: float = 30.0) -> list:
"""Execute a query, retrying only on transient failures."""
for attempt in range(retries + 1):
try:
return await pool.execute(query, params)
except mssql_python.Error as e:
if not is_transient_error(e) or attempt == retries:
raise
delay = min(base_delay * (2 ** attempt), max_delay)
delay *= 0.5 + random.random() # Add jitter
await asyncio.sleep(delay)