Wzorce asynchroniczne w mssql-python

Sterownik mssql-python korzysta z synchronicznego wejścia/wyjścia i nie zapewnia natywnego async/await wsparcia. Natywna obsługa asynchroniczności znajduje się w planie rozwoju sterownika. Do tego czasu możesz zintegrować mssql-python z aplikacjami asynchronicznymi, korzystając z następujących wzorców obejść:

  • Wykonawcy puli wątków do odciążania wywołań blokujących.
  • Asynchroniczne opakowania dla operacji synchronicznych.
  • Integracja z frameworkami asynchronicznymi, takimi jak FastAPI.

Note

Wzorce przedstawione w tym artykule wykorzystują ThreadPoolExecutor do uruchamiania synchronicznych wywołań mssql-python w wątkach w tle. To podejście zwiększa narzut wątkowy w porównaniu do natywnych sterowników asynchronicznych. W przypadku obciążeń bazodanowych związanych z I/O, narzut jest zazwyczaj akceptowalny.

Kiedy stosować wzory asynchroniczne

Podejście z pulą wątków sprawdza się dobrze, gdy:

  • Twoja aplikacja już korzysta z asyncio (np. FastAPI, aiohttp lub botów Discord) i musisz integrować wywołania bazy danych bez blokowania pętli zdarzeń.
  • Zapytania bazowe są ograniczone względem I/O, a nie CPU. Pula wątków pozwala pętli zdarzeń obsługiwać inne żądania podczas oczekiwania na Microsoft SQL.
  • Masz umiarkowaną równocześność (dziesiątki jednoczesnych zapytań, nie tysiące).

W przypadku czysto synchronicznych zastosowań pomiń te wzorce. Używaj sterownika bezpośrednio z kodem synchronicznym, aby zapewnić bezpośrednie wykonywanie przy mniejszym narzucie.

Wzorzec wykonawcy puli wątków

Poniższe przykłady pokazują, jak opakować synchroniczne wywołania mssql-python w elemencie ThreadPoolExecutor do użycia z asyncio.

Podstawowe opakowanie asynchroniczne

Stwórz prostą funkcję pomocniczą, która wykonuje synchroniczne operacje mssql-python w puli wątków i czeka na wynik.

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

Ten przykład czeka na dwa zapytania jedno po drugim, więc uruchamiają się one sekwencyjnie. Słowo kluczowe await zwalnia pętlę zdarzeń, umożliwiając wykonywanie innych zadań, gdy każde zapytanie czeka na wynik, ale nie powoduje równoległego wykonywania tych dwóch zapytań względem siebie. Aby uruchamiać niezależne zapytania jednocześnie, zaplanuj je razem za pomocą asyncio.gather, jak pokazano w sekcji Asynchroniczna pula połączeń:

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"),
)

Asynchroniczna pula połączeń

Ta sekcja pokazuje, jak utworzyć otoczkę przystosowaną do operacji asynchronicznych, opartą na wbudowanej puli połączeń mssql-python, do użytku w aplikacjach asyncio.

Note

Sterownik mssql-python zawiera wbudowaną pulę połączeń. Pokazana tutaj pula asynchroniczna opakowuje synchroniczne połączenia z puli w asynchroniczne menedżery kontekstu do użycia w aplikacjach asyncio. Nie musisz zarządzać własną pulą, jeśli wywołujesz mssql-python wyłącznie w executorze puli wątków.

Klasa puli asynchronicznej bazy danych

Zbuduj ponownie używalną klasę asynchronicznej puli połączeń, która zarządza kolejką połączeń mssql-python i zapewnia asynchroniczne metody wykonywania zapytań.

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())

Integracja z FastAPI

Menedżer kontekstu lifespan w FastAPI automatycznie obsługuje inicjalizację oraz czyszczenie puli.

Async FastAPI z mssql-python

Ten przykład opiera się na sekcji puli połączeń asynchronicznych . Zapisz kod tej sekcji w pliku o nazwie db.py, a następnie stwórz następującą aplikację FastAPI w pliku o nazwie main.py obok. Przykład puli zabezpiecza swój kod demonstracyjny instrukcją if __name__ == "__main__":, dzięki czemu zaimportowanie db.py nie uruchamia kodu demonstracyjnego. Ta aplikacja używa menedżera lifespan kontekstu do inicjalizacji puli przy uruchomieniu i czyszczenia go po wyłączeniu.

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
    }

