Continuously consume from queue with aio-pika

Viewed 316

I am building an asynchronous consumer class with aio-pika to continuously retrieve messages from a RabbitMQ queue. Based on the documentation, my code looks like the following :

class Consumer:
    ... __init__() and stuff ...

    async def run(self):
        ... Connection stuff ...
        await self._queue.consume(self.handle_message)

        # This part I don't understand
        await self._loop.create_future()

    def close(self):
        self._channel.close()
        self._loop.stop()


# Main part
loop = asyncio.get_event_loop()
consumer = Consumer(loop)

try:
    loop.create_task(consumer.run())
    loop.run_forever()
except KeyboardInterrupt:
    consumer.close()

My first issue is that I don't understand why adding the line await self._loop.create_future() is necessary. It looks strange to have an empty future running. If I remove it, it seems to be working fine.

My second issue is that this line causes a warning when I stop the program :

Task was destroyed but it is pending!
task: <Task pending name='Task-1' coro=<Consumer.run() running at scripts/consumer.py:59> wait_for=<Future pending cb=[<TaskWakeupMethWrapper object at 0x7fdbfa0134f0>()]>>

Oddly, when I store the future in an instance attribute, the warning disappears :

self._exit = self._loop.create_future()
await self._exit

To make it cleaner, I also cancel the future upon exiting, but it does not seem to be changing anything :

def close(self):
    self._exit.cancel()
    self._channel.close()
    self._loop.stop()

My question is then : what exactly is happening in this line and is it right to remove it completely ?

0 Answers
Related