Nota:
El acceso a esta página requiere autorización. Puede intentar iniciar sesión o cambiar directorios.
El acceso a esta página requiere autorización. Puede intentar cambiar los directorios.
El controlador mssql-python usa E/S síncrona y no ofrece soporte nativo async/await . La asincronía nativa está en la hoja de ruta del controlador. Hasta entonces, puedes integrar mssql-python con aplicaciones asíncronas usando estos patrones alternativos:
- Ejecutores de pool de hilos para descargar llamadas bloqueadoras.
- Envolturas asíncronas alrededor de operaciones síncronas.
- Integración con frameworks asincrónicos como FastAPI.
Note
Los patrones de este artículo usan ThreadPoolExecutor para ejecutar llamadas síncronas de mssql-python en hilos en segundo plano. Este enfoque añade sobrecarga de hilos en comparación con los controladores asincrónicos nativos. Para cargas de trabajo de bases de datos vinculadas a E/S, la sobrecarga suele ser aceptable.
Cuándo usar patrones asincrónicos
El enfoque de pool de hilos funciona bien cuando:
- Tu aplicación ya utiliza
asyncio(por ejemplo, bots FastAPI, aiohttp o Discord) y necesitas integrar llamadas a bases de datos sin bloquear el bucle de eventos. - Las consultas a la base de datos están limitadas por E/S, no por CPU. El pool de hilos permite que el bucle de eventos gestione otras solicitudes mientras espera Microsoft SQL.
- Tienes una concurrencia moderada (decenas de consultas concurrentes, no miles).
Para aplicaciones puramente síncronas, salta estos patrones. Usa el controlador directamente con código síncrono para una ejecución directa y con menor sobrecarga.
Patrón de ejecutor con grupo de hilos
Los siguientes ejemplos muestran cómo envolver llamadas mssql-python síncronas en un ThreadPoolExecutor para su uso con asyncio.
Envoltorio asíncrono básico
Crea una función auxiliar sencilla que ejecute operaciones sincrónicas mssql-python en el pool de hilos y espere el resultado.
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
Este ejemplo espera que las dos consultas se repitan una tras otra, por lo que se ejecutan de forma secuencial. La await palabra clave libera el bucle de eventos para ejecutar otras tareas mientras cada consulta espera, pero no solapa estas dos consultas entre sí. Para ejecutar consultas independientes al mismo tiempo, prográmalas junto con asyncio.gather, como se muestra en la sección Grupo de conexiones asíncronas:
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 conexión asíncrona
Esta sección muestra cómo crear un envoltorio compatible con operaciones asíncronas sobre el pool de conexiones integrado de mssql-python para usarlo en aplicaciones asyncio.
Note
El controlador mssql-python incluye agrupación de conexiones integrada. El pool asíncrono que se muestra aquí envuelve conexiones agrupadas síncronas con gestores de contexto asíncronos para su uso en asyncio aplicaciones. No necesitas gestionar un pool personalizado si solo llamas a mssql-python desde un ejecutor de pool de hilos.
Clase de base de datos asíncrona agrupada
Construye una clase de pool de conexiones asíncronas reutilizable que gestione una cola de conexiones mssql-python y proporcione métodos asíncronos para la ejecución de consultas.
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())
Integración con FastAPI
El gestor de contexto de lifespan FastAPI gestiona automáticamente la inicialización y limpieza del pool.
FastAPI asíncrona con mssql-python
Este ejemplo se basa en la sección Grupo de conexiones asíncronas. Guarda el código de esa sección en un archivo llamado db.py, y luego crea la siguiente aplicación FastAPI en un archivo con el nombre main.py de al lado. El ejemplo del pool protege su demo con if __name__ == "__main__":, por lo que importar db.py no ejecuta la demo. Esta aplicación utiliza el lifespan gestor de contexto para inicializar el pool al arrancar y limpiarlo al apagarlo.
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
}
Instala las dependencias y ejecuta la app con un servidor ASGI como Uvicorn. Ejecuta este comando desde la carpeta que contiene main.py y db.py:
pip install fastapi uvicorn mssql-python
uvicorn main:app --reload
Con el servidor en ejecución, abre http://127.0.0.1:8000/products, http://127.0.0.1:8000/products/1 o http://127.0.0.1:8000/stats para llamar a cada extremo.
Tareas en segundo plano
Ejecuta operaciones periódicas de base de datos en un calendario sin bloquear el bucle de eventos de la aplicación.
Trabajador asíncrono en segundo plano
Implementa un ejecutor de tareas que ejecute operaciones registradas de base de datos en intervalos especificados, evitando ejecuciones concurrentes duplicadas. Este ejemplo se basa en la sección del pool de conexiones asíncronas , así que guarda el código de esa sección como db.py. Luego, guarda el siguiente código como worker.py junto a este. Configura el registro para que cada ejecución informe su resultado.
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())
Ejecuta el proceso de trabajo:
python worker.py
Cada tarea registrada registra cuando se ejecuta, así que ves la salida repetida cada pocos segundos:
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
Ejecución concurrente de consultas
Úsalo asyncio.gather con un semáforo para limitar el número de consultas que se ejecutan simultáneamente. Los ejemplos de esta sección se basan en la sección del pool de conexiones asíncronas , así que guarda el código de esa sección como db.py y ejecuta cada ejemplo en su propio archivo junto a él.
Consultas paralelas con semáforo
Ejecuta varias consultas simultáneamente utilizando un semáforo para limitar el número de operaciones simultáneas, evitando la saturación del pool de hilos.
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())
Cada consulta se ejecuta simultáneamente, y la salida informa cuántas filas ha devuelto cada una:
Query 1: 32 rows
Query 2: 43 rows
Query 3: 31465 rows
Query 4: 3520 rows
Para cargar grandes volúmenes de filas, no recurras a la concurrencia. Utiliza menos viajes de ida y vuelta, como se describe en la copia Bulk.
Resultados grandes en streaming
Obtener filas de grandes conjuntos de resultados usando OFFSET/FETCH la paginación para mantener el uso de memoria acotado. Este ejemplo se basa en la sección grupo de conexiones asíncronas, así que guarda el código de esa sección como db.py y luego ejecuta este ejemplo junto a él.
Generador asíncrono para grandes conjuntos de datos
Implementa una función generadora asincrónica que recupere páginas de resultados bajo demanda, lo que permite a los usuarios iterar sobre grandes conjuntos de datos sin cargarlo todo en memoria.
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())
El generador recupera una página a la vez, por lo que la memoria permanece acotada sin importar lo grande que sea el conjunto de resultados:
Streamed 31465 orders in chunks of 500.
procedimientos recomendados
Aplica estas pautas para mantener los patrones asincrónicos seguros y eficientes.
Dimensionado correcto del ejecutor
Dimensiona el pool de hilos para que coincida con las cargas de trabajo limitadas a E/S, no solo con el número 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)
Apagado ordenado
Detener el procesador de tareas, dejar que termine el trabajo en curso y después cerrar el grupo y el ejecutor, en ese orden.
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)
Gestión de errores
Reintenta solo ante errores transitorios y aplica un retroceso con un retraso máximo y variación aleatoria. Reutiliza el is_transient_error clasificador de la lógica de reintento y la resiliencia de la conexión para que los fallos permanentes, como credenciales incorrectas o errores de sintaxis, fallen de inmediato en lugar de volver a intentarlo.
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)