Zainstaluj zależności i uruchom aplikację na serwerze ASGI, takim jak Uvicorn. Wykonaj to polecenie z folderu, który zawiera main.py i db.py:

pip install fastapi uvicorn mssql-python
uvicorn main:app --reload

Po uruchomieniu serwera otwórz http://127.0.0.1:8000/products, http://127.0.0.1:8000/products/1 lub http://127.0.0.1:8000/stats, aby wywołać każdy punkt końcowy.

Zadania w tle

Wykonuj okresowe operacje bazowe według harmonogramu bez blokowania pętli zdarzeń aplikacji.

Asynchroniczny proces roboczy działający w tle

Zaimplementuj mechanizm uruchamiania zadań, który wykonuje zarejestrowane operacje na bazie danych w określonych odstępach czasu, zapobiegając ich jednoczesnemu, wielokrotnemu uruchamianiu. Ten przykład opiera się na sekcji puli połączeń asynchronicznych, więc zapisz kod tej sekcji jako .db.py Następnie zapisz poniższy kod jako worker.py obok niego. Konfiguruje rejestrowanie tak, aby każde uruchomienie raportowało swój wynik.

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())

Uruchom pracownika:

python worker.py

Każde zarejestrowane zadanie rejestruje moment, więc co kilka sekund widzisz powtarzające się wyniki:

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

Wykonywanie zapytań współbieżnych

Używaj asyncio.gather z semaforem, aby ograniczyć liczbę jednoczesnych zapytań. Przykłady w tej sekcji bazują na sekcji Asynchroniczna pula połączeń, więc zapisz kod z tej sekcji jako db.py i uruchamiaj każdy przykład w osobnym pliku obok niego.

Zapytania równoległe z użyciem semaforu

Wykonaj wiele zapytań jednocześnie, używając semaforu, aby ograniczyć liczbę operacji równocześnie, zapobiegając nasyceniu puli wątków.

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żde zapytanie jest wykonywane równolegle, a wyjście raportuje, ile wierszy każde z nich zwróciło:

Query 1: 32 rows
Query 2: 43 rows
Query 3: 31465 rows
Query 4: 3520 rows

Aby załadować duże wolumeny wierszy, nie stosuj współbieżności. Zamiast tego używaj mniejszej liczby podróży w obie strony, jak opisano w kopii zbiorczej.

Streaming dużych wyników

Otrzymuj wiersze z dużych zbiorów wyników za pomocą OFFSET/FETCH paginacji, aby utrzymać ograniczenie wykorzystania pamięci. Ten przykład bazuje na sekcji Asynchroniczna pula połączeń, więc zapisz kod z tej sekcji jako db.py, a następnie uruchom ten przykład obok niego.

Generator asynchroniczny dla dużych zbiorów danych

Zaimplementuj funkcję generatora asynchronicznego, która pobiera strony wyników na żądanie, pozwalając wywołującym iterować duże zbiory danych bez ładowania wszystkiego do pamięci.

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())

Generator pobiera po jednej stronie naraz, więc pamięć pozostaje ograniczona niezależnie od wielkości zbioru wyników:

Streamed 31465 orders in chunks of 500.

Najlepsze rozwiązania

Stosuj te wytyczne, aby wzorce asynchroniczne były bezpieczne i efektywne.

Właściwy dobór rozmiaru egzekutora

Rozmiar puli wątków odpowiada obciążeniom powiązanym z I/O, a nie tylko liczbie 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)

Łagodne zamykanie

Zatrzymaj mechanizm uruchamiania zadań, zaczekaj na zakończenie zadań będących w toku, a następnie zamknij po kolei pulę i obiekt wykonawczy.

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)

Obsługa błędów

Ponawiaj próby tylko w przypadku przejściowych błędów oraz stosuj opóźnienie z górnym limitem i losową zmiennością. Ponownie wykorzystaj klasyfikator is_transient_error z sekcji Logika ponawiania i odporność połączeń, aby trwałe błędy, takie jak nieprawidłowe poświadczenia lub błędy składni, były od razu zgłaszane zamiast ponawiania prób.

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)