Google Cloud Run批量Upsert大CSV至PostgreSQL性能骤降求助
问题描述
我有一个大型价格CSV文件,需要Upsert到PostgreSQL表中。本地通过Pandas分块读取,先上传至临时表,再将临时表数据同步至目标表后删除临时表,性能表现优异:每个块耗时18秒,1800万行文件仅需14分钟。但在Google Cloud Run环境中,前3个块耗时约17秒,第4个块增至42秒,第5个块达102秒,后续块耗时稳定在100-120秒甚至更久。
配置信息:
- Cloud Run:2核CPU、并发数200、1GB内存,仅在请求处理时分配CPU
- PostgreSQL实例:1核CPU、628MB内存
代码实现
主执行代码
import time import pandas as pd from sqlalchemy.pool import NullPool from aldjemy.core import get_engine import uuid from typing import List import pandas as pd import sqlalchemy as sa engine = get_engine(connect_args={'connect_timeout': 3600, 'pool_size': NullPool}) i = 0 with pd.read_csv("prices.csv", chunksize=350000) as reader: for prices in reader: time.sleep(0.01) i=i+1 print("chunk number %s" %str(i)) t = time.time() try: _ = df_upsert(data_frame=prices, engine=engine, table_name=ModelPrice._meta.db_table, match_columns=['id', 'date']) except Exception as e: raise Exception(f"{self.__class__.__name__}: {str(e)}") elapsed = time.time() - t print('price data upserted. TOTAL elapsed = %d' %elapsed)
df_upsert函数实现
def df_upsert(data_frame: pd.DataFrame, table_name: str, engine: sa.engine.Engine, match_columns: List[str]=None): """ Perform an "upsert" on a PostgreSQL table from a DataFrame. Constructs an INSERT … ON CONFLICT statement, uploads the DataFrame to a temporary table, and then executes the INSERT. Parameters ---------- data_frame : pandas.DataFrame The DataFrame to be upserted. table_name : str The name of the target table. Note that this string value is injected directly into the SQL statements, so proper quoting is required for table names that contain spaces, etc. A schema can be specified as well, e.g., 'my_schema."my table"'. engine : sa.engine.Engine The SQLAlchemy Engine to use. match_columns : list of str, optional A list of the column name(s) on which to match. If omitted, the primary key columns of the target table will be used. Note that these names *are* automatically quoted in the INSERT statement, so do not quote them in this list, e.g., ["my column"], not ['"my column"']. """ df_columns = list(data_frame.columns) if not match_columns: insp = sa.inspect(engine) match_columns = insp.get_pk_constraint(table_name)[ "constrained_columns" ] temp_table_name = f"temp_{uuid.uuid4().hex[:6]}" columns_to_update = [col for col in df_columns if col not in match_columns] insert_col_list = ", ".join([f'"{col_name}"' for col_name in df_columns]) stmt = f"INSERT INTO {table_name} ({insert_col_list})\n" stmt += f"SELECT {insert_col_list} FROM {temp_table_name}\n" match_col_list = ", ".join([f'"{col}"' for col in match_columns]) stmt += f"ON CONFLICT ({match_col_list}) DO UPDATE SET\n" stmt += ", ".join( [f'"{col}" = EXCLUDED."{col}"' for col in columns_to_update] ) with engine.begin() as conn: t = time.time() conn.exec_driver_sql( f"CREATE TEMPORARY TABLE {temp_table_name} AS SELECT * FROM {table_name} WHERE false" ) elapsed = time.time() - t print('TEMPORARY TABLE created. elapsed = %d' %elapsed) t = time.time() data_frame.to_sql(temp_table_name, conn, if_exists="append", index=False) elapsed = time.time() - t print('dataframe written to TEMPORARY TABLE. elapsed = %d' %elapsed) t = time.time() conn.exec_driver_sql(stmt) elapsed = time.time() - t print('TEMPORARY TABLE copied into price table. elapsed = %d' %elapsed) t = time.time() engine.execute(f'DROP TABLE "{temp_table_name}"') elapsed = time.time() - t print('TEMPORARY TABLE dropped. elapsed = %d' %elapsed)
已尝试的无效优化措施
- 减小分块大小至10万,性能反而更差
- 删除价格表并重建索引后在空表上执行Upsert
- 提升Cloud Run的CPU与内存配置
- 为Cloud Run配置持续分配CPU
时间拆解细节
前3个块
- 临时表创建:0秒
- 数据写入临时表:9秒
- 临时表同步至价格表:7-9秒
- 临时表删除:0秒
后续块
- 临时表创建:0秒
- 数据写入临时表:48秒
- 临时表同步至价格表:74秒
- 临时表删除:0秒
内容的提问来源于stack exchange,提问作者Courvoisier
相关产品推荐
相关产品推荐

