Poznámka:
Přístup k této stránce vyžaduje autorizaci. Můžete se zkusit přihlásit nebo změnit adresáře.
Přístup k této stránce vyžaduje autorizaci. Můžete zkusit změnit adresáře.
Ovladač mssql-python používá synchronní I/O a neposkytuje nativní async/await podporu. Nativní podpora asynchronních operací je součástí plánu vývoje ovladače. Do té doby můžete integrovat mssql-python s asynchronními aplikacemi pomocí těchto obcházecích vzorců:
- Spouštěče fondu vláken pro přesunutí blokujících volání.
- Asynchronní obaly kolem synchronních operací.
- Integrace s asynchronními frameworky jako FastAPI.
Note
Postupy v tomto článku používají ThreadPoolExecutor ke spouštění synchronních volání mssql-python ve vláknech na pozadí. Tento přístup způsobuje vyšší režii spojenou s vlákny ve srovnání s nativními asynchronními drivery. U databázových zátěží vázaných na I/O je režie obvykle přijatelná.
Kdy použít asynchronní vzory
Přístup thread-pool funguje dobře, když:
- Vaše aplikace už používá
asyncio(například FastAPI, aiohttp nebo Discord boty) a musíte integrovat volání databáze bez blokování událostní smyčky. - Databázové dotazy jsou vázané na I/O, nikoli na CPU. Thread pool umožňuje event loopu zpracovávat další požadavky během čekání na Microsoft SQL.
- Máte střední souběžnost (desítky souběžných dotazů, ne tisíce).
Pro čistě synchronní aplikace tyto vzory přeskočte. Používejte driver přímo v synchronním kódu pro přímé provádění s nižší režijní zátěží.
Vzor executoru thread pool
Následující příklady ukazují, jak zabalit synchronní volání mssql-python do ThreadPoolExecutor pro použití s asyncio.
Základní asynchronní obal
Vytvořte jednoduchou pomocnou funkci, která vykonává synchronní operace mssql-python v poolu vláken a čeká na výsledek.
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
Tento příklad čeká na dva dotazy jeden po druhém, takže běží postupně. Klíčové await slovo uvolňuje smyčku událostí, aby mohla spouštět další úkoly, zatímco každý dotaz čeká, ale tyto dva dotazy se nepřekrývají. Pro současné spuštění nezávislých dotazů je naplánujte společně s asyncio.gather, jak je uvedeno v sekci Async connection pool :
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"),
)
Asynchronní fond připojení
Tato sekce ukazuje, jak vytvořit asynchronický wrapper kolem vestavěného poolu připojení mssql-python pro použití v asyncio aplikacích.
Note
Ovladač mssql-python obsahuje vestavěné poolování připojení. Asynchronní pool zobrazený zde obaluje synchronní poolovaná spojení s asynchronními správci kontextu pro použití v asyncio aplikacích. Nemusíte spravovat vlastní pool, pokud voláte mssql-python pouze z executoru thread poolu.
Třída asynchronní databáze ve skupině
Vytvořte opakovaně použitelnou třídu asynchronního poolu připojení, která spravuje frontu připojení mssql-python a poskytuje asynchronní metody pro provádění dotazů.
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())
Integrace s FastAPI
Správce kontextu FastAPI lifespan automaticky zajišťuje inicializaci a čištění poolu.
Asynchronní FastAPI s mssql-python
Tento příklad navazuje na sekci Async connection pool . Uložte kód této sekce do souboru s názvem db.py, a poté vytvořte následující aplikaci FastAPI v souboru pojmenovaném main.py vedle ní. Příklad poolu chrání své demo pomocí if __name__ == "__main__":, takže import db.py dema nespouští. Tato aplikace používá lifespan správce kontextu k inicializaci poolu při spuštění a jeho vyčištění při vypnutí.
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
}
Nainstalujte závislosti a spusťte aplikaci na ASGI serveru, například Uvicorn. Spusť tento příkaz ze složky, která obsahuje main.py a db.py:
pip install fastapi uvicorn mssql-python
uvicorn main:app --reload
Když je server spuštěný, otevřete http://127.0.0.1:8000/products, http://127.0.0.1:8000/products/1 nebo http://127.0.0.1:8000/stats, abyste zavolali jednotlivé koncové body.
Úlohy na pozadí
Spouštějte periodické databázové operace podle plánu bez blokování smyčky událostí aplikace.
Asynchronní pracovník na pozadí
Implementujte task runner, který vykonává registrované databázové operace v určených intervalech, čímž se zabraňuje duplicitním souběžným běhům. Tento příklad vychází ze sekce Async connection pool, takže kód této sekce uložte jako .db.py Poté uložte následující kód jako worker.py vedle něj. Konfiguruje logování tak, aby každý běh hlásil svůj výsledek.
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())
Běžte pracovníka:
python worker.py
Každý registrovaný úkol se zaznamenává při svém běhu, takže vidíte opakující se výstup každých pár sekund:
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
Souběžné provádění dotazů
Použijte asyncio.gather se semaforem k omezení počtu současně běžících dotazů. Příklady v této sekci navazují na sekci Async connection pool, takže kód z této sekce uložte jako db.py a každý příklad spusťte v samostatném souboru ve stejném adresáři.
Paralelní dotazy se semaforem
Provádět více dotazů současně při použití semaforu k omezení počtu současných operací a zabránit tak nasycení vlákna.
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())
Každý dotaz běží současně a výstup hlásí, kolik řádků každý z nich vrátil:
Query 1: 32 rows
Query 2: 43 rows
Query 3: 31465 rows
Query 4: 3520 rows
Pro načítání velkého množství řádků nepoužívejte souběžné zpracování. Místo toho používejte méně cest tam a zpět, jak je popsáno v hromadném kopírování.
Streamování velkých výsledků
Načítejte řádky z velkých sad výsledků pomocí stránkování OFFSET/FETCH, aby využití paměti zůstalo omezené. Tento příklad navazuje na část Asynchronní fond připojení, takže kód z této části uložte jako db.py a poté vedle něj spusťte i tento příklad.
Asynchronní generátor pro velké datové sady
Implementujte asynchronní generátorovou funkci, která na vyžádání načítá stránky s výsledky a volajícím umožňuje procházet rozsáhlé datové sady, aniž by bylo nutné načíst vše do paměti.
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())
Generátor načítá jednu stránku po druhé, takže paměť zůstává omezená bez ohledu na velikost výsledné množiny:
Streamed 31465 orders in chunks of 500.
Osvědčené postupy
Aplikujte tyto pokyny, abyste udrželi asynchronní vzory bezpečné a efektivní.
Správná velikost exekutoru
Velikost poolu vláken tak, aby odpovídala zátěžím vázaným na I/O, ne jen počtu 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)
Bezpečné vypnutí
Zastavte spouštěč úloh, nechte doběhnout rozpracované úlohy a poté v tomto pořadí uzavřete fond a vykonavatele.
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)
Zpracování chyb
Zkoušejte to znovu jen při přechodných selháních, a pak zpomalte s limitovaným zpožděním a chvěním. Znovu použijte is_transient_error klasifikátor z logiky opakování a odolnosti připojení, aby trvalá selhání, jako jsou chybné přihlašovací údaje nebo chyby syntaxe, selhala okamžitě místo opakovaných pokusů.
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)