Подтверждение обработки сообщения из очереди RabbitMQ
Организую для взаимосвязи приложений очередь.
Использую библиотеку aio-pika для того, чтобы получать сообщения.
Проблема:
Сообщения успешно получаю и обрабатываю их, НО по средствам мониторинга вижу, что сообщения переходят из статуса Ready в статус Unranked и всё.
Т.е. мой код не отправлят подтверждение обработки. Поэтому при перезапуске слушателя, все сообщения отправлять по новой.
Вопрос:
Как сделать подтверждение обработки сообщение в библиотеке aio-pika?
Или здесь нужно действовать иначе?
Инструменты:
aio-pika (latest):
https://aio-pika.readthedocs.io/en/latest/
RabbitMQ: 3.8.9
Erlang: 23.1.3
Авторизован как админ.
Код отправителя:
import asyncio
import aio_pika
from aio_pika import DeliveryMode
async def main(loop):
connection = await aio_pika.connect_robust(
"amqp://admin:[email protected]/", loop=loop
)
routing_key = 'events_queue'
channel = await connection.channel()
data = {
'key': 'val',
}
import json
await channel.default_exchange.publish(
aio_pika.Message(
body=json.dumps(data).encode(),
app_id='999',
delivery_mode=DeliveryMode.PERSISTENT,
),
routing_key=routing_key,
)
await connection.close()
if __name__ == "__main__":
loop = asyncio.get_event_loop()
loop.run_until_complete(main(loop))
loop.close()
Код слушателя:
import asyncio
import json
from aio_pika import IncomingMessage, connect
from settings import (
RABBIT_URL, ROUTING_KEY
)
async def on_message(message: IncomingMessage):
print('Get: {}'.format(message.body))
try:
message_data = json.loads(message.body)
except json.decoder.JSONDecodeError:
return
print(message_data)
async def main(_loop):
connection = await connect(RABBIT_URL, loop=_loop)
channel = await connection.channel()
await channel.set_qos()
queue = await channel.declare_queue(
ROUTING_KEY, durable=True,
)
await queue.consume(on_message)
if __name__ == '__main__':
loop = asyncio.get_event_loop()
loop.create_task(main(loop))
loop.run_forever()
Запускаю RabbitMQ в докере:
version: '3.1'
services:
rabbit:
image: "rabbitmq:3-management"
environment:
- RABBITMQ_DEFAULT_USER=admin
- RABBITMQ_DEFAULT_PASS=pass
ports:
- 15673:15672
- 5672:5672