mssql-python ile asenkron kalıplar

mssql-python sürücüsü senkron I/O kullanıyor ve yerel async/await destek sağlamıyor. Doğal asenkron destek, sürücünün yol haritasında yer alıyor. O zamana kadar, şu geçici çözümleri kullanarak mssql-python'u asenkron uygulamalarla entegre edebilirsiniz:

  • Bloklayıcı çağrıları aktarmak için iş parçacığı havuzu yürütücüleri.
  • Senkron işlemler için asenkron sarmalayıcılar.
  • FastAPI gibi asenkron çerçevelerle entegrasyon.

Note

Bu makaledeki desenler, arka plan iş parçacıklarında senkron mssql-python çağrılarını çalıştırmak için ThreadPoolExecutor kullanır. Bu yaklaşım, yerel asenkron sürücülere kıyasla iş parçacıllığı yükü artırıyor. I/O-bağlı veritabanı iş yükleri için ek yük genellikle kabul edilebilir.

Asenkron desenleri ne zaman kullanmalı?

İş parçacığı havuzu yaklaşımı şu durumlarda iyi çalışır:

  • Uygulamanız zaten asyncio kullanıyor (örneğin FastAPI, aiohttp veya Discord botları) ve veritabanı çağrılarını olay döngüsünü engellemeden entegre etmeniz gerekiyor.
  • Veritabanı sorguları I/O-bağlı, CPU ile sınırlı değil. İş parçacığı havuzu, Microsoft SQL’in yanıt vermesini beklerken olay döngüsünün diğer istekleri işlemesini sağlar.
  • Orta düzeyde bir eşzamanlılık söz konusu (binlerce değil, onlarca eşzamanlı sorgu).

Saf senkron uygulamalar için bu kalıpları atlayın. Doğrudan, daha düşük ek yükle yürütme için sürücüyü senkron kodla kullanın.

İş parçacığı havuzu yürütücü deseni

Aşağıdaki örnekler, eşzamanlı mssql-python çağrılarının ThreadPoolExecutor içine, asyncio ile kullanılmak üzere nasıl sarılacağını göstermektedir.

Temel asenkron sarmalayıcı

İş parçacığı havuzunda eşzamanlı mssql-python işlemlerini çalıştıran ve sonucu bekleyen basit bir yardımcı işlev oluşturun.

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

Bu örnek, iki sorgunun tamamlanmasını art arda bekler; bu nedenle sıralı olarak çalışırlar. await anahtar sözcüğü, her bir sorgu beklerken olay döngüsünü diğer görevleri çalıştırabilmesi için serbest bırakır, ancak bu iki sorgunun birbiriyle eşzamanlı yürütülmesini sağlamaz. Bağımsız sorguları aynı anda çalıştırmak için, bunları asyncio.gather bölümünde gösterildiği gibi ile birlikte programlayın:

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

Asenkron bağlantı havuzu

Bu bölüm, uygulamalarda asyncio kullanılmak üzere mssql-python'un yerleşik bağlantı havuzu etrafında asenkron dostu bir wrapper nasıl oluşturulacağını gösterir.

Note

mssql-python sürücüsü, yerleşik bağlantı havuzu içerir. Burada gösterilen asenkron havuz, eşzamanlı havuzlanmış bağlantıları, asyncio uygulamalarında kullanılmak üzere asenkron bağlam yöneticileriyle sarmalar. Yalnızca mssql-python'u bir iş parçacığı havuzu yürütücüsü içinden çağırıyorsanız, özel bir havuzu yönetmeniz gerekmez.

Havuzlu asenkron veritabanı sınıfı

mssql-python bağlantılarının kuyruğunu yöneten ve sorgu yürütme için asenkron yöntemler sağlayan yeniden kullanılabilir bir asenkron bağlantı havuzu sınıfı oluşturun.

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

FastAPI ile Entegrasyon

FastAPI'nin bağlam yöneticisi, lifespan havuz başlatma ve temizliği otomatik olarak yönetir.

mssql-python ile asenkronize FastAPI

