JMeter如何仅在所有线程完成循环后执行一次最终校验
RabbitMQ负载测试后全量消息校验方案
针对你遇到的问题——多线程发送消息后无法在合适时机执行全量ID校验,核心解决思路是先统一收集所有发送的Message ID,等待所有发送任务完成后,再单次执行全量校验,具体实现分3步:
1. 线程安全存储所有发送的Message ID
- 用线程安全的容器(比如Java的
CopyOnWriteArrayList、Python带锁的列表)来存所有发送的消息ID,每个发送线程在成功发送一条消息后,立即把对应的ID存入容器。 - 必须保证容器的线程安全性,避免多线程写入时出现数据丢失或重复。比如Python里用
threading.Lock包裹写入操作,Java直接用线程安全集合。
2. 等待所有发送任务完全结束
- 如果用JMeter这类测试工具:
- 把发送请求放在一个独立的
Thread Group,校验逻辑放在Teardown Thread Group里——Teardown组会在所有主线程(发送线程)执行完毕后自动启动,天然满足“等待发送完成”的需求。 - 或者在主线程末尾添加
Flow Control Action,结合__jexl3函数判断所有线程是否完成发送,再触发校验。
- 把发送请求放在一个独立的
- 如果是自定义代码(Java/Python):
- Java用
CountDownLatch:初始化计数器为总发送条数(线程数×单线程循环数),每个发送任务完成后调用countDown(),校验逻辑调用await()等待计数器归0。 - Python用
threading.Semaphore或join():先启动所有发送线程,调用每个线程的join()方法,等所有发送线程结束后再执行校验代码。
- Java用
3. 单次执行全量校验逻辑
- 校验部分只执行一次,步骤如下:
- 从RabbitMQ队列拉取全部消息:可以用
basic.get循环拉取直到队列为空,或者设置足够大的预取数后批量消费。注意如果需要保留队列中的消息,消费时不要确认(no_ack=False),校验完成后可以选择重新入队;如果不需要保留,消费后直接确认。 - 提取所有拉取到的消息ID,存入一个集合(方便快速查找)。
- 遍历之前收集的发送ID集合,逐个断言每个ID都存在于拉取到的集合中,统计缺失的ID并输出结果。
- 从RabbitMQ队列拉取全部消息:可以用
示例代码(Python)
import threading import pika # 线程安全存储发送的消息ID sent_msg_ids = [] lock = threading.Lock() # 总发送任务数:10线程×100次循环=1000条 total_tasks = 1000 # 用信号量等待所有任务完成 task_semaphore = threading.Semaphore(total_tasks) def send_msg(): conn = pika.BlockingConnection(pika.ConnectionParameters("localhost")) channel = conn.channel() channel.queue_declare(queue="load_test_queue", durable=True) for idx in range(100): msg_id = f"test_msg_{threading.get_ident()}_{idx}" # 发送消息并携带message_id channel.basic_publish( exchange="", routing_key="load_test_queue", body=f"content_{msg_id}", properties=pika.BasicProperties(message_id=msg_id, delivery_mode=2) ) # 线程安全写入ID with lock: sent_msg_ids.append(msg_id) # 完成一个任务,信号量减1 task_semaphore.acquire() conn.close() def validate_all_msgs(): # 等待所有发送任务完成 for _ in range(total_tasks): task_semaphore.acquire() conn = pika.BlockingConnection(pika.ConnectionParameters("localhost")) channel = conn.channel() channel.queue_declare(queue="load_test_queue", durable=True) received_ids = set() # 拉取队列中所有消息 while True: method_frame, header_frame, body = channel.basic_get(queue="load_test_queue") if not method_frame: break received_ids.add(header_frame.message_id) # 若需保留消息,注释下面的确认语句;若无需保留则打开 # channel.basic_ack(method_frame.delivery_tag) # 全量校验 missing_ids = [msg_id for msg_id in sent_msg_ids if msg_id not in received_ids] if missing_ids: print(f"校验失败,缺失{len(missing_ids)}条消息,ID列表:{missing_ids}") else: print("所有发送的消息均已存在于队列,校验通过") conn.close() # 启动10个发送线程 for _ in range(10): threading.Thread(target=send_msg).start() # 启动单次校验线程 validate_thread = threading.Thread(target=validate_all_msgs) validate_thread.start() validate_thread.join()
内容的提问来源于stack exchange,提问作者Amphata
相关产品推荐
相关产品推荐

