MongoDB CursorNotFound错误求助(附Python批量更新代码)
修复MongoDB CursorNotFound错误的方案
问题原因
CursorNotFound(错误码43)是因为MongoDB服务器的游标超时被自动关闭。默认情况下,MongoDB游标在10分钟无操作后会被回收,而你的函数在遍历文档时,每个文档都需要调用generate_embeddings生成向量,这个过程如果耗时较长(比如集合文档量大、嵌入模型推理速度慢),就会导致游标在遍历完成前被服务器关闭。
修复方案
方案1:延长游标超时时间并分批执行批量更新
通过设置游标不超时,同时每处理一定数量的文档就执行一次批量更新,避免游标长时间闲置,也减少内存占用。
import pymongo import json def add_collection_content_vector_field(collection_name: str): ''' Add a new field to the collection to hold the vectorized content of each document. ''' collection = db[collection_name] bulk_operations = [] # 获取游标并禁用自动超时 cursor = collection.find() cursor.no_cursor_timeout = True try: for doc in cursor: if "contentVector" in doc: del doc["contentVector"] content = json.dumps(doc, default=str) content_vector = generate_embeddings(content) bulk_operations.append(pymongo.UpdateOne( {"_id": doc["_id"]}, {"$set": {"contentVector": content_vector}}, upsert=True )) # 每处理100个文档就执行批量更新,清空操作列表 if len(bulk_operations) >= 100: collection.bulk_write(bulk_operations) bulk_operations = [] # 处理剩余未执行的批量操作 if bulk_operations: collection.bulk_write(bulk_operations) finally: # 必须手动关闭游标,释放服务器资源 cursor.close()
注意:使用no_cursor_timeout后,一定要在finally块中手动关闭游标,否则服务器会残留无效游标,占用资源。
方案2:分批遍历文档(更安全的长期方案)
避免一次性遍历整个集合,改用分批获取文档的方式,每批处理完成后重新获取新的游标,彻底避免超时问题。
import pymongo import json def add_collection_content_vector_field(collection_name: str, batch_size: int = 100): ''' Add a new field to the collection to hold the vectorized content of each document. ''' collection = db[collection_name] total_docs = collection.count_documents({}) processed_count = 0 while processed_count < total_docs: bulk_operations = [] # 分批获取文档:跳过已处理的,取指定数量的新文档 for doc in collection.find().skip(processed_count).limit(batch_size): if "contentVector" in doc: del doc["contentVector"] content = json.dumps(doc, default=str) content_vector = generate_embeddings(content) bulk_operations.append(pymongo.UpdateOne( {"_id": doc["_id"]}, {"$set": {"contentVector": content_vector}}, upsert=True )) # 执行当前批次的更新 collection.bulk_write(bulk_operations) processed_count += batch_size print(f"已处理 {processed_count}/{total_docs} 个文档")
优化建议:如果集合是动态的(有新增/删除文档),用skip()+limit()可能会出现重复或遗漏,这时可以改用基于_id范围的分页:
last_id = None while True: query = {} if last_id: query = {"_id": {"$gt": last_id}} docs = list(collection.find(query).limit(batch_size)) if not docs: break # 处理文档... last_id = docs[-1]["_id"]
额外优化点
- 调整
batch_size:根据你的服务器内存和嵌入模型速度,选择合适的批量大小,平衡内存占用和处理效率。 - 异步处理:如果
generate_embeddings支持异步,可以用异步IO并行生成向量,减少单文档处理时间,降低游标超时风险。
内容的提问来源于stack exchange,提问作者Max Getuba
相关产品推荐
相关产品推荐

