【问题标题】:Parallel loading of Input Files in Pandas Dataframe在 Pandas Dataframe 中并行加载输入文件
【发布时间】:2019-06-16 00:02:43
【问题描述】:

我有一个需求,我有三个输入文件,需要将它们加载到 Pandas 数据框中,然后将其中两个文件合并到一个数据框中。

文件扩展名总是变化的,一次可能是 .txt,另一次可能是 .xlsx 或 .csv。

我怎样才能并行运行这个过程,以节省等待/加载时间?

这是我目前的代码,

from time import time # to measure the time taken to run the code
start_time = time()

Primary_File = "//ServerA/Testing Folder File Open/Report.xlsx"
Secondary_File_1 = "//ServerA/Testing Folder File Open/Report2.csv"
Secondary_File_2 = "//ServerA/Testing Folder File Open/Report2.csv"

import pandas as pd # to work with the data frames
Primary_df = pd.read_excel (Primary_File)
Secondary_1_df = pd.read_csv (Secondary_File_1)
Secondary_2_df = pd.read_csv (Secondary_File_2)

Secondary_df = Secondary_1_df.merge(Secondary_2_df, how='inner', on=['ID'])
end_time = time()

print(end_time - start_time)

加载 primary_df 和 secondary_df 大约需要 20 分钟。因此,我正在寻找一种可能使用并行处理来节省时间的有效解决方案。 我是通过Reading操作计时的,大部分时间大约是18分45秒。

硬件配置:- Intel i5 处理器、16 GB 内存和 64 位操作系统

有资格获得赏金的问题:- 因为我正在寻找工作 包含详细步骤的代码 - 在 anaconda 中使用 包 支持加载我的输入文件的环境 并行和 将它们分别存储在熊猫数据框中。这最终应该 节省时间。

【问题讨论】:

  • 你至少有3个选项; asyncio,线程,多进程,但我不确定这些选项是否会给你足够的性能。您需要查看读取操作是否占用大部分时间(在这种情况下,上述选项应该对您有所帮助)或者是否在内存中创建数据帧占用大部分时间。
  • 考虑使用 dask (docs.dask.org/en/latest/why.html) 作为 pandas 的替代品。
  • 你有什么样的硬件?如果瓶颈是磁盘 I/O,我不确定如何解决此问题。
  • @Logan 是的,英特尔 i5 处理器,16 GB 内存和 64 位操作系统。
  • 您受 IO 限制,无法绕过它。加载时间就是加载时间。

标签: python pandas anaconda


【解决方案1】:

试试这个:

from time import time 
import pandas as pd
from multiprocessing.pool import ThreadPool


start_time = time()

pool = ThreadPool(processes=3)

Primary_File = "//ServerA/Testing Folder File Open/Report.xlsx"
Secondary_File_1 = "//ServerA/Testing Folder File Open/Report2.csv"
Secondary_File_2 = "//ServerA/Testing Folder File Open/Report2.csv"


# Define a function for the thread
def import_xlsx(file_name):
    df_xlsx = pd.read_excel(file_name)
    # print(df_xlsx.head())
    return df_xlsx


def import_csv(file_name):
    df_csv = pd.read_csv(file_name)
    # print(df_csv.head())
    return df_csv

# Create two threads as follows

Primary_df = pool.apply_async(import_xlsx, (Primary_File, )).get() 
Secondary_1_df = pool.apply_async(import_csv, (Secondary_File_1, )).get() 
Secondary_2_df = pool.apply_async(import_csv, (Secondary_File_2, )).get() 

Secondary_df = Secondary_1_df.merge(Secondary_2_df, how='inner', on=['ID'])
end_time = time()

【讨论】:

  • 我将 _thread 作为线程导入。它在这段代码上给了我一个错误:- thread.start_new_thread(import_xlsx, [Primary_File]) 它说,第二个参数必须是一个元组..
  • 好的,现在试试。第二个参数现在是一个元组。
  • 现在,它给了我这个错误 - import_xlsx() 缺少 1 个必需的位置参数:'file_name'
  • 这给我带来了问题。数据帧 Primary_df、Secondary_1_df、Secondary_2_df 都返回单个整数值,而不是包含所有列和相应列值的数据帧。
  • @Sid29 是的,现在它返回 DataFrame,它在 3 个单独的线程上工作。他们每个人都在导入数据框。
【解决方案2】:

尝试使用@Cezary.Sz 代码但使用(删除对.get() 的调用),而不是:

Primary_df_job = pool.apply_async(import_xlsx, (Primary_File, ))
Secondary_1_df_job = pool.apply_async(import_csv, (Secondary_File_1, ))
Secondary_2_df_job = pool.apply_async(import_csv, (Secondary_File_2, ))

然后

Secondary_1_df = Secondary_1_df_job.get()
Secondary_2_df = Secondary_2_df_job.get()

您可以在加载Primary_df_job 时使用数据框。

Secondary_df = Secondary_1_df.merge(Secondary_2_df, how='inner', on=['ID'])

当您的代码中需要 Primary_df 时,请使用

Primary_df = Primary_df_job.get()

这将阻止执行,直到 Primary_df_job 完成。

【讨论】:

    【解决方案3】:

    为什么不用asyncio 而不是multiprocessing

    您可能希望首先使用an Async CSV Dict Reader(可以使用multiprocessing 并行化多个文件)来利用I/O 级别,而不是使用多个线程。之后,您可以连接字典,然后将这些字典加载到 pandas 中,或者将单个字典加载到 pandas 中并在那里连接。 但是,pandas 不支持asyncio,因此您在某些时候会出现性能损失。

    【讨论】:

      【解决方案4】:

      很遗憾,由于 Python 中的 GIL(全局解释器锁),多个线程不会同时运行 - 所有线程都使用同一个 CPU 的核心。这意味着如果您创建多个线程来加载文件,总时间将等于(或实际上更长)一个接一个地加载这些文件所需的时间。

      更多关于 GIL:https://wiki.python.org/moin/GlobalInterpreterLock

      为了加快加载时间,您可以尝试从 csv/excel 切换到 pickle 文件(或 HDF)。

      【讨论】:

      • 多处理绕过了 GIL,所以这是一个无效的反对意见。
      • 我说的是多线程,而不是多处理。对于多处理,您需要在消耗内存的进程之间传输大量数据。
      • 对于 I/O 密集型任务,GIL 的问题较少。它正在阻塞 CPU 密集型任务。
      【解决方案5】:

      您提供了硬件详细信息,但没有提供最有趣的部分:您拥有的磁盘数量、您拥有的 RAID 类型以及您正在读取的文件系统。

      如果您只有一个磁盘、没有 RAID 和一个常规文件系统(ext4、XFS 等),就像您通常在笔记本电脑上一样,您将无法简单地通过抛出 CPU(多线程或多进程)来增加带宽) 的问题。使用多线程或异步 I/O 将有助于稍微掩盖延迟,但不会增加带宽,因为您可能已经用单个读取器进程饱和了它。

      因此,使用@Cezary.Sz 建议的代码,尝试将其中一个文件移动到 USB3.0 外部存储或 SDSX 存储。如果您在大型工作站上运行,请查看硬件详细信息以查看是否有多个磁盘可用,如果您在大型集群上运行,请查找并行文件系统(BeeGFS、Lustre 等)

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2012-03-28
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2021-12-12
        相关资源
        最近更新 更多