【问题标题】:Pandas dataframe to Object instances array efficiency for bulk DB insert用于批量数据库插入的 Pandas 数据帧到对象实例数组效率
【发布时间】:2020-05-17 02:34:35
【问题描述】:

我有一个 Pandas 数据框,格式为:

Time    Temperature    Voltage    Current
0.0     7.8            14         56
0.1     7.9            12         58
0.2     7.6            15         55
... So on for a few hundred thousand rows...

我需要尽可能快地将数据批量插入 PostgreSQL 数据库。这是一个 Django 项目,我目前正在使用 ORM 进行数据库操作和构建查询,但如果有更有效的方法来完成任务,欢迎提出建议。

我的数据模型如下所示:

class Data(models.Model):
    time = models.DateTimeField(db_index=True)
    parameter = models.ForeignKey(Parameter, on_delete=models.CASCADE)
    parameter_value = models.FloatField()

所以Time 是DataFrame 的row[0],然后对于每个标题列,我获取与其对应的值,将标题用作parameter。所以示例表的row[0] 将在我的数据库中生成3 个Data 对象:

Data(time=0.0, parameter="Temperature", parameter_value=7.8)
Data(time=0.0, parameter="Voltage", parameter_value=14)
Data(time=0.0, parameter="Current", parameter_value=56)

我们的应用程序允许用户解析以毫秒为单位的数据文件。所以我们从一个文件中生成了很多单独的数据对象。我当前的任务是改进解析器,使其更加高效,直到我们在硬件级别达到 I/O 限制。

我当前的解决方案是遍历每一行,为time + parameter + value 上的每一行创建一个Data 对象,并将所述对象附加到一个数组中,这样我就可以通过Django 进行Data.objects.bulk_create(all_data_objects)。当然我知道这是低效的,可能会有很多改进。

使用此代码:

# Convert DataFrame to dict
df_records = df.to_dict('records')

# Start empty dta array
all_data_objects = []

# Go through each row creating objects and appending to data array
for row in df_records:
    for parameter, parameter_value in row.items():
        if parameter != "Time":
            all_data_objects.append(Data(
                    time=row["Time"],
                    parameter_value=parameter_value,
                    parameter=parameter))

# Commit data to Postgres DB
Data.objects.bulk_create(all_data)

目前整个操作,没有包括DB插入操作(写入磁盘),即只生成Data对象数组,对于一个55mb的文件,生成大约600万个单独的@ 987654334@ 个对象大约需要 370 秒。仅df_records = df.to_dict('records') 行就需要 83 秒。使用time.time()在每个部分的两端测量时间并计算差异。

我该如何改善这些时间?

【问题讨论】:

    标签: python pandas postgresql django-database


    【解决方案1】:

    如果您真的需要快速解决方案,我建议您直接使用 pandas 来哑化表格。

    首先让我们为您的示例创建数据:

    import pandas as pd
    
    data = {
        'Time': {0: 0.0, 1: 0.1, 2: 0.2},
        'Temperature': {0: 7.8, 1: 7.9, 2: 7.6},
        'Voltage': {0: 14, 1: 12, 2: 15},
        'Current': {0: 56, 1: 58, 2: 55}
    }
    df = pd.DataFrame(data)
    

    现在您应该转换数据框,以便您拥有带有melt 的所需列:

    df = df.melt(["Time"], var_name="parameter", value_name="parameter_value")
    

    此时您应该将parameter 值映射到外部id。我将以params为例:

    params = {"Temperature": 1, "Voltage": 2, "Current": 3}
    df["parameter"] = df["parameter"].map(params)
    

    此时数据框将如下所示:

       Time  parameter  parameter_value
    0   0.0          1              7.8
    1   0.1          1              7.9
    2   0.2          1              7.6
    3   0.0          2             14.0
    4   0.1          2             12.0
    5   0.2          2             15.0
    6   0.0          3             56.0
    7   0.1          3             58.0
    8   0.2          3             55.0
    

    现在要使用 pandas 导出,您可以使用:

    import sqlalchemy as sa
    engine = sa.create_engine("use your connection data")
    df.to_sql(name="my_table", con=engine, if_exists="append", index=False)
    

    但是,当我使用它时,它的速度不足以满足我们的要求。所以我建议你使用cursor.copy_from insted,因为它更快:

    from io import StringIO
    
    output = StringIO()
    df.to_csv(output, sep=';', header=False, index=False, columns=df.columns)
    output.getvalue()
    # jump to start of stream
    output.seek(0)
    
    # Insert df into postgre
    connection = engine.raw_connection()
    with connection.cursor() as cursor:
        cursor.copy_from(output, "my_table", sep=';', null="NULL", columns=(df.columns))
        connection.commit()
    

    我们尝试了几百万次,这是使用 PostgreSQL最快的方法

    【讨论】:

    • cursor.copy_from(output, "data", sep=';', null="NULL", columns=(df.columns)) 行尝试此方法后出现错误。回溯显示:Expected bytes or unicode string, got numpy.float64 instead,我想这是因为没有在某处提供正确的数据值。 (data 是我要插入的表的名称)。我以前没有使用过这种方法,你知道这里发生了什么吗? (此时我已经走出了自己的舒适区,以前从未使用过 sqlalchemy 或 stringIO,所以我在尝试学习的同时几乎复制/粘贴了你的 sn-ps)
    • 我不太确定您为什么会收到此错误,但您似乎有一些值作为 numpy 数字而不是字符串。您有可能在数据库中将其中一列定义为字符串?我建议您在使用cursor.copy_from 替代方案之前先尝试df.to_sql 选项。这个选项可能对您来说足够快。
    • 我可以发布df.to_csv() 的输出示例,如果有帮助,将在正文中进行
    • 我也可能遗漏了一些用于复制操作的参数,例如,我的标题不是参数的名称,而是一个 ID。我该如何指定?
    • 我不确定你的意思。能否添加数据应该去的SQL表的定义?
    【解决方案2】:

    您不需要为所有行创建数据对象。 SqlAlchemy 也支持这种方式的批量插入:

    data.insert().values([
                        dict(time=0.0, parameter="Temperature", parameter_value=7.8),
                        dict(time=0.0, parameter="Voltage", parameter_value=14)
                    ])
    

    更多详情请见https://docs.sqlalchemy.org/en/13/core/dml.html?highlight=insert%20values#sqlalchemy.sql.expression.ValuesBase.values

    如果您只需要插入数据,则不需要 pandas,并且可以为您的数据文件使用不同的解析器(或编写您自己的解析器,具体取决于您的数据文件的格式)。此外,将数据集拆分为更小的部分并并行化插入命令可能是有意义的。

    【讨论】:

      猜你喜欢
      • 2015-11-06
      • 2011-03-09
      • 2020-10-21
      • 2021-08-09
      • 1970-01-01
      • 1970-01-01
      • 2016-10-13
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多