Aszinkron minták az mssql-pythonnal

Az mssql-python illesztőprogram szinkron I/O-t használ, és nem nyújt natív async/await támogatást. A natív aszinkron szerepel az illesztőprogram ütemtervében. Addig is, az alábbi megoldások segítségével integrálhatod az mssql-pythont asszinkron alkalmazásokkal:

  • Thread pool végrehajtók blokkoló hívások leterítésére.
  • Aszinkron burkolófüggvények szinkron műveletekhez.
  • Integráció aszinkron keretrendszerekkel, mint a FastAPI.

Megjegyzés:

A cikkben bemutatott minták a ThreadPoolExecutor elemet használják a szinkron mssql-python hívások háttérszálakban történő futtatására. Ez a megközelítés növeli a szálakozási túlterhelést a natív aszinkron illesztőszerekhez képest. Az I/O-hoz kötött adatbázis-terheléseknél a többletköltség általában elfogadható.

Mikor kell használni aszinkron mintákat

A szál-pool megközelítés jól működik, ha:

  • Az alkalmazásod már használ asyncio (például FastAPI, aiohttp vagy Discord botokat), és integrálni kell az adatbázis-hívásokat anélkül, hogy blokkolnád az eseményciklust.
  • Az adatbázis lekérdezések I/O-hoz kötöttek, nem CPU-hoz kötöttek. A szál pool lehetővé teszi, hogy az event loop kezelje a többi kérést, miközben a Microsoft SQL-re vár.
  • Önnek közepes mértékű párhuzamossága van (több tíz egyidejű lekérdezés, nem több ezer).

Tiszta szinkron alkalmazásoknál hagyjuk ki ezeket a mintákat. Használd közvetlenül az illesztőprogramot szinkron kóddal a közvetlen, alacsonyabb terheléses végrehajtáshoz.

Szálkészlet-végrehajtó tervezési minta

Az alábbi példák bemutatják, hogyan lehet a szinkron mssql-python hívásokat egy ThreadPoolExecutor-ba becsomagolni, hogy a asyncio használhatók legyenek.

Egyszerű aszinkron burkoló

Létrehozz egy egyszerű segítő függvényt, amely szinkron mssql-python műveleteket hajt végre a szál poolban, és megvárja az eredményt.

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

Megjegyzés:

Ez a példa a két lekérdezést egymás után várja, így azok egymás után futnak. A await kulcsszó felszabadítja az eseményhurkot, hogy más feladatokat futtasson le, miközben az egyes lekérdezések várak, de nem fedi át ezeket a két lekérdezést egymással. Ahhoz, hogy független lekérdezéseket egyszerre futtassunk, ütemezzük őket egymással , asyncio.gatherahogy az Async kapcsolati pool szakaszban látható:

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

Aszinkron kapcsolati pool

Ez a rész bemutatja, hogyan lehet aszinkronbarát wrappert építeni az mssql-python beépített kapcsolati poolja köré alkalmazások alkalmazásához asyncio .

Megjegyzés:

Az mssql-python illesztőprogram beépített kapcsolati poolingot tartalmaz. Az itt bemutatott aszinkron készlet a szinkron kapcsolat-készletből származó kapcsolatokat aszinkron kontextuskezelőkbe csomagolja, hogy asyncio alkalmazásokban használhatók legyenek. Nem kell egy egyéni poolt kezelned, ha az mssql-python-t csak egy szálkészlet-végrehajtóból hívod.

Pufferelt aszinkron adatbázisosztály

Építs egy újrahasználható aszinkron kapcsolati pool osztályt, amely egy mssql-python kapcsolatok sorát kezeli, és asszinkron módszereket biztosít a lekérdezések végrehajtásához.

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

Integráció a FastAPI-val

A FastAPI kontextuskezelője lifespan automatikusan kezeli a pool inicializálását és a tisztítást.

Aszinkron FastAPI mssql-python használatával

Ez a példa az Async connection pool szekcióra épül. Mentsd el a szekció kódját egy fájlba, amelynek címe db.py, majd hozd létre a következő FastAPI alkalmazást egy mellette lévő main.py fájlban. A poolpélda a demót if __name__ == "__main__": használatával védi ki, így db.py importálásakor nem fut le a demó. Ez az alkalmazás a lifespan kontextuskezelőt használja a készlet indításkori inicializálására, majd leállításkor annak felszabadítására.

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
    }

