【问题标题】:How to increase write speed on inserts, pymongo?如何提高插入的写入速度,pymongo?
【发布时间】:2020-10-08 01:53:07
【问题描述】:

我有以下代码将文档插入 MongoDB,问题是它很慢,因为我无法对其进行多处理器处理,并且考虑到我必须检查插入的每个文档是否已经存在,我相信这是不可能的使用批量插入。我想知道是否有更快的方法来解决这个问题。在对下面进行分析后,我发现check record()update_upstream() 是两个非常耗时的函数。所以优化它们会提高整体速度。任何有关如何在下面进行优化的输入将不胜感激。谢谢!


import os
import pymongo


from directory import Directory
from pymongo import ASCENDING
from pymongo import DESCENDING
from pymongo import MongoClient
from storage_config import StorageConfig
from tqdm import tqdm

dir = Directory()

def DB_collections(collection_type):
    types = {'p': 'player_stats',
             't': 'team_standings',
             'f': 'fixture_stats',
             'l': 'league_standings',
             'pf': 'fixture_players_stats'}
    return types.get(collection_type)



class DB():

    def __init__(self, league, season, func=None):
        self.db_user = os.environ.get('DB_user')
        self.db_pass = os.environ.get('DB_pass')
        self.MONGODB_URL = f'mongodb+srv://{self.db_user}:{self.db_pass}@cluster0-mbqxj.mongodb.net/<dbname>?retryWrites=true&w=majority'
        self.league = league
        self.season = str(season)
        self.client = MongoClient(self.MONGODB_URL)
        self.DATABASE = self.client[self.league + self.season]


        self.pool = multiprocessing.cpu_count()
        self.playerfile = f'{self.league}_{self.season}_playerstats.json'
        self.teamfile = f'{self.league}_{self.season}_team_standings.json'
        self.fixturefile = f'{self.league}_{self.season}_fixturestats.json'
        self.leaguefile = f'{self.league}_{self.season}_league_standings.json'
        self.player_fixture = f'{self.league}_{self.season}_player_fixture.json'
        self.func = func

    def execute(self):
        if self.func is not None:
            return self.func(self)


def import_json(file):
    """Imports a json file in read mode
        Args:
            file(str): Name of file
    """
    return dir.load_json(file , StorageConfig.DB_DIR)

def load_file(file):
    try:
        loaded_file = import_json(file)
        return loaded_file
    except FileNotFoundError:
        print("Please check that", file, "exists")

def check_record(collection, index_dict):
    """Check if record exists in collection
        Args:
            index_dict (dict): key, value
    """
    return collection.find_one(index_dict)

def collection_index(collection, index, *args):
    """Checks if index exists for collection, 
    and return a new index if not

        Args:
            collection (str): Name of collection in database
            index (str): Dict key to be used as an index
            args (str): Additional dict keys to create compound indexs
    """
    compound_index = tuple((arg, ASCENDING) for arg in args)
    if index not in collection.index_information():
        return collection.create_index([(index, DESCENDING), *compound_index], unique=True)

def push_upstream(collection, record):
    """Update record in collection
        Args:
            collection (str): Name of collection in database
            record_id (str): record _id to be put for record in collection
            record (dict): Data to be pushed in collection
    """
    return collection.insert_one(record)

def update_upstream(collection, index_dict, record):
    """Update record in collection
        Args:
            collection (str): Name of collection in database
            index_dict (dict): key, value
            record (dict): Data to be updated in collection
    """
    return collection.update_one(index_dict, {"$set": record}, upsert=True)

def executePushPlayer(db):

    playerstats = load_file(db.playerfile)
    collection_name = DB_collections('p')
    collection = db.DATABASE[collection_name]
    collection_index(collection, 'p_id')
    for player in tqdm(playerstats):
        existingPost = check_record(collection, {'p_id': player['p_id']})
        if existingPost:
            update_upstream(collection, {'p_id': player['p_id']}, player)
        else:
            push_upstream(collection, player)

if __name__ == '__main__':
    db = DB('EN_PR', '2019')
        executePushPlayer(db)

【问题讨论】:

    标签: python python-3.x pymongo


    【解决方案1】:

    您可以使用 upsert=True 将检查/插入/更新逻辑合并到单个 update_one() 命令中,然后将批量运算符与以下内容一起使用:

    updates = []
    
    for player in tqdm(playerstats):
        updates.append(UpdateOne({'p_id': player['p_id']}, player, upsert=True))
    
    collection.bulk_write(updates)
    

    最后,请在 MongoDB shell 中使用以下命令检查您的索引是否正在使用:

    db.mycollection.aggregate([{ $indexStats: {} }])
    

    并查看 accesses.ops 指标。

    【讨论】:

    • 感谢您的意见。 UpdateOne 究竟是如何工作的?我不确定我是否完全理解如何实现这一点。
    • UpdateOneupdate_one() api.mongodb.com/python/current/api/pymongo/… 的批量操作符,upsert=True,如果记录匹配,它将执行更新,否则执行插入。
    • 谢谢你的澄清,我不知道!通过一些修改,我能够使您的代码工作,并将对其进行测试以查看它的速度。谢谢!
    猜你喜欢
    • 2021-11-15
    • 2016-08-14
    • 2011-10-20
    • 2012-07-18
    • 1970-01-01
    • 1970-01-01
    • 2012-12-01
    • 2013-09-30
    • 1970-01-01
    相关资源
    最近更新 更多