You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.14 15:08:00