【问题标题】:Python - multiprocessing multiple large size files using pandasPython - 使用熊猫多处理多个大文件
【发布时间】:2022-06-12 02:33:04
【问题描述】:

我有一个y.csv 文件。文件大小为 10 MB,包含来自Jan 2020 to May 2020 的数据。

我每个月也有一个单独的文件。例如data-2020-01.csv。它包含详细的数据。每个月文件的文件大小在1 GB左右。

我将y.csv 按月份拆分,然后通过加载相关月份文件来处理数据。当我花了很多个月的时间时,这个过程花费的时间太长了。例如24 个月。

我想更快地处理数据。我可以访问具有 32 vCPU128 GB 内存的 AWS m6i.8xlarge 实例。

我是多处理的新手。那么有人可以在这里指导我吗?

这是我当前的代码。

import pandas as pd

periods = [(2020, 1), (2020, 2), (2020, 3), (2020, 4), (2020, 5)]

y = pd.read_csv("y.csv", index_col=0, parse_dates=True).fillna(0)  # Filesize: ~10 MB


def process(_month_df, _index):
    idx = _month_df.index[_month_df.index.get_loc(_index, method='nearest')]
    for _, value in _month_df.loc[idx:].itertuples():

        up_delta = 200
        down_delta = 200

        up_value = value + up_delta
        down_value = value - down_delta

        if value > up_value:
            y.loc[_index, "result"] = 1
            return

        if value < down_value:
            y.loc[_index, "result"] = 0
            return


for x in periods:
    filename = "data-" + str(x[0]) + "-" + str(x[1]).zfill(2)  # data-2020-01
    filtered_y = y[(y.index.month == x[1]) & (y.index.year == x[0])]  # Only get the current month records
    month_df = pd.read_csv(f'{filename}.csv', index_col=0, parse_dates=True)  # Filesize: ~1 GB (data-2020-01.csv)

    for index, row in filtered_y.iterrows():
        process(month_df, index)

【问题讨论】:

  • 对同一主题感兴趣,遗憾的是还没有多进程经验,因此无法提供建议。只是一个观察,.iterrows(): 的最后一个块正在大大减慢你的进程。 stackoverflow.com/a/65356169/8805842也调查那部分
  • 这里的问题是您不能真正跨多个进程共享数据帧(由 y 引用)。它可以在多个线程之间共享,但这是一个有争议的问题,原因有两个:1)这是 CPU 绑定,所以多线程不合适 2)pandas 数据帧不是线程安全的
  • @NoobVB 因为我的filtered_y 很小,所以它不是这里的瓶颈。但是由于我这里只对索引感兴趣,所以我将它切换为itertuples。感谢您指出。
  • @LancelotduLac 我可以优化代码以不共享 y。我的 y 有唯一索引。
  • @John 请记住,10Mb 并不重要,对于 .iterrows() 或 itertuples(),行数是主要问题,所以只需检查您的 filters_y 的形状是否好奇.当然,请用你的 multiP 解决方案更新这个线程,-好奇 :)

标签: python python-3.x pandas multiprocessing


【解决方案1】:

多线程池非常适合在线程之间共享y 数据帧(无需使用共享内存),但不太擅长并行运行 CPU 密集型处理。多处理池非常适合进行 CPU 密集型处理,但在跨进程共享数据方面不太好,而没有提出 y 数据帧的碎片内存表示。

在这里,我重新排列了您的代码,以便我使用多线程池为每个周期创建 filtered_y 是一个 CPU 密集型操作,但 pandas 确实会释放全局解释器锁操作——希望是这个)。然后我们只将一个月的数据传递给多处理池,而不是整个y 数据帧,以使用工作函数process_month 处理该月。但由于每个池进程无权访问 y 数据帧,它只返回需要用要替换的值更新的索引。

import pandas as pd
from multiprocessing.pool import Pool, ThreadPool, cpu_count

def process_month(period, filtered_y):
    """
    returns a list of tuples consisting of (index, value) pairs
    """
    filename = "data-" + str(period[0]) + "-" + str(period[1]).zfill(2)  # data-2020-01
    month_df = pd.read_csv(f'{filename}.csv', index_col=0, parse_dates=True)  # Filesize: ~1 GB (data-2020-01.csv)
    results = []
    for index, row in filtered_y.iterrows():   
        idx = month_df.index[month_df.index.get_loc(index, method='nearest')]
        for _, value in month_df.loc[idx:].itertuples():
    
            up_delta = 200
            down_delta = 200
    
            up_value = value + up_delta
            down_value = value - down_delta
    
            if value > up_value:
                results.append((index, 1))
                break
    
            if value < down_value:
                results.append((index, 0))
                break
    return results

def process(period):
    filtered_y = y[(y.index.month == period[1]) & (y.index.year == period[0])]  # Only get the current month records
    for index, value in multiprocessing_pool.apply(process_month, (period, filtered_y)):
        y.loc[index, "result"] = value

def main():
    global y, multiprocessing_pool

    periods = [(2020, 1), (2020, 2), (2020, 3), (2020, 4), (2020, 5)]
    y = pd.read_csv("y.csv", index_col=0, parse_dates=True).fillna(0)  # Filesize: ~10 MB

    MAX_THREAD_POOL_SIZE = 100
    thread_pool_size = min(MAX_THREAD_POOL_SIZE, len(periods))
    multiprocessing_pool_size = min(thread_pool_size, cpu_count())
    with Pool(multiprocessing_pool_size) as multiprocessing_pool, \
    ThreadPool(thread_pool_size) as thread_pool:
        thread_pool.map(process, periods)
        
    # Presumably y gets written out again as a CSV file here?

# Required for Windows:
if __name__ == '__main__':
    main()

【讨论】:

  • main() 函数中,我看不到results 变量。如何访问该变量?
  • results 变量只返回给使用(index, value) 元组更新yprocess 工作函数,这是您最终想要做的。为什么main 需要这个元组列表?
  • 好的,我现在明白了。那么当这条线被执行y.loc[index, "result"] = value时,它在进程之外吗?我在某处读到无法访问进程内的全局变量。
  • 代码y.loc[index, "result"] = value 正在由运行在多线程池中的工作函数process 执行,该线程与y 被定义为全局的主进程在同一进程中运行。工作函数process_month 在多处理池(单独的进程)中运行,并使用传递的过滤月份生成这些元组,因为y 对其不可见,因此必须返回需要更新的列表。明白了吗?您是否真的运行过这个,因为我没有数据,因此我无法
  • 任何运气测试?很好奇这些.itertuples 和 multiP 是怎么回事
【解决方案2】:

正如在多个 pandas/threading 问题中所评论的,CSV 文件是 IO 绑定的,您可以从使用 ThreadPoolExecutor 中获得一些好处。

同时,如果您要执行聚合操作,请考虑在您的处理器内部执行read_csv,并改用ProcessPoolExecutor

如果您要在多进程之间传递大量数据,您还需要适当的内存共享方法。

但是我看到iterrowsitertuples 的用法 总的来说,这两个指令让我的眼睛流血了。您确定不能以矢量化模式处理数据吗?

这个特定部分我不确定它应该做什么,并且有 M 行会使其非常变慢。

def process(_month_df, _index):
    idx = _month_df.index[_month_df.index.get_loc(_index, method='nearest')]
    for _, value in _month_df.loc[idx:].itertuples():

        up_delta = 200
        down_delta = 200

        up_value = value + up_delta
        down_value = value - down_delta

        if value > up_value:
            y.loc[_index, "result"] = 1
            return

        if value < down_value:
            y.loc[_index, "result"] = 0
            return

在矢量化代码下方查找它是上升还是下降,以及在哪一行

df=pd.DataFrame({'vals': np.random.random(int(10))*1000+5000}).astype('int64')
print(df.vals.values)

up_value = 6000
down_value = 3000
valsup = df.vals.values + 200*np.arange(df.shape[0])+200
valsdown = df.vals.values - 200*np.arange(df.shape[0])-200

#! argmax returns 0 if all false
# idx_up = np.argmax(valsup > up_value)
# idx_dwn= np.argmax(valsdown < down_value)

idx_up = np.argwhere(valsup > up_value)
idx_dwn= np.argwhere(valsdown < down_value)
idx_up = idx_up[0][0] if len(idx_up) else -1
idx_dwn = idx_dwn[0][0] if len(idx_dwn) else -1


if idx_up < 0 and idx_dwn<0:
    print(f" Not up nor down")
if idx_up < idx_dwn or idx_dwn<0:
    print(f" Result is positive, in position {idx_up}")
else: 
    print(f" Result is negative, in position {idx_dwn}")

为了完整起见,对 1000 个元素进行基准测试 itertuples()argwhere 方法:

  • .itertuples(): 757µs
  • arange + argwhere: 60µs

【讨论】:

  • 我绝对更喜欢矢量化模式。但是,我相信在我的用例中这是不可能的,因为我正在检查 up_value 或 down_value 是否首先命中。所以顺序很重要。
  • 使用cumsum 并获得第一个索引怎么样?如果您提供一些样本数据我们也可以测试
  • 为此,我应该能够按照值的精确顺序 pd.cut 我的数据。我相信目前在熊猫中这是不可能的。如果您有任何想法,请告诉我。
  • 是的,很好,问题是关于 MP 的。 我的观点是,代码通常在没有优化的情况下被并行化
猜你喜欢
  • 2015-01-03
  • 2016-08-16
  • 1970-01-01
  • 2016-09-02
  • 1970-01-01
  • 2016-09-20
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多