Locust搭配pyodbc时任务阻塞问题求助(ThreadPool无此异常)
Locust结合pyodbc任务阻塞问题的解决方案
问题原因
Locust基于gevent实现并发,gevent通过猴子补丁替换Python标准库中的同步IO操作来实现协程切换,但pyodbc是基于C扩展的同步数据库驱动,其底层的网络IO和数据库操作并未被gevent的猴子补丁处理,导致cursor.execute()这类操作会阻塞整个gevent事件循环,使得所有任务串行执行。而ThreadPool版本使用真实线程,不受此限制,因此能正常并行。
解决方案
方案1:用gevent线程池包装pyodbc操作
将阻塞的数据库操作提交到线程池中执行,让gevent在等待线程结果时切换到其他协程,实现并发。
修改后的Locust任务代码:
import logging import sys import time from concurrent.futures import ThreadPoolExecutor import pyodbc from locust import task, User, TaskSet, events # SQL Server Connection Details server = 'localhost' database = 'test' username = 'test' password = 'test' driver = 'ODBC Driver 17 for SQL Server' # 创建线程池,可根据并发需求调整大小 thread_pool = ThreadPoolExecutor(max_workers=20) count = 0 table_name = f"SampleTable_{int(time.time())}" class SampleTaskSet(TaskSet): def __init__(self, parent: User): super().__init__(parent) global count count += 1 self.counter = count def _execute_query(self): """实际数据库操作,放到线程中执行""" logging.info("User {} started".format(self.counter)) global table_name query = "select * from {}".format(table_name) conn = pyodbc.connect(driver=driver, host=server, Database=database, UID=username, PWD=password, Trusted_Connection='no') logging.info("User {} connection established".format(self.counter)) cursor = conn.cursor() try: logging.info("User %s Executing Query %s ", self.counter, query) cursor.execute(query) rows = cursor.fetchall() resp_len = len(rows) logging.info("User %s completed Query", self.counter) finally: cursor.close() conn.close() @task def execute_query(self): # 提交任务到线程池,非阻塞 future = thread_pool.submit(self._execute_query) # 等待结果,gevent会在等待时切换协程 future.result() class TestUser(User): tasks = [SampleTaskSet] @events.test_start.add_listener def prep_data(environment, **kwargs): global table_name conn = pyodbc.connect(driver=driver, host=server, Database=database, UID=username, PWD=password, Trusted_Connection='no', autocommit=True) cursor = conn.cursor() logging.info("Creating table %s", table_name) try: create_table_query = f"CREATE TABLE {table_name} (id INT, name VARCHAR(255), email VARCHAR(255))" cursor.execute(create_table_query) logging.info("Created table %s", table_name) logging.info("Loading data to table %s", table_name) for i in range(100_000): insert_query = f"INSERT INTO {table_name} (id, name, email) VALUES (?, ?, ?)" cursor.execute(insert_query, (i + 1, f"Name {i + 1}", f"email_{i + 1}@example.com")) logging.info("Loaded data to table %s", table_name) finally: cursor.close() conn.close()
方案2:改用异步数据库驱动(aiodbc)
使用支持asyncio的aiodbc库,适配Locust的异步模式,从根本上避免阻塞问题。
首先安装依赖:
pip install aiodbc
异步版本Locust代码:
import logging import time import asyncio import aiodbc from locust import task, User, events # SQL Server Connection Details server = 'localhost' database = 'test' username = 'test' password = 'test' driver = 'ODBC Driver 17 for SQL Server' connection_string = f"DRIVER={driver};SERVER={server};DATABASE={database};UID={username};PWD={password}" count = 0 table_name = f"SampleTable_{int(time.time())}" class AsyncTestUser(User): abstract = True async def connect_db(self): conn = await aiodbc.connect(connection_string=connection_string) return conn class SampleTaskSet(AsyncTestUser): def __init__(self, parent): super().__init__(parent) global count count += 1 self.counter = count @task async def execute_query(self): logging.info("User {} started".format(self.counter)) global table_name query = "select * from {}".format(table_name) conn = await self.connect_db() logging.info("User {} connection established".format(self.counter)) cursor = await conn.cursor() try: logging.info("User %s Executing Query %s ", self.counter, query) await cursor.execute(query) rows = await cursor.fetchall() resp_len = len(rows) logging.info("User %s completed Query", self.counter) finally: await cursor.close() await conn.close() @events.test_start.add_listener async def prep_data(environment, **kwargs): global table_name conn = await aiodbc.connect(connection_string=connection_string, autocommit=True) cursor = await conn.cursor() logging.info("Creating table %s", table_name) try: create_table_query = f"CREATE TABLE {table_name} (id INT, name VARCHAR(255), email VARCHAR(255))" await cursor.execute(create_table_query) logging.info("Created table %s", table_name) logging.info("Loading data to table %s", table_name) for i in range(100_000): insert_query = f"INSERT INTO {table_name} (id, name, email) VALUES (?, ?, ?)" await cursor.execute(insert_query, (i + 1, f"Name {i + 1}", f"email_{i + 1}@example.com")) logging.info("Loaded data to table %s", table_name) finally: await cursor.close() await conn.close()
内容的提问来源于stack exchange,提问作者toothferry
相关产品推荐
相关产品推荐

