Remarque
L’accès à cette page nécessite une autorisation. Vous pouvez essayer de vous connecter ou de modifier des répertoires.
L’accès à cette page nécessite une autorisation. Vous pouvez essayer de modifier des répertoires.
Le pilote mssql-python utilise des E/S synchrones et ne fournit pas de support natif async/await . L’asynchrone natif est sur la feuille de route des conducteurs. D’ici là, vous pouvez intégrer mssql-python avec des applications asynchrones en utilisant ces schémas de contournement :
- Exécuteurs de pool de threads pour décharger les appels de blocage.
- Des enveloppes asynchrones autour des opérations synchrones.
- Intégration avec des frameworks asynchrones comme FastAPI.
Note
Les modèles de cet article utilisent ThreadPoolExecutor pour exécuter des appels synchrones à mssql-python dans des fils d’exécution en arrière-plan. Cette approche ajoute un surcoût lié au multithreading par rapport aux pilotes asynchrones natifs. Pour les charges de travail liées aux bases de données liées aux E/S, la surcharge est généralement acceptable.
Quand utiliser des motifs asynchrones
L’approche par pool de threads fonctionne bien lorsque :
- Votre application utilise
asynciodéjà (par exemple, FastAPI, aiohttp ou Discord bots) et vous devez intégrer les appels de base de données sans bloquer la boucle d’événements. - Les requêtes de base de données sont limitées en E/S, pas en CPU. Le pool de threads permet à la boucle d’événements de gérer d’autres requêtes en attendant Microsoft SQL.
- Vous avez une concurrence modérée (des dizaines de requêtes simultanées, pas des milliers).
Pour les applications purement synchrones, ignorez ces motifs. Utilisez le pilote directement avec du code synchrone pour une exécution directe et à moindre surcharge.
Modèle d’exécuteur de pool de threads
Les exemples suivants montrent comment enrouler des appels mssql-python synchrones dans un ThreadPoolExecutor pour une utilisation avec asyncio.
Enveloppe asynchrone de base
Créez une fonction d’assistance simple qui exécute des opérations mssql-python synchrones dans le pool de threads et attend le résultat.
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
Cet exemple attend les deux requêtes l’une après l’autre, donc elles s’exécutent séquentiellement. Le await mot-clé libère la boucle d’événements pour exécuter d’autres tâches pendant que chaque requête attend, mais il ne recoupe pas ces deux requêtes. Pour exécuter des requêtes indépendantes en même temps, planifiez-les ensemble avec asyncio.gather, comme indiqué dans la section Pool de connexion asynchrone :
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"),
)
Pool de connexion asynchrone
Cette section explique comment créer un encapsulage adapté à l’asynchrone autour du pool de connexions intégré de mssql-python pour une utilisation dans les applications asyncio.
Note
Le pilote mssql-python inclut un pooling de connexions intégré. Le pool asynchrone présenté ici encapsule des connexions synchrones mises en pool avec des gestionnaires de contexte asynchrones pour être utilisées dans les applications asyncio. Vous n’avez pas besoin de gérer un pool personnalisé si vous n’appelez que mssql-python depuis un exécuteur de pool de threads.
Classe de base de données asynchrone en pool
Construisez une classe de pool de connexion asynchrone réutilisable qui gère une file d’attente de connexions mssql-python et fournit des méthodes asynchrones pour l’exécution des requêtes.
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())
Intégration avec FastAPI
Le gestionnaire de contexte de lifespan FastAPI gère automatiquement l’initialisation et le nettoyage des pools.
Async FastAPI avec mssql-python
Cet exemple s’appuie sur la section du pool de connexion asynchrone . Sauvegardez le code de cette section dans un fichier nommé db.py, puis créez l’application FastAPI suivante dans un fichier nommé main.py à côté. L’exemple du pool protège sa démo avec if __name__ == "__main__":, donc l’importation db.py ne fait pas tourner la démo. Cette application utilise le gestionnaire de lifespan contexte pour initialiser le pool au démarrage et le nettoyer à l’arrêt.
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
}
Installez les dépendances et lancez l’application avec un serveur ASGI comme Uvicorn. Exécutez cette commande depuis le dossier qui contient main.py et db.py:
pip install fastapi uvicorn mssql-python
uvicorn main:app --reload
Le serveur étant en cours d’exécution, ouvrez http://127.0.0.1:8000/products, http://127.0.0.1:8000/products/1 ou http://127.0.0.1:8000/stats pour appeler chaque point de terminaison.
Tâches en arrière-plan
Exécutez périodiquement des opérations de base de données selon un calendrier sans bloquer la boucle d’événements de l’application.
Travailleur asynchrone en arrière-plan
Implémentez un exécuteur de tâches qui exécute des opérations de base de données enregistrées à des intervalles spécifiés, empêchant les doublons des exécutions concurrentes. Cet exemple s’appuie sur la section du pool de connexions asynchrones , donc enregistrez le code de cette section sous db.pyforme de . Ensuite, enregistrez le code suivant sous worker.py dans le même dossier. Il configure la journalisation pour que chaque exécution rapporte son résultat.
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())
Faites courir l’ouvrier :
python worker.py
Chaque tâche enregistrée enregistre son exécution, donc vous voyez une sortie répétée toutes les quelques secondes :
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
Exécution concurrente de requêtes
Utilisez asyncio.gather avec un sémaphore pour limiter le nombre de requêtes exécutées simultanément. Les exemples de cette section s’appuient sur la section du pool de connexions asynchrones , donc enregistrez le code de cette section sous db.py et exécutez chaque exemple dans son propre fichier à côté.
Requêtes parallèles avec sémaphore
Exécutez plusieurs requêtes simultanément tout en utilisant un sémaphore pour limiter le nombre d’opérations simultanées, évitant ainsi la saturation du pool de threads.
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())
Chaque requête s’exécute simultanément, et la sortie indique combien de lignes chacune a renvoyées :
Query 1: 32 rows
Query 2: 43 rows
Query 3: 31465 rows
Query 4: 3520 rows
Pour charger de gros volumes de lignes, n’optez pas pour le parallélisme. Privilégiez plutôt moins d’allers-retours, comme décrit dans la copie en bloc.
Streaming de grands volumes de résultats
Obtenir des lignes à partir de grands ensembles de résultats en utilisant OFFSET/FETCH la pagination pour limiter l’utilisation de la mémoire. Cet exemple s’appuie sur la section Pool de connexions asynchrone, donc enregistrez le code de cette section sous db.py et exécutez ensuite cet exemple à côté de celui-ci.
Générateur asynchrone pour de grands ensembles de données
Implémentez une fonction génératrice asynchrone qui récupère les pages de résultats à la demande, permettant aux appelants d’itérer sur de grands ensembles de données sans tout charger en mémoire.
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())
Le générateur récupère une page à la fois, donc la mémoire reste bornée, peu importe la taille de l’ensemble de résultats :
Streamed 31465 orders in chunks of 500.
Bonnes pratiques
Appliquez ces conseils pour garder les schémas asynchrones sûrs et efficaces.
Dimensionnement approprié de l’exécuteur
Dimensionnez le pool de threads pour correspondre aux charges de travail liées aux E/S, pas seulement au nombre de 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)
Arrêt approprié
Arrêtez l’exécuteur de tâches, laissez s’achever les tâches en cours, puis fermez le pool et l’exécuteur dans cet ordre.
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)
Gestion des erreurs
Réessayer uniquement en cas d’échecs temporaires, en augmentant le délai entre les tentatives jusqu’à un maximum et en y ajoutant une part d’aléa. Réutilisez le is_transient_error classificateur issu de la logique de tentative et de la résilience de connexion afin que les défaillances permanentes comme de mauvaises identifiantes ou des erreurs de syntaxe échouent rapidement au lieu de réessayer.
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)