-
Notifications
You must be signed in to change notification settings - Fork 0
/
Copy pathapp.py
72 lines (61 loc) · 2.44 KB
/
app.py
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
from threading import Thread
from random import randint
from time import sleep, time
import sys
import pika
class ReceiverThread(Thread):
def __init__(self, conn_params, queue_number, delay_range_ms):
Thread.__init__(self)
self.conn_params = conn_params
self.received = 0
self.last_received = 0
self.delay_range_ms = delay_range_ms
self.queue_number = queue_number
def run(self):
self.conn = pika.BlockingConnection(self.conn_params)
while True:
ch = self.conn.channel()
if self.try_register_consumer(ch):
print('Will consume from queue number %03d' % self.queue_number)
ch.start_consuming()
else:
print('Queue number %03d is locked, will try later' % self.queue_number)
sleep(10)
def try_register_consumer(self, ch):
queue_name='q.order-test.shard.%03d' % self.queue_number
try:
ch.basic_consume(
queue=queue_name,
consumer_tag='app_%s_cons_%03d' % (sys.argv[1], self.queue_number),
exclusive=True,
on_message_callback=self.on_message)
return True
except pika.exceptions.ChannelClosedByBroker as e:
if e.reply_code == 403:
return False
else:
raise
def on_message(self, ch, method, properties, body):
self.received += 1
ch.basic_publish(
exchange='', routing_key='q.order-test.results',
body='%s_app_%s_q_%03d_apptime_%f' % (body, sys.argv[1], self.queue_number, time()))
delay_ms = randint(self.delay_range_ms[0], self.delay_range_ms[1])
print('Received: %s, total count: %d, will process for %d ms' % (str(body), self.received, delay_ms))
sleep(float(delay_ms) / 1000)
ch.basic_ack(method.delivery_tag)
if __name__ == '__main__':
conn_params = pika.ConnectionParameters(
'localhost', 5672, '/',
pika.PlainCredentials('guest', 'guest'))
delay_min = randint(20, 300)
delay_max = delay_min + randint(10, 100)
print('Randomized delay range: %d to %d ms' % (delay_min, delay_max))
receivers = []
for i in range(1, 11):
receiver = ReceiverThread(conn_params, i, (delay_min, delay_max))
receiver.start()
receivers.append(receiver)
for receiver in receivers:
receiver.join()
print('Done!')