-
Notifications
You must be signed in to change notification settings - Fork 0
/
Copy pathworker.py
30 lines (21 loc) · 1 KB
/
worker.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
import json
from decouple import config
from pika import BlockingConnection
from pika import ConnectionParameters
BROKER_URL = config("BROKER_URL")
BROKER_EXCHANGE_NAME = config("BROKER_EXCHANGE_NAME")
GAZU_WORKER_QUEUE_NAME = config("GAZU_WORKER_QUEUE_NAME")
def digest_event(channel, delivery, properties, payload):
dispatch_event(delivery.routing_key, json.loads(payload))
channel.basic_ack(delivery_tag=delivery.delivery_tag)
def dispatch_event(routing_key, payload):
print(routing_key, payload)
params = ConnectionParameters(BROKER_URL)
connection = BlockingConnection(params)
channel = connection.channel()
channel.exchange_declare(exchange=BROKER_EXCHANGE_NAME, exchange_type='topic')
channel.queue_declare(queue=GAZU_WORKER_QUEUE_NAME, durable=True)
channel.queue_bind(exchange=BROKER_EXCHANGE_NAME, queue=GAZU_WORKER_QUEUE_NAME, routing_key="#")
channel.basic_qos(prefetch_count=20)
channel.basic_consume(queue=GAZU_WORKER_QUEUE_NAME, on_message_callback=digest_event)
channel.start_consuming()