Python serve RabbitMq messages using Flask as SSE

Viewed 115

I'm working on a microservice endpoint which "only" consumes messages from RabbitMq and then serves those messages as SSE events.

This is my code for the endpoint:

data = queue.Queue()
def sseRouteHandler(reqId):


    def consume():
        connection = pika.BlockingConnection(pika.ConnectionParameters(host=rmqHostName,
                                            port=rmqPort,
                                            virtual_host='/',
                                            credentials=pika.PlainCredentials(username=rmqUserName, password=rmqPassword),
                                            connection_attempts=retryCount,
                                            retry_delay=retryInterval))

        channel = connection.channel()

        channel.queue_declare(queue=consumerQueues, auto_delete=False, exclusive=False, arguments=None)

        channel.queue_bind(queue=consumerQueues, exchange="eis.ds", routing_key=consumerQueues)

        def callback(ch, method, properties, body):
            # print(body)
            data.put(body)
            # ch.basic_ack(delivery_tag = method.delivery_tag)


        channel.basic_consume(queue=consumerQueues, on_message_callback=callback, exclusive=False, arguments=None)
        channel.start_consuming()

thread = Thread(target=consume)
thread.start()


    # Check if queue is empty, if not then pop the element else continue. Not working... I only get values after I close the server using keyboard interrupt
    def xcallback():
        while True:
            if not data.empty():
                yield data.get(block=False)
            else:
                continue


return Response(xcallback(), mimetype="text/event-stream")

Check if queue is empty, if not then get the element else continue. This is not working... I only get values after I close the server using keyboard interrupt

^C^CTraceback (most recent call last):


File "python3.6/site-packages/waitress/server.py", line 307, in run
    use_poll=self.adj.asyncore_use_poll,
  File "python3.6/site-packages/waitress/wasyncore.py", line 222, in loop
    poll_fun(timeout, map)
  File "python3.6/site-packages/waitress/wasyncore.py", line 152, in poll
    r, w, e = select.select(r, w, e, timeout)
KeyboardInterrupt


curl -X GET 'http://localhost:9080/v1/req/1ljhzckm3'
curl: (18) transfer closed with outstanding read data remaining
{"requestId": "1212122", "message": "Dat"}

What could be the solution for this?

0 Answers
Related