MCPcopy Create free account
hub / github.com/PaddlePaddle/FastDeploy / consumer_task

Function consumer_task

benchmarks/benchmark_fmq.py:64–96  ·  view source on GitHub ↗
(consumer_id, total_msgs, result_q, consumer_event)

Source from the content-addressed store, hash-verified

62# Consumer Task
63# ============================================================
64async 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
99def consumer_process(consumer_id, total_msgs, result_q, consumer_event):

Callers 1

runFunction · 0.85

Calls 8

queueMethod · 0.95
FMQClass · 0.90
writeMethod · 0.80
setMethod · 0.45
getMethod · 0.45
updateMethod · 0.45
closeMethod · 0.45
putMethod · 0.45

Tested by

no test coverage detected