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

Python连接Docker MongoDB写入数据时RabbitMQ消费异常问题

问题根因

你现在连消息接收的打印日志都完全没有,说明receive_msg回调函数根本没有被触发,问题出在消费者启动阶段,和消息处理逻辑本身无关,常见诱因有两个:

  • 全局作用域的MongoDB连接逻辑没有做异常捕获,一旦出现MongoDB服务未启动、连接参数错误、依赖缺失的情况,程序会在启动阶段直接崩溃退出,根本走不到后续注册消费回调、监听队列的步骤,自然不会有任何日志输出。
  • 你贴出的代码缺失了RabbitMQ消费者的核心逻辑:包括RabbitMQ连接建立、队列声明、回调函数绑定、消费循环启动的代码,如果这部分漏写或者写错,消费者根本不会监听队列,自然收不到消息。

就算回调能正常触发,你现有代码的逻辑顺序也有严重问题:你在完成JSON解析、MongoDB写入之前就提前执行了消息ack操作,一旦后续步骤抛异常,消息会被RabbitMQ直接删除,造成永久数据丢失,且未捕获的异常会直接打崩消费者进程。

修复方案
  • 给启动阶段的所有连接逻辑加异常校验,提前暴露连接类错误,避免程序静默崩溃
  • 调整回调内的执行顺序:按照「接收消息→解析→写入数据库→所有操作成功后ack消息」的顺序执行,全流程加异常捕获,单条消息处理失败不要搞崩整个消费者
  • 补全RabbitMQ消费者的核心启动逻辑,确保队列名和生产者侧完全一致

修复后的可运行参考代码:

import pika
import pymongo
import json
import time

# MongoDB初始化,添加连接超时避免无限阻塞
myclient = pymongo.MongoClient(
    "mongodb://localhost:27017/",
    serverSelectionTimeoutMS=5000
)
# 启动时主动校验MongoDB连接状态,失败直接退出
try:
    myclient.admin.command("ping")
    print("[启动日志] MongoDB连接正常")
except Exception as e:
    print(f"[启动错误] MongoDB连接失败: {str(e)}")
    exit(1)

db = myclient.database_sample
my_collection = db["database"]

def receive_msg(ch, method, properties, body):
    try:
        msg_content = body.decode("utf-8")
        print(f"[消费日志] 收到消息: {msg_content}")
        time.sleep(2)
        # 先解析数据
        data = json.loads(body)
        # 写入MongoDB
        insert_result = my_collection.insert_one(data)
        print(f"[消费日志] 数据写入完成,文档ID: {insert_result.inserted_id}")
        # 所有操作成功后再确认消息
        ch.basic_ack(delivery_tag=method.delivery_tag)
        print("[消费日志] 消息已确认")
    except Exception as e:
        print(f"[消费错误] 消息处理失败: {str(e)}")
        # 处理失败拒绝消息,可根据业务需求决定是否重新入队
        ch.basic_nack(delivery_tag=method.delivery_tag, requeue=False)

if __name__ == "__main__":
    # 建立RabbitMQ连接
    rabbit_conn = pika.BlockingConnection(pika.ConnectionParameters("localhost"))
    channel = rabbit_conn.channel()
    # 声明队列,注意队列名必须和生产者发送消息的队列名完全一致
    channel.queue_declare(queue="msg_queue", durable=True)
    # 配置预取数,避免单消费者堆积过多消息
    channel.basic_qos(prefetch_count=1)
    # 绑定消费回调
    channel.basic_consume(queue="msg_queue", on_message_callback=receive_msg)
    print("[启动日志] 消费者已启动,等待接收消息...")
    channel.start_consuming()
额外排查项
  • 确认本地RabbitMQ、MongoDB服务都处于正常运行状态,对应端口没有被防火墙或其他进程占用
  • 核对生产者发送消息时指定的队列名,和消费者监听的队列名完全一致,不要出现拼写、大小写错误
  • 启动消费者时先看终端的启动日志,如果连启动阶段的打印都没有,直接看抛出的异常栈信息定位问题即可,不要跳过启动阶段的报错直接排查消费逻辑

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 10:18:22