如何让多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
相关产品推荐
相关产品推荐

