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}条记录,程序执行完成!")
关键修改点说明
- 非阻塞获取元素:使用
queue.get(block=False)替代默认的阻塞获取,当队列空时直接抛出Queue.Empty异常,我们捕获这个异常并跳出收集循环,避免线程无限等待。 - 任务完成标记:每次获取元素后调用
queue.task_done(),这是队列的关键机制——它会告诉队列当前任务已经处理完成,主线程的queue.join()才能准确等待所有任务结束。 - 动态批量大小:不再强制凑够50条才写入,当队列剩余38条时,会直接收集这38条并执行插入,完美解决最后一批数据的写入问题。
- 线程正常退出:当收集不到任何数据时(队列空),线程会跳出循环并结束运行,不会一直占用资源。
内容的提问来源于stack exchange,提问作者Lucas Amos
相关产品推荐
相关产品推荐

