-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathutils.py
More file actions
61 lines (51 loc) · 1.71 KB
/
Copy pathutils.py
File metadata and controls
61 lines (51 loc) · 1.71 KB
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
from pika import BlockingConnection, ConnectionParameters, PlainCredentials
class MessageQueueUtils:
MQ_DIRECT_EXCHANGE_TYPE = "direct"
MQ_TOPIC_EXCHANGE_TYPE = "topic"
MQ_DIRECT_EXCHANGE_NAME = "1606875806_DIRECT"
MQ_TOPIC_EXCHANGE_NAME = "1606875806_TOPIC"
MQ_HOST = "152.118.148.95"
MQ_PORT = 5672
MQ_USERNAME = "0806444524"
MQ_PASSWORD = "0806444524"
MQ_VIRTUAL_HOST = "/0806444524"
def __init__(self, routing_key, exchange_mode):
self.routing_key = routing_key
self.exchange_mode = exchange_mode
self.exchange_name = self.MQ_DIRECT_EXCHANGE_NAME if self.exchange_mode == self.MQ_DIRECT_EXCHANGE_TYPE \
else self.MQ_TOPIC_EXCHANGE_NAME
self.pika_connection = BlockingConnection(
ConnectionParameters(
host=self.MQ_HOST,
virtual_host=self.MQ_VIRTUAL_HOST,
port=self.MQ_PORT,
credentials=PlainCredentials(self.MQ_USERNAME, self.MQ_PASSWORD)
)
)
self.connection_channel = self.pika_connection.channel()
self.connection_channel.exchange_declare(
exchange=self.exchange_name,
exchange_type=self.exchange_mode
)
self.connection_channel.queue_declare(queue=self.routing_key)
def send_message(self, message):
self.connection_channel.basic_publish(
exchange=self.exchange_name,
routing_key=self.routing_key,
body=message
)
def consume_message(self, callback_function):
self.connection_channel.queue_bind(
exchange=self.exchange_name,
queue=self.routing_key
)
def __callback(ch, method, properties, body):
callback_function(body)
self.connection_channel.basic_consume(
queue=self.routing_key,
on_message_callback=__callback,
auto_ack=True
)
self.connection_channel.start_consuming()
def close_connection(self):
self.pika_connection.close()