Bu örnek, Async bağlantı havuzu bölümü üzerine inşa edilmiştir. O bölümün kodunu , adlı db.pybir dosyaya kaydedin ve ardından yanında adı main.py verilen bir dosyada aşağıdaki FastAPI uygulamasını oluşturun. Havuz örneği, demosunu if __name__ == "__main__": ile korur; böylece db.py içe aktarıldığında demo çalışmaz. Bu uygulama, başlangıçta havuzu başlatmak ve kapanışta temizlemek için lifespan bağlam yöneticisini kullanır.

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
    }

Bağımlılıkları yükleyip uygulamayı Uvicorn gibi bir ASGI sunucusuyla çalıştır. Bu komutu main.py ve db.py içeren klasörden çalıştırın:

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

Sunucu çalışırken, her uç noktayı açın http://127.0.0.1:8000/products, http://127.0.0.1:8000/products/1, veya http://127.0.0.1:8000/stats çağırın.

Arka plan görevleri

Uygulama olay döngüsünü engellemeden periyodik veritabanı işlemlerini bir program üzerinde çalıştırın.

Asenkron arka plan çalışanı

Kayıtlı veritabanı işlemlerini belirli aralıklarla çalıştıran bir görev çalıştırıcısı uygulayın, böylece eşzamanlı çalıştırmaların tekrarlanmasını önleyin. Bu örnek, Async bağlantı havuzu bölümünü temel alır; bu nedenle o bölümün kodunu db.py olarak kaydedin. Sonra, aşağıdaki kodu worker.py yanına kaydedin. Kayıtları yapılandırıyor, böylece her çalışma sonucunu bildiriyor.

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

İşçiyi çalıştırın:

python worker.py

Her kayıtlı görev, çalıştığında günlüğe kaydedilir; böylece her birkaç saniyede bir tekrarlanan çıktı görürsünüz:

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

Eşzamanlı sorgu yürütme

Aynı anda çalışan sorgu sayısını sınırlamak için semafor ile kullanın asyncio.gather . Bu bölümdeki örnekler Async bağlantı havuzu bölümü üzerine inşa ediliyor, bu yüzden o bölümün kodunu şu şekilde db.py kaydedin ve her örneği kendi dosyasında çalıştırın.

Semafor ile paralel sorgular

Eşzamanlı işlem sayısını sınırlamak için semafor kullanarak birden fazla sorgu aynı anda çalıştırın, böylece iş parçacığının doygunluğunu önleyin.

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

Her sorgu eşzamanlı çalışır ve çıktı, her birinin kaç satır döndürdüğünü bildirir:

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

Çok sayıda satır yüklemek için eşzamanlı işleme başvurmayın. Bunun yerine daha az gidiş-dönüş kullanın, Toplu kopya bölümünde açıklandığı gibi.

Büyük sonuçların akış olarak iletilmesi

Bellek kullanımını sınırlı tutmak için büyük sonuç kümelerinden satırları OFFSET/FETCH sayfalama kullanarak döndürün. Bu örnek Async bağlantı havuzu bölümü üzerine inşa ediliyor, bu yüzden o bölümün kodunu kaydedin db.py ve bu örneği yanına çalıştırın.

Büyük veri kümeleri için asenkron üreteç

Sonuç sayfalarını talep üzerine getiren, arayanların her şeyi belleğe yüklemeden büyük veri setleri üzerinde yineleme yapabilmelerini sağlayan bir asenkron üreteç fonksiyonu uygulayın.

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

Üretici her seferinde bir sayfa getirir, bu yüzden sonuç kümesi ne kadar büyük olursa olsun bellek sınırlı kalır:

Streamed 31465 orders in chunks of 500.

En iyi uygulamalar

Bu yönergeleri uygulayarak asenkron kalıpları güvenli ve verimli tutun.

Doğru yürütücü boyutlandırması

İş parçacığı havuzunu sadece CPU sayısına değil, I/O-bağlı iş yüklerine göre boyutlandırı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)

Sorunsuz kapatma

Görev koşucusunu durdur, uçuş sırasında işi boşalt, sonra havuzu ve yürütücüyü sırayla kapat.

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)

Hata yönetimi

Sadece geçici arızalarda tekrar deneyin, gecikme ve titreme ile geri çekilin. is_transient_error Yeniden deneme mantığı ve bağlantı dayanıklılığı bölümündeki sınıflandırıcıyı yeniden kullanın; böylece hatalı kimlik bilgileri veya söz dizimi hataları gibi kalıcı hatalar yeniden denenmek yerine hemen başarısız olur.

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)