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

如何让多Locust Worker节点从单一数据源获取不重复数据(云环境)

Locust多Worker节点无重复数据获取方案指导

方案1:数据库排他性查询(推荐高并发场景)

利用数据库行级锁特性,让每个Worker获取未被锁定的任务数据,天然避免重复。

以MySQL 8.0+为例,使用SKIP LOCKED语法实现无锁等待的排他性读取:

from locust import TaskSet, task
import pymysql

class UserBehavior(TaskSet):
    def on_start(self):
        self.db_conn = pymysql.connect(host="your-db-host", user="user", password="pass", db="test")
        self.cursor = self.db_conn.cursor(pymysql.cursors.DictCursor)
    
    @task
    def fetch_and_process_data(self):
        # 每次获取10条未处理且未被锁定的数据
        self.cursor.execute("""
            SELECT id, data FROM test_data 
            WHERE status = 'pending' 
            LIMIT 10 FOR UPDATE SKIP LOCKED
        """)
        data_list = self.cursor.fetchall()
        if not data_list:
            return
        # 模拟业务处理逻辑
        for item in data_list:
            print(f"Processing data {item['id']}")
        # 更新数据状态为已处理
        ids = [item['id'] for item in data_list]
        placeholders = ','.join(['%s']*len(ids))
        self.cursor.execute(f"UPDATE test_data SET status = 'processed' WHERE id IN ({placeholders})", ids)
        self.db_conn.commit()
    
    def on_stop(self):
        self.cursor.close()
        self.db_conn.close()

优点:无需额外中间件,数据库原生支持,适配高并发场景;缺点:依赖数据库版本特性,小批量取数会增加DB查询次数。

方案2:预分片分配(适合数据量固定的测试场景)

测试前根据Worker节点数量,将数据源按规则分片,每个Worker只处理专属分片的数据,彻底避免竞争。

比如按数据ID取模分片,假设共有3个Worker节点:

from locust import TaskSet, task, runners

class UserBehavior(TaskSet):
    def on_start(self):
        self.db_conn = pymysql.connect(host="your-db-host", user="user", password="pass", db="test")
        self.cursor = self.db_conn.cursor(pymysql.cursors.DictCursor)
        # 获取当前Worker的索引(从0开始)
        self.worker_index = runners.locust_runner.worker_index
        # 总Worker数量
        self.total_workers = runners.locust_runner.num_workers
    
    @task
    def fetch_and_process_data(self):
        # 只取ID模总Worker数等于当前Worker索引的数据
        self.cursor.execute("""
            SELECT id, data FROM test_data 
            WHERE status = 'pending' AND id % %s = %s
            LIMIT 10
        """, (self.total_workers, self.worker_index))
        data_list = self.cursor.fetchall()
        if not data_list:
            return
        # 处理数据+更新状态逻辑
        for item in data_list:
            print(f"Worker {self.worker_index} processing data {item['id']}")
        ids = [item['id'] for item in data_list]
        placeholders = ','.join(['%s']*len(ids))
        self.cursor.execute(f"UPDATE test_data SET status = 'processed' WHERE id IN ({placeholders})", ids)
        self.db_conn.commit()

优点:无DB锁竞争,查询效率高;缺点:测试前需保证数据已初始化,Worker数量固定后不能动态调整。

方案3:Master节点集中分发(适合数据动态生成场景)

由Master节点统一从数据库拉取数据,维护任务队列,Worker节点向Master请求数据,Master每次分发一批未处理的数据。

Master端核心逻辑:

# Master端维护全局任务队列
task_queue = []

def master_init():
    # 启动时从DB加载所有待处理数据
    db_conn = pymysql.connect(host="your-db-host", user="user", password="pass", db="test")
    cursor = db_conn.cursor(pymysql.cursors.DictCursor)
    cursor.execute("SELECT id, data FROM test_data WHERE status = 'pending'")
    global task_queue
    task_queue = cursor.fetchall()
    cursor.close()
    db_conn.close()

# 处理Worker的数据请求
def handle_worker_request(batch_size=10):
    global task_queue
    if not task_queue:
        return []
    batch = task_queue[:batch_size]
    task_queue = task_queue[batch_size:]
    return batch

Worker端核心逻辑:

from locust import TaskSet, task
import requests

class UserBehavior(TaskSet):
    @task
    def fetch_and_process_data(self):
        # 向Master节点请求数据批次
        response = requests.get("http://master-host:8080/get-task-batch", params={"batch_size":10})
        data_list = response.json()
        if not data_list:
            return
        # 模拟业务处理
        for item in data_list:
            print(f"Processing data {item['id']}")
        # 处理完成后更新数据库状态
        ids = [item['id'] for item in data_list]
        placeholders = ','.join(['%s']*len(ids))
        self.cursor.execute(f"UPDATE test_data SET status = 'processed' WHERE id IN ({placeholders})", ids)
        self.db_conn.commit()

优点:数据分发逻辑集中,Worker无需直接操作DB;缺点:Master可能成为性能瓶颈,需保证Master与Worker的网络稳定性。


内容的提问来源于stack exchange,提问作者B-Sam

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 14:05:22