2015-10-22 01:09:30 +00:00
|
|
|
#!/usr/bin/env python
|
|
|
|
|
2016-03-09 01:59:43 +00:00
|
|
|
from kombu import Connection, Queue
|
2015-10-22 01:17:50 +00:00
|
|
|
from kombu.mixins import ConsumerProducerMixin
|
2015-10-22 01:09:30 +00:00
|
|
|
|
|
|
|
rpc_queue = Queue('rpc_queue')
|
|
|
|
|
|
|
|
|
|
|
|
def fib(n):
|
|
|
|
if n == 0:
|
|
|
|
return 0
|
|
|
|
elif n == 1:
|
|
|
|
return 1
|
|
|
|
else:
|
|
|
|
return fib(n-1) + fib(n-2)
|
|
|
|
|
|
|
|
|
2015-10-22 01:17:50 +00:00
|
|
|
class Worker(ConsumerProducerMixin):
|
2015-10-22 01:09:30 +00:00
|
|
|
|
|
|
|
def __init__(self, connection):
|
|
|
|
self.connection = connection
|
|
|
|
|
|
|
|
def get_consumers(self, Consumer, channel):
|
|
|
|
return [Consumer(
|
|
|
|
queues=[rpc_queue],
|
|
|
|
on_message=self.on_request,
|
|
|
|
accept={'application/json'},
|
|
|
|
prefetch_count=1,
|
|
|
|
)]
|
|
|
|
|
|
|
|
def on_request(self, message):
|
|
|
|
n = message.payload['n']
|
|
|
|
print(' [.] fib({0})'.format(n))
|
|
|
|
result = fib(n)
|
|
|
|
|
2015-10-22 01:17:50 +00:00
|
|
|
self.producer.publish(
|
|
|
|
{'result': result},
|
|
|
|
exchange='', routing_key=message.properties['reply_to'],
|
|
|
|
correlation_id=message.properties['correlation_id'],
|
|
|
|
serializer='json',
|
|
|
|
retry=True,
|
|
|
|
)
|
2015-10-22 01:09:30 +00:00
|
|
|
message.ack()
|
|
|
|
|
|
|
|
|
|
|
|
def start_worker(broker_url):
|
|
|
|
connection = Connection(broker_url)
|
|
|
|
print " [x] Awaiting RPC requests"
|
|
|
|
worker = Worker(connection)
|
|
|
|
worker.run()
|
|
|
|
|
|
|
|
|
|
|
|
if __name__ == '__main__':
|
2015-10-22 01:17:50 +00:00
|
|
|
try:
|
|
|
|
start_worker('pyamqp://')
|
|
|
|
except KeyboardInterrupt:
|
|
|
|
pass
|