You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.27 04:17:13