I have a code based on this example https://github.com/pika/pika/blob/0.12.0/examples/basic_consumer_threaded.py which receive data from queues and process them on another thread (as the process is really long). But even if I am using this example from pika I'm still getting an error :
1654011396.901431 - pika.channel - WARNING - channel._on_close_from_broker - Received remote Channel.Close (406): 'PRECONDITION_FAILED - delivery acknowledgement on channel 1 timed out. Timeout value used: 1800000 ms. This timeout value can be configured, see consumers doc guide to learn more' on <Channel number=1 OPEN conn=<SelectConnection OPEN transport=<pika.adapters.utils.io_services_utils._AsyncPlaintextTransport object at 0x7faf9b547670> params= ConnectionParameters host=* port=* virtual_host=/ ssl=False>>>
Traceback (most recent call last): File "rabbit.py", line 91, in channel.start_consuming() File "/opt/app/project/venv/lib/python3.8/site-packages/pika/adapters/blocking_connection.py", line 1865, in start_consuming self._process_data_events(time_limit=None) File "/opt/app/project/venv/lib/python3.8/site-packages/pika/adapters/blocking_connection.py", line 2031, in _process_data_events raise self._closing_reason # pylint: disable=E0702 pika.exceptions.ChannelClosedByBroker: (406, 'PRECONDITION_FAILED - delivery acknowledgement on channel 1 timed out. Timeout value used: 1800000 ms. This timeout value can be configured, see consumers doc guide to learn more')
My code looks like :
import pika
from config import INGESTION_EXCHANGE
from config import DATA_QUEUE, FIRST_DATA_BINDING_KEY, SECOND_DATA_BINDING_KEY, POISON_BINDING_KEY
from config import RABBIT_CREDENTIALS, RABBIT_HOST, RABBIT_PORT
from config import logger, CONTAINER_ID
from config import path_first_data, path_second_data
import json
import logging
from process_data import process_data
import argparse
import threading
import functools
logging.getLogger("pika").setLevel(logging.WARNING)
logger.info("{} - Waiting to receive data".format(CONTAINER_ID))
f_buffer, s_buffer = [], []
def ack_message(channel, delivery_tag):
if channel.is_open:
channel.basic_ack(delivery_tag)
else:
logger.info("{} - Channel closed".format(CONTAINER_ID))
pass
def do_work(connection, channel, delivery_tag):
process_data(client="*", type_data="*")
cb = functools.partial(ack_message, channel, delivery_tag)
connection.add_callback_threadsafe(cb)
def on_message(channel, method_frame, properties, body, args):
(connection, threads) = args
delivery_tag = method_frame.delivery_tag
global f_buffer
global s_buffer
if body.decode("utf-8") == " ":
logger.info("{} - First Message received - {}".format(CONTAINER_ID, str(len(f_buffer))))
logger.info("{} - Second Message received - {}".format(CONTAINER_ID, str(len(s_buffer))))
if (len(f_buffer) == 0) or (len(s_buffer) == 0):
logger.error("{} - Data received is empty".format(CONTAINER_ID))
else:
with open(path_first_data, "w") as f:
json.dump(f_buffer, f)
with open(path_second_data, "w") as f:
json.dump(s_buffer, f)
f_buffer, s_buffer = [], []
t = threading.Thread(target=do_work, args=(connection, channel, delivery_tag))
t.start()
threads.append(t)
else:
data_received = json.loads(body.decode("utf-8"))
if "first_id" in data_received.keys():
if data_received["first_id"] is None:
pass
else:
f_buffer.append(data_received)
if "second_id" in data_received.keys():
if data_received["second_id"] is None:
pass
else:
s_buffer.append(data_received)
channel.basic_ack(delivery_tag=delivery_tag)
credentials = pika.PlainCredentials(RABBIT_CREDENTIALS["username"], RABBIT_CREDENTIALS["password"])
parameters = pika.ConnectionParameters(RABBIT_HOST, RABBIT_PORT, '/', credentials)
connection = pika.BlockingConnection(parameters=parameters)
channel = connection.channel()
channel.exchange_declare(exchange=INGESTION_EXCHANGE, exchange_type='direct', durable=True)
channel.queue_declare(queue=DATA_QUEUE, durable=True)
channel.queue_bind(exchange=INGESTION_EXCHANGE, queue=DATA_QUEUE, routing_key=FIRST_DATA_BINDING_KEY)
channel.queue_bind(exchange=INGESTION_EXCHANGE, queue=DATA_QUEUE, routing_key=SECOND_DATA_BINDING_KEY)
channel.queue_bind(exchange=INGESTION_EXCHANGE, queue=DATA_QUEUE, routing_key=POISON_BINDING_KEY)
channel.basic_qos(prefetch_count=1)
threads = []
on_message_callback = functools.partial(on_message, args=(connection, threads))
channel.basic_consume(queue=DATA_QUEUE, on_message_callback=on_message_callback, auto_ack=False)
try:
channel.start_consuming()
except KeyboardInterrupt:
channel.stop_consuming()
# Wait for all to complete
for thread in threads:
thread.join()
Do someone know what can I do to avoid this error without changing the configuration of rabbitmq ?
Thank you.