【问题标题】:Update multiple rows of SQL table from Python script从 Python 脚本更新多行 SQL 表
【发布时间】:2020-03-07 12:34:44
【问题描述】:

我有一个巨大的表(超过 100B 条记录),我在其中添加了一个空列。如果所需的字符串可用,我会从另一个字段(字符串)解析字符串,从该字段中提取一个整数,并希望在新列中为具有该字符串的所有行更新它。

目前,在数据被解析并本地保存在数据框中之后,我对其进行迭代以使用干净的数据更新 Redshift 表。这需要大约 1 秒/次迭代,这太长了。

我当前的代码示例:

conn = psycopg2.connect(connection_details)
cur = conn.cursor()
clean_df = raw_data.apply(clean_field_to_parse)
for ind, row in clean_df.iterrows():
  update_query = build_update_query(row.id, row.clean_integer1, row.clean_integer2)
  cur.execute(update_query)

其中update_query 是生成更新查询的函数:

def update_query(id, int1, int2):
  query = """
  update tab_tab
  set 
  clean_int_1 = {}::int,
  clean_int_2 = {}::int,
  updated_date = GETDATE()
  where id = {}
  ;
  """
  return query.format(int1, int2, id)

clean_df 的结构如下:

id . field_to_parse . clean_int_1 . clean_int_2
1  . {'int_1':'2+1'}.      3      .    np.nan
2  . {'int_2':'7-0'}.     np.nan  .      7

有没有办法批量更新特定的表字段,这样就不需要一次执行一个查询了?

我正在解析字符串并从 Python 运行更新语句。数据库存储在 Redshift 上。

【问题讨论】:

  • 100B 是 1000 亿?您真的不想 逐行解析该 df,更不用说运行单个查询了。 parse_single_row() 是做什么的?任何解决方案都必须先解决这个问题,然后再进行批量插入
  • pandas 的重点不是迭代数据帧。单独进行批处理可能会快几个数量级,那是在您提高表的更新效率之前。无论如何,您应该包含该功能,我无法仅从描述中理解。此外,“这将需要比几天更长的时间”是对今年的轻描淡写:P
  • 谢谢。另外,很抱歉,我误读了您的最后一条评论
  • 所以,批处理DF,将批处理中的列作为一个整体处理,将批处理转换为内存文件,创建临时表,将文件复制到临时表,然后跨连接到主表。有点花哨,但希望明显更快

标签: python sql database pandas amazon-redshift


【解决方案1】:

如前所述,考虑纯 SQL 并避免迭代数十亿行,方法是将 Pandas 数据帧作为临时表推送到 Postgres,然后跨两个表运行一个 UPDATE。使用 SQLAlchemy,您可以使用 DataFrame.to_sql 创建数据框的表副本。甚至添加连接字段的索引 id,并在末尾删除非常大的临时表。

from sqlalchemy import create_engine

engine = create_engine("postgresql+psycopg2://myuser:mypwd!@myhost/mydatabase")

# PUSH TO POSTGRES (SAME NAME AS DF)
clean_df.to_sql(name="clean_df", con=engine, if_exists="replace", index=False)

# SQL UPDATE (USING TRANSACTION)
with engine.begin() as conn:     

    sql = "CREATE INDEX idx_clean_df_id ON clean_df(id)"
    conn.execute(sql)

    sql = """UPDATE tab_tab t
             SET t.clean_int_1 = c.int1,
                 t.clean_int_2 = c.int2,
                 t.updated_date = GETDATE()
             FROM clean_df c
             WHERE c.id = t.id
          """
    conn.execute(sql)

    sql = "DROP TABLE IF EXISTS clean_df"
    conn.execute(sql)

engine.dispose()

【讨论】:

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