| 62 | # Consumer Task |
| 63 | # ============================================================ |
| 64 | async def consumer_task(consumer_id, total_msgs, result_q, consumer_event): |
| 65 | fmq = FMQ() |
| 66 | q = fmq.queue("mp_bench_latency", role="consumer") |
| 67 | consumer_event.set() |
| 68 | |
| 69 | latencies = [] |
| 70 | recv = 0 |
| 71 | |
| 72 | # tqdm 显示进度 |
| 73 | pbar = tqdm(total=total_msgs, desc=f"Consumer-{consumer_id}", position=consumer_id + 1, leave=True, disable=False) |
| 74 | |
| 75 | first_recv = None |
| 76 | last_recv = None |
| 77 | |
| 78 | while recv < total_msgs: |
| 79 | msg = await q.get() |
| 80 | recv_ts = time.perf_counter() |
| 81 | if msg is None: |
| 82 | pbar.write("recv None") |
| 83 | continue |
| 84 | if first_recv is None: |
| 85 | first_recv = recv_ts |
| 86 | last_recv = recv_ts |
| 87 | send_ts = msg.payload["send_ts"] |
| 88 | latencies.append((recv_ts - send_ts) * 1000) # ms |
| 89 | pbar.update(1) |
| 90 | recv += 1 |
| 91 | |
| 92 | pbar.close() |
| 93 | |
| 94 | result_q.put( |
| 95 | {"consumer_id": consumer_id, "latencies": latencies, "first_recv": first_recv, "last_recv": last_recv} |
| 96 | ) |
| 97 | |
| 98 | |
| 99 | def consumer_process(consumer_id, total_msgs, result_q, consumer_event): |