Catatan
Akses ke halaman ini memerlukan otorisasi. Anda dapat mencoba masuk atau mengubah direktori.
Akses ke halaman ini memerlukan otorisasi. Anda dapat mencoba mengubah direktori.
Driver mssql-python menggunakan I/O sinkron dan tidak memberikan dukungan asli async/await . Asinkron asli ada di peta jalan pengemudi. Sampai saat itu, Anda dapat mengintegrasikan mssql-python dengan aplikasi asinkron dengan menggunakan pola solusi ini:
- Eksekutor kumpulan utas untuk mengalihkan pemanggilan yang memblokir.
- Pembungkus asinkron di sekitar operasi sinkron.
- Integrasi dengan kerangka kerja asinkron seperti FastAPI.
Note
Pola dalam artikel ini menggunakan ThreadPoolExecutor untuk menjalankan panggilan mssql-python sinkron di utas latar belakang. Pendekatan ini menghasilkan overhead tambahan akibat penggunaan thread dibandingkan dengan driver asinkron asli. Untuk beban kerja database terikat I/O, overhead biasanya dapat diterima.
Kapan menggunakan pola asinkron
Pendekatan pool thread berfungsi dengan baik ketika:
- Aplikasi Anda sudah menggunakan
asyncio(misalnya, bot FastAPI, aiohttp, atau Discord) dan Anda perlu mengintegrasikan panggilan database tanpa memblokir perulangan peristiwa. - Kueri database terikat I/O, bukan terikat CPU. Pool thread memungkinkan event loop untuk menangani permintaan lain sambil menunggu respons dari Microsoft SQL.
- Anda memiliki tingkat konkurensi sedang (puluhan kueri secara bersamaan, bukan ribuan).
Untuk aplikasi sinkron murni, lewati pola ini. Gunakan driver secara langsung dengan kode sinkron untuk eksekusi langsung dan overhead yang lebih rendah.
Pola pelaksana kumpulan utas
Contoh berikut menunjukkan cara membungkus pemanggilan sinkron mssql-python dalam ThreadPoolExecutor agar dapat digunakan dengan asyncio.
Pembungkus asinkron dasar
Buat fungsi pembantu sederhana yang menjalankan operasi mssql-python sinkron di kumpulan utas dan menunggu hasilnya.
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
Contoh ini menunggu dua kueri satu demi satu, sehingga berjalan secara berurutan. Kata kunci await membuat loop peristiwa dapat menjalankan tugas lain saat masing-masing kueri menunggu, tetapi tidak membuat kedua kueri ini berjalan tumpang tindih satu sama lain. Untuk menjalankan kueri independen secara bersamaan, jadwalkan bersama dengan asyncio.gather, seperti yang ditunjukkan di bagian kumpulan koneksi Asinkron :
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"),
)
Kumpulan koneksi asinkron
Bagian ini menunjukkan cara membuat pembungkus ramah asinkron di sekitar kumpulan koneksi bawaan mssql-python untuk digunakan dalam asyncio aplikasi.
Note
Driver mssql-python menyertakan pengumpulan koneksi bawaan. Kumpulan asinkron yang ditampilkan di sini membungkus koneksi gabungan sinkron dengan pengelola konteks asinkron untuk digunakan dalam asyncio aplikasi. Anda tidak perlu mengelola pool kustom jika hanya memanggil mssql-python dari thread pool executor.
Kelas database asinkron yang dikumpulkan
Bangun kelas kumpulan koneksi asinkron yang dapat digunakan kembali yang mengelola antrean koneksi mssql-python dan menyediakan metode asinkron untuk eksekusi kueri.
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())
Integrasi dengan FastAPI
Manajer konteks FastAPI lifespan menangani inisialisasi dan pembersihan kumpulan secara otomatis.
FastAPI asinkron dengan mssql-python
Contoh ini dibangun di atas bagian kumpulan koneksi asinkron . Simpan kode bagian itu dalam file bernama db.py, lalu buat aplikasi FastAPI berikut dalam file bernama main.py di sebelahnya. Contoh pool melindungi kode demonya dengan if __name__ == "__main__":, sehingga mengimpor db.py tidak menjalankan demo. Aplikasi ini menggunakan lifespan manajer konteks untuk menginisialisasi pool saat memulai dan membersihkannya saat penghentian.
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
}
Instal dependensi dan jalankan aplikasi dengan server ASGI seperti Uvicorn. Jalankan perintah ini dari folder yang berisi main.py dan db.py:
pip install fastapi uvicorn mssql-python
uvicorn main:app --reload
Dengan server berjalan, buka http://127.0.0.1:8000/products, http://127.0.0.1:8000/products/1, atau http://127.0.0.1:8000/stats untuk memanggil setiap titik akhir.
Tugas latar belakang
Jalankan operasi database berkala sesuai jadwal tanpa memblokir perulangan peristiwa aplikasi.
Pekerja latar belakang asinkron
Terapkan pelaku tugas yang menjalankan operasi database terdaftar pada interval tertentu, mencegah eksekusi bersamaan duplikat. Contoh ini dibangun di bagian kumpulan koneksi Asinkron , jadi simpan kode bagian tersebut sebagai db.py. Kemudian, simpan kode berikut di worker.py sebelahnya. Ini mengonfigurasi pengelogan sehingga setiap eksekusi melaporkan hasilnya.
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())
Jalankan pekerja:
python worker.py
Setiap tugas terdaftar mencatat saat berjalan, sehingga Anda melihat output berulang setiap beberapa detik:
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
Eksekusi kueri bersamaan
Gunakan asyncio.gather dengan semaphore untuk membatasi jumlah kueri yang berjalan secara bersamaan. Contoh-contoh di bagian ini mengacu pada bagian pool koneksi asinkron, jadi simpan kode pada bagian tersebut sebagai db.py dan jalankan setiap contoh dalam file terpisah di direktori yang sama.
Kueri Paralel dengan Semaphore
Jalankan beberapa kueri secara bersamaan dengan menggunakan semaphore untuk membatasi jumlah operasi yang berjalan secara bersamaan, sehingga mencegah kumpulan utas menjadi jenuh.
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())
Setiap kueri berjalan secara bersamaan, dan output melaporkan berapa banyak baris yang ditampilkan masing-masing:
Query 1: 32 rows
Query 2: 43 rows
Query 3: 31465 rows
Query 4: 3520 rows
Untuk memuat baris dalam jumlah besar, jangan gunakan konkurensi. Gunakan lebih sedikit perjalanan pulang pergi, seperti yang dijelaskan dalam Salinan massal.
Streaming hasil dalam jumlah besar
Mengembalikan baris dari kumpulan hasil besar menggunakan paginasi OFFSET/FETCH agar penggunaan memori tetap terkendali. Contoh ini dibangun di atas bagian kumpulan koneksi Asinkron , jadi simpan kode bagian tersebut sebagai dan db.py jalankan contoh ini di sebelahnya.
Generator asinkron untuk himpunan data besar
Terapkan fungsi generator asinkron yang mengambil halaman hasil sesuai permintaan, memungkinkan penelepon untuk mengulangi himpunan data besar tanpa memuat semuanya ke dalam memori.
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 mengambil satu halaman pada satu waktu, sehingga memori tetap terbatas tidak peduli seberapa besar kumpulan hasilnya:
Streamed 31465 orders in chunks of 500.
Praktik terbaik
Terapkan panduan ini untuk menjaga pola asinkron tetap aman dan efisien.
Penentuan ukuran eksekutor yang tepat
Sesuaikan ukuran thread pool agar sesuai dengan beban kerja yang bergantung pada I/O, bukan hanya berdasarkan jumlah 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)
Penonaktifan yang anggun
Hentikan pelari tugas, kuras pekerjaan dalam penerbangan, lalu tutup kumpulan dan pelaksana secara berurutan.
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)
Penanganan kesalahan
Lakukan percobaan ulang hanya saat terjadi kegagalan sementara, dan beri jeda yang dibatasi nilai maksimumnya serta jitter. Gunakan kembali pengklasifikasi is_transient_error dari Logika percobaan ulang dan ketahanan koneksi agar kegagalan permanen seperti kredensial tidak valid atau kesalahan sintaks segera gagal alih-alih dicoba ulang.
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)