You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在双线程中运行Pika消费者与生产者并正确停止关闭?

解决Pika多线程下停止消费者并关闭通道/连接的AssertionError问题

问题分析

你遇到的AssertionError: ('_AsyncTransportBase._initate_abort() expected non-_STATE_COMPLETED', 4)错误,主要源于以下几点:

  • 消费者线程未启动消费循环,后续停止操作无有效目标
  • 多线程环境下未按正确顺序关闭资源,且未保证线程安全
  • 停止消费和关闭操作的时机不对,导致连接状态异常

修复方案及代码示例

关键修复点

  • 消费者线程必须启动消费循环,同时预留停止信号检测逻辑
  • 通过线程事件对象实现安全的线程间通信,避免跨线程直接操作通道
  • 严格遵循停止消费 → 关闭通道 → 关闭连接的资源释放顺序

修改后的完整代码

import time
import pika
import threading

url = 'amqp://user:user@localhost:5672/'  # 修正拼写错误:localhosh → localhost

def run_consumer(connection, tag, stop_event):
    def on_message(channel, method, properties, body):
        print(f"收到消息: {body.decode()}")
        channel.basic_ack(delivery_tag=method.delivery_tag)  # 手动确认消息,避免堆积

    channel = connection.channel()
    channel.queue_declare(queue='queue', auto_delete=True)
    channel.basic_consume(
        queue='queue',
        auto_ack=False,
        on_message_callback=on_message,
        consumer_tag=tag
    )

    # 循环处理事件,同时检测停止信号
    while not stop_event.is_set():
        try:
            # 限时处理事件,避免永久阻塞
            channel.process_data_events(time_limit=1)
        except (pika.exceptions.ConnectionClosedByBroker, pika.exceptions.AMQPChannelError):
            break
    # 取消消费者并关闭通道
    channel.basic_cancel(consumer_tag=tag)
    channel.close()
    print("消费者通道已关闭")

def run_publisher(connection, stop_event):
    channel = connection.channel()
    time.sleep(5)  # 缩短等待时间,便于测试
    for i in range(3):
        channel.basic_publish(
            exchange='',
            routing_key='queue',
            body=f'Test message {i+1}'.encode(),
        )
        print(f"发送消息: Test message {i+1}")
        time.sleep(2)
    # 发送完成后触发停止信号
    stop_event.set()
    channel.close()
    print("生产者通道已关闭")

if __name__ == '__main__':
    parameters = pika.URLParameters(url)
    connection = pika.BlockingConnection(parameters)
    consumer_tag = 'consumer'
    # 创建事件对象用于线程间停止信号传递
    stop_event = threading.Event()

    thread_consumer = threading.Thread(
        target=run_consumer,
        args=(connection, consumer_tag, stop_event),
        name='Consumer'
    )
    thread_consumer.start()

    thread_publisher = threading.Thread(
        target=run_publisher,
        args=(connection, stop_event),
        name='Publisher'
    )
    thread_publisher.start()

    # 等待所有线程执行完毕
    thread_consumer.join()
    thread_publisher.join()

    # 最后关闭连接
    if connection.is_open:
        connection.close()
    print("连接已关闭")

核心说明

  1. 线程安全控制:用threading.Event作为停止信号,避免跨线程直接调用通道方法引发状态冲突
  2. 消费循环优化:使用process_data_events(time_limit=1)替代start_consuming(),让消费者线程能定期检查停止信号,而非永久阻塞
  3. 资源释放顺序:生产者完成发送后触发停止信号,消费者收到信号后先取消消费者、关闭通道,最后主线程关闭连接
  4. 消息可靠性:添加basic_ack手动确认消息,防止队列中堆积未确认的消息

内容的提问来源于stack exchange,提问作者AlexanderKondrat

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.21 10:37:31