【问题标题】:Loading different data frames in multiple threads at the same在多个线程中同时加载不同的数据帧
【发布时间】:2019-02-08 06:55:30
【问题描述】:

我有一个烧瓶服务器,它对数据帧执行读写查询。我有一个缓存机制(使用 cacheout 库)在我收到请求时缓存数据帧,然后在收到对同一数据帧的请求时使用缓存的数据帧。

目前我正在使用一个锁,它使所有线程按顺序加载它们的(不同的)数据帧,然后进一步处理加载的数据帧。

我想要的是,当我收到对不同数据帧的多个请求时,每个线程(对于每个请求)应该同时将数据帧(使用 pandas.read_excel)加载到内存中,而不是按顺序加载。

目前我正在使用一个简单的锁来确保相同的数据帧不会被加载两次,但我还需要并行加载多个数据帧。

`def read_query_request(query, file_path, sheet_name, source_id): logger.info('处理源的读取请求' + sheet_name + '_' + source_id)

try:
    data_frame_identifier = sheet_name + '_' + source_id

    # Load df with lock ensuring data frame loads only once.
    with lock:
        start_l=time.time()
        load_data_frame(file_path, sheet_name, source_id)
        end_l=time.time()
        logger.info('BENCHMARKING INFO: Read Request, Data frame load time ---' + str(end_l - start_l))

    #cache_state()
    # Executing query on loaded data frame
    # sheetName = getSheetName( query )
    query = query.replace('dataframe', data_frame_identifier)
    start_e = time.time()
    queryResult = ps.sqldf(query)
    end_e = time.time()
    logger.info('BENCHMARKING INFO: Read Request, psql query execution time ---' + str(end_e - start_e))

    start_j = time.time()
    queryResult = queryResult.to_json(orient='records')
    res = {"isErrored":"False", "results": json.loads(queryResult)}
    result = json.dumps(res)
    end_j = time.time()
    logger.info('BENCHMARKING INFO: Read Request, json conversion time ---' + str(end_j - start_j))

    logger.info(LRU_cache.keys())
    return result`

【问题讨论】:

    标签: python multithreading flask


    【解决方案1】:

    我从您的代码中了解到,您正在为整个应用程序使用一个锁,这限制了一次只能处理一个数据帧,并且您希望并行处理多个数据帧。首先,Python 中的线程(因为GIL)不能并行运行,而是按顺序运行。因此,如果您想要并行执行,则需要多处理。最简单的实现是使用stdlib 的multiprocessing pool。但是您仍然需要一些同步以避免一次处理多个 df。为此,您可以保留当前正在处理的 df 的注册表:

    ...
    registry_change_lock = Lock()
    registry = set()  # you can use list, but search in list is O(n) and in set is O(1)
    while True:
        registry_change_lock.acquire()
        if source_id in registry:
            # This id is already processed, release the lock 
            # to let other threads register their ids and avoid 
            # deadlocks
            registry_change_lock.release()
            time.sleep(0.5)
            continue
        else:
            registry.add(source_id)
            registry_change_lock.release()
    

    附:这不是解决问题的唯一方法,而是更简单的方法之一。

    【讨论】:

      猜你喜欢
      • 2021-09-08
      • 1970-01-01
      • 2020-06-21
      • 2018-12-28
      • 1970-01-01
      • 2019-09-08
      • 1970-01-01
      • 1970-01-01
      • 2022-08-09
      相关资源
      最近更新 更多