PostgreSQL到Azure SQL大字符串列数据高效迁移方案咨询(附当前迁移脚本)
PostgreSQL到Azure SQL大字符串列数据高效迁移方案咨询(附当前迁移脚本)
我现在在做一个数据迁移项目,需要把PostgreSQL里的3700万条数据迁移到Azure SQL。其中有一列是存为jsonb类型的JSON payload,长度经常超过4000字符;另外还要解析这些payload提取特定元素,需要做一些数据转换。
我写了个脚本,目前迁移速度大概是每秒300条左右,而且Azure SQL的DTU已经跑到100%了。本来我用了execute_many,以为速度会不错——毕竟只迁移两列,之前用类似批量插入的方法(execute_many/批量插入)能跑到每秒1万条,但这次却只有240-300条/秒,算下来要30多个小时才能完成,这速度实在太慢了,完全达不到要求。
我现在的脚本里暂时没加JSON解析逻辑,先想解决大字段列的快速迁移这个核心问题。我想知道有没有更高效的批量插入方法?比如用DataFrame?或者换其他技术?还是说我的脚本已经够高效了,瓶颈就是Azure SQL的DTU?我的目标是提升插入速度,缩短整体迁移时间。
下面是我当前的迁移脚本,先只做数据传输,没包含解析步骤:
import time import psycopg2 import pyodbc import config # Configuration # PostgreSQL connection details # Azure SQL connection details # Source and target info # Number of rows per batch BATCH_SIZE = 1000 COLUMNS = ["id", "payload"] # ------------------------------------------------------------------ # PostgreSQL: Connect and create server-side cursor # ------------------------------------------------------------------ def get_postgres_cursor(): conn = psycopg2.connect( host=PG_HOST, port=PG_PORT, dbname=PG_DBNAME, user=PG_USER, password=PG_PASSWORD ) conn.autocommit = False cursor = conn.cursor(name="statement_cursor") cursor.execute(f""" SELECT id, payload::text AS payload FROM {PG_SCHEMA}.{PG_TABLE} ORDER BY statement_id """) return conn, cursor # ------------------------------------------------------------------ # Azure SQL: Connect (pyodbc) # ------------------------------------------------------------------ def get_azure_cursor(): conn_str = ( "DRIVER={ODBC Driver 18 for SQL Server};" f"SERVER={AZURESQL_HOST};" f"DATABASE={AZURESQL_DBNAME};" f"UID={AZURESQL_USER};" f"PWD={AZURESQL_PASSWORD};" ) conn = pyodbc.connect(conn_str, autocommit=False) cursor = conn.cursor() cursor.fast_executemany = True # Handle large string columns cursor.setinputsizes([ (pyodbc.SQL_WVARCHAR, 255, 0), # 'id' with max length 255 (pyodbc.SQL_WVARCHAR, 0, 0) # 'payload' as large text ]) return conn, cursor # ------------------------------------------------------------------ # Migration Function # ------------------------------------------------------------------ def migrate_data(): pg_conn, pg_cursor = get_postgres_cursor() azure_conn, azure_cursor = get_azure_cursor() # Build the INSERT statement column_list = ", ".join(COLUMNS) placeholders = ", ".join(["?"] * len(COLUMNS)) insert_sql = f"INSERT INTO {AZURE_SCHEMA}.{AZURE_TABLE} ({column_list}) VALUES ({placeholders})" total_migrated = 0 start_time = time.time() print("Starting migration...") try: while True: # Fetch next batch from PostgreSQL rows = pg_cursor.fetchmany(BATCH_SIZE) if not rows: break # no more data # Convert rows to list of tuples insert_data = list(rows) try: azure_cursor.executemany(insert_sql, insert_data) azure_conn.commit() except pyodbc.Error as ex: azure_conn.rollback() print(f"Error in batch insert: {ex}") # Optionally, implement retry logic or log the failed batch # Report progress total_migrated += len(insert_data) elapsed = time.time() - start_time speed = total_migrated / (elapsed if elapsed else 1.0) print(f"Migrated {total_migrated} rows at ~{speed:.1f} rows/sec") except Exception as e: print(f"An error occurred during migration: {e}") finally: # Cleanup PostgreSQL pg_cursor.close() pg_conn.commit() pg_conn.close() # Cleanup Azure SQL azure_cursor.close() azure_conn.close() print(f"Done! Total rows migrated: {total_migrated}.") if __name__ == "__main__": migrate_data()
备注:内容来源于stack exchange,提问作者user29370337
相关产品推荐
相关产品推荐