Telepítsd a függőségeket, és futtasd az alkalmazást egy ASGI-kiszolgálóval, például az Uvicornnal. Futtassa ezt a parancsot abban a mappában, amely a main.py és db.py elemeket tartalmazza:

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

A szerver futása közben nyisd meg a http://127.0.0.1:8000/products, http://127.0.0.1:8000/products/1 vagy http://127.0.0.1:8000/stats elemeket az egyes végpontok meghívásához.

Háttérfeladatok

Indítsa el időszakos adatbázis-műveleteket ütemezésen anélkül, hogy blokkolnád az alkalmazás eseményciklusát.

Aszinkron háttérmunkás

Valósítson meg egy feladatfutót, amely meghatározott időközönként hajt végre regisztrált adatbázis-műveleteket, így elkerülve a párhuzamos futásokat is. Ez a példa az Async kapcsolatkészlet szakaszra épül, ezért mentse a szakasz kódját db.py néven. Ezután mentsd el a következő kódot worker.py néven mellé. Úgy konfigurálja a naplózást, hogy minden futás naplózza az eredményét.

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

Futtatd a munkást:

python worker.py

Minden regisztrált feladat naplózik, amikor fut, így néhány másodpercenként ismétlődő kimenetet láthatsz:

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

Egyidejű lekérdezés végrehajtása

Használd a asyncio.gather elemet szemafórral az egyidejűleg futó lekérdezések számának korlátozására. A példák ebben a szakaszban az Async kapcsolati pool szakaszra épülnek, így mentsd el az adott szekció kódját, db.py és futtatd az egyes példákat a saját fájljában mellette.

Párhuzamos lekérdezések szemafóral

Több lekérdezést hajtson végre egyidejűleg, miközben szemafort használ az egyidejű műveletek számának korlátozására, így megelőzve a szálkészlet telítődését.

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

Minden lekérdezés egyidejűleg fut, és a kimenet jelzi, hogy mindegyik hány sort adott vissza:

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

Nagy mennyiségű sor betöltéséhez ne a párhuzamosítást alkalmazd. Használj kevesebb oda-vissza utat, ahogy azt a Bulk copy is leírja.

Nagy eredmények streamelése

Nagy eredményhalmazokból származó sorokat kapunk oldalozással OFFSET/FETCH a memóriahasználat korlátozására. Ez a példa az Async kapcsolati pool szakaszra épül, ezért mentsd el az adott szekció kódját, db.py és futtatd ezt a példát mellette.

Aszinkron generátor nagy adathalmazokhoz

Valósíts meg egy aszinkron generátor függvényt, amely igény szerint hozza le az eredményoldalakat, lehetővé téve a hívók számára, hogy nagy adathalmazon iteráljanak anélkül, hogy mindent betöltenének a memóriába.

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

A generátor egyszerre egy oldalt kér le, így a memóriahasználat korlátozott marad, függetlenül attól, hogy mekkora az eredményhalmaz:

Streamed 31465 orders in chunks of 500.

Bevált gyakorlatok

Alkalmazd ezeket az irányelveket, hogy az aszinkron minták biztonságosak és hatékonyak legyenek.

Megfelelő végrehajtó méretezés

A szálkészlet méretét az I/O-hoz kötött terheléshez igazítsd, ne kizárólag a CPU-k száma alapján.

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)

Kecses leállítás

Állítsd le a feladatütemezőt, hagyd kifutni a folyamatban lévő feladatokat, majd zárd le az erőforráskészletet és a végrehajtót ebben a sorrendben.

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)

Hibakezelés

Csak átmeneti hibák esetén próbáld újra, és korlátozott késleltetéssel és jitterrel zárj vissza. Használd újra a is_transient_error osztályozót a újrapróbálkozási logika és kapcsolatellenállóság részből, hogy az olyan tartós hibák, mint a hibás hitelesítő adatok vagy a szintaktikai hibák, azonnal meghiúsuljanak ahelyett, hogy újrapróbálkoznának.

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)