Docker环境下Pika+RabbitMQ跨连接消息确认失败问题咨询
Alright, let's tackle your two questions and figure out why your current approach isn't working, plus how to fix it.
Q1: Can I acknowledge a message received via Connection A using Connection B?
Nope, that's not possible with AMQP. Delivery tags are tied to the exact connection-channel pair that received the message. Each connection maintains its own independent sequence of delivery tags for each channel—RabbitMQ doesn't recognize a delivery tag from one connection as valid for another, even if you reuse the same channel number. The acknowledgment has to come from the original connection and channel that initially received the message.
Q2: Why is your current solution failing?
The error PRECONDITION_FAILED - unknown delivery tag 1 makes total sense now. When you create a new BlockingConnection (let's call this Connection B) and open a channel with number 1, you're starting a brand new session. The delivery tag 1 belongs exclusively to Connection A's channel 1 session; Connection B's channel 1 has its own delivery tag sequence starting at 1, but it never received the message you're trying to ack. RabbitMQ sees this as an invalid request and rejects it outright.
A Better Way to Handle Long-Running Tasks with Pika
Since you're using BlockingConnection and need to keep the heartbeat alive while processing slow tasks, here's a safe approach that doesn't require spinning up a new connection:
Pika isn't thread-safe, but it provides the add_callback_threadsafe method to safely run operations on the connection's main event loop from a worker thread. This way, you can offload your slow work to a daemon thread, then schedule the acknowledgment to run on the original connection's thread where it's valid.
Here's a code example to implement this:
import threading import pika def process_message(body, delivery_tag, connection): # Your time-consuming business logic goes here print(f"Processing message content: {body.decode()}") # Simulate slow work (replace with your actual task) # time.sleep(30) # Safely schedule the ack operation on the original connection's thread def ack_task(): if connection.is_open: channel = connection.channel() channel.basic_ack(delivery_tag=delivery_tag) connection.add_callback_threadsafe(ack_task) def on_message(ch, method, properties, body): delivery_tag = method.delivery_tag # Launch the slow task in a daemon thread threading.Thread( target=process_message, args=(body, delivery_tag, ch.connection), daemon=True ).start() # Setup your consumer connection and queue connection = pika.BlockingConnection(pika.ConnectionParameters(host="rabbitmq")) channel = connection.channel() channel.queue_declare(queue='myqueue', durable=False) channel.basic_consume(queue='myqueue', on_message_callback=on_message, auto_ack=False) try: channel.start_consuming() except KeyboardInterrupt: channel.stop_consuming() connection.close()
Alternatively, Use Thread-Safe AMQP Libraries
As you mentioned, libraries like rabbitpy and AMQPStorm are built to be thread-safe out of the box. With these tools, you can safely share a connection or channel across threads and call basic_ack directly from your worker thread without needing to use add_callback_threadsafe. This will simplify your code and reduce the chance of session-related errors.
内容的提问来源于stack exchange,提问作者plazmakeks

