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

Python多线程MongoDB批量写入程序挂起问题求助

解决MongoDB多线程批量写入挂起问题

我来帮你搞定这个问题!你的程序挂起的核心原因非常清晰:当队列里剩余记录不足50条时,线程的run方法还在硬循环执行50次self.queue.get()——而这个方法默认是阻塞式的,如果队列空了,它会一直等待新元素进来,但此时你已经没有更多数据要入队了,所以线程就卡在这一步,永远完不成最后一批38条记录的写入。

核心解决方案思路

我们需要修改线程的批量获取逻辑:

  • 不再固定强制收集50条记录,而是收集最多50条,当队列剩余记录不足时,收集剩下的所有记录
  • 处理队列空的情况,让线程能正常退出,避免无限阻塞
  • 必须调用queue.task_done(),确保队列能正确跟踪任务完成状态,让主线程可以等待所有任务结束

修改后的完整代码

import threading
import Queue
import json
from pymongo import MongoClient, InsertOne
import datetime

# 设置线程数量
n_thread = 10
# 创建任务队列
queue = Queue.Queue()
# 连接MongoDB数据库
client = MongoClient("mongodb://mydatabase.com")
db = client.threads

class ThreadClass(threading.Thread):
    def __init__(self, queue):
        threading.Thread.__init__(self)
        self.queue = queue

    def run(self):
        while True:
            objects = []
            # 最多收集50条记录,直到队列为空
            for _ in range(50):
                try:
                    # 非阻塞获取队列元素,空队列时抛出Empty异常
                    item = self.queue.get(block=False)
                    objects.append(item)
                    # 标记该任务已处理完成
                    self.queue.task_done()
                except Queue.Empty:
                    break  # 队列空了,停止收集
            
            # 只有当有数据时才执行批量插入
            if objects:
                db.threads.insert_many(objects)
            else:
                # 队列已空,没有更多数据,线程退出
                break

# ----------------------
# 这里是填充队列的示例代码(替换成你实际的数据源)
total_records = 530838
for i in range(total_records):
    queue.put({
        "record_id": i,
        "content": f"测试内容_{i}",
        "create_time": datetime.datetime.now()
    })
# ----------------------

# 启动所有线程
threads = []
for _ in range(n_thread):
    t = ThreadClass(queue)
    t.start()
    threads.append(t)

# 等待队列中所有任务都被处理完毕
queue.join()

# 等待所有线程正常退出
for t in threads:
    t.join()

print(f"成功写入{total_records}条记录,程序执行完成!")

关键修改点说明

  1. 非阻塞获取元素:使用queue.get(block=False)替代默认的阻塞获取,当队列空时直接抛出Queue.Empty异常,我们捕获这个异常并跳出收集循环,避免线程无限等待。
  2. 任务完成标记:每次获取元素后调用queue.task_done(),这是队列的关键机制——它会告诉队列当前任务已经处理完成,主线程的queue.join()才能准确等待所有任务结束。
  3. 动态批量大小:不再强制凑够50条才写入,当队列剩余38条时,会直接收集这38条并执行插入,完美解决最后一批数据的写入问题。
  4. 线程正常退出:当收集不到任何数据时(队列空),线程会跳出循环并结束运行,不会一直占用资源。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 03:53:31