快速啟動:使用 Python 的 mssql-python 驅動程式批量複製

在這個快速入門中,你可以使用 mssql-python 驅動程式在資料庫間批量複製資料。 應用程式會使用 Apache Arrow 從來源資料庫結構下載資料表到本地 Parquet 檔案,然後使用高效能 bulkcopy_arrow 方法上傳至目標資料庫。 你可以利用此模式在 SQL Server、Azure SQL 資料庫與 Fabric 中的 SQL 資料庫之間遷移、複製或轉換資料。

在 Windows 電腦上,mssql-python 驅動程式不需要任何外部相依性。 驅動程式會透過單一 pip 安裝來安裝所需的所有內容,讓您可以將最新版本的驅動程式用於新腳本,而不會中斷您沒有時間升級和測試的其他腳本。

MSSQL-Python 文件 | MSSQL-Python 原始碼 | 套件(PyPI) | uv

先決條件

  • Python 3.10 或更新版本

  • 如果你還沒有 Python,建議安裝 Python 執行環境 和 Pip 套件管理器 ,從 python.org 安裝。

  • 不想使用自己的系統環境? 遵循容器與本地開發,建立可重現的開發容器或 GitHub Codespaces 環境。

  • Visual Studio Code 使用下列擴充套件:

  • Azure Command-Line 介面(CLI)用於 macOS 和 Linux 的無密碼認證。

  • 如果你還沒有安裝<該軟體>,請依照 安裝說明進行操作。

  • 一個包含AdventureWorksLT範例結構和有效 連接字串 的原始資料庫。

  • 具備有效連線字串的目的地資料庫。 使用者必須擁有建立及寫入資料表的權限。 如果你沒有第二個資料庫,也可以用同一個資料庫,但目的資料表用不同的結構。

安裝一次性作業系統特定先決條件。 Windows 使用者可以跳過此步驟。 完整平台細節請參見 安裝 mssql-python。

apk add libtool krb5-libs krb5-dev

建立 SQL 資料庫

在以下平台建立或連接 SQL 資料庫:

建立專案並執行程式碼

  1. 建立新專案
  2. 新增相依性
  3. 啟動 Visual Studio Code
  4. 更新 pyproject.toml
  5. 更新 main.py
  6. 保存連接線
  7. 使用 uv run 執行腳本

建立新專案

  1. 在開發目錄中開啟命令提示字元。 如果你沒有,請建立一個新的目錄,例如 python 或 scripts。 避免在 OneDrive 上放置資料夾,因為同步可能會干擾虛擬環境的管理。

  2. 建立新的專案 使用uv。

    uv init mssql-python-bcp-qs
    cd mssql-python-bcp-qs
    

新增依賴性

在相同的目錄中,安裝 mssql-python、 python-dotenv和 pyarrow 套件。

uv add mssql-python python-dotenv pyarrow

啟動 Visual Studio Code

在相同的目錄中,執行下列命令。

code .

更新 pyproject.toml

  1. pyproject.toml 包含專案的中繼資料。 在您最喜歡的編輯器中打開該文件。

  2. 檢閱檔案的內容。 它應該類似於這個例子。 請注意 Python 的版本和相依性,mssql-python 使用 >= 來定義最低版本。 如果您偏好確切的版本,請將版本號碼之前的 變更 >= 為 ==。 然後,每個套件的解析版本會儲存在 uv.lock 中。 鎖定檔案可確保處理專案的開發人員使用一致的套件版本。 它也會確保在將套件分發給一般使用者時,使用完全相同的套件版本集。 您不應該編輯檔案 uv.lock 。

    [project]
    name = "mssql-python-bcp-qs"
    version = "0.1.0"
    description = "Add your description here"
    readme = "README.md"
    requires-python = ">=3.11"
    dependencies = [
        "mssql-python>=1.15.0",
        "python-dotenv>=1.1.1",
        "pyarrow>=19.0.0",
    ]
    
  3. 更新描述以更具描述性。

    description = "Bulk copies data between SQL databases using mssql-python and Apache Arrow"
    
  4. 儲存並關閉檔案。

更新 main.py

  1. 開啟名為 main.py的檔案。 它應該類似於這個例子。

    def main():
        print("Hello from mssql-python-bcp-qs!")
    
    if __name__ == "__main__":
        main()
    
  2. 將 main.py 的內容替換為以下程式碼區塊。 每個區塊都是建立在前一個區塊的基礎上,應該依序排列 main.py 。

    小提示

    如果 Visual Studio Code 在解析套件時遇到問題,您必須 更新解譯器才能使用虛擬環境。

  3. 在 main.py 的頂部加入導入和常數。 該腳本使用 mssql_python 進行資料庫連線與 Arrow 擷取,使用 pyarrow 和 pyarrow.parquet 進行欄式資料處理與 Parquet 檔案 I/O,使用 python-dotenv 從 .env 檔案載入連線字串,並使用已編譯的正規表示式模式驗證 SQL 識別碼,以防止 SQL 注入。

    """Round-trip: download tables from a source DB/schema to parquet, upload to a destination DB/schema."""
    
    import os
    import re
    import time
    
    import pyarrow as pa
    import pyarrow.parquet as pq
    from dotenv import load_dotenv
    import mssql_python
    
    BATCH_SIZE = 64_000
    _SAFE_IDENT = re.compile(r"^[A-Za-z0-9_]+$")
    
    
    def _validate_ident(name: str) -> str:
        if not _SAFE_IDENT.match(name):
            raise ValueError(f"Unsafe SQL identifier: {name!r}")
        return name
    
  4. 在匯入的下方,新增從 SQL 到 Arrow 類型的映射。 此字典會將 SQL Server 欄位類型轉換為其 Apache Arrow 對應的欄位,以確保在寫入 Parquet 時保持資料的完整性。 輔助函式會從NVARCHAR(100)元資料建立精確的 SQL 型別字串(例如 DECIMAL(18,2) 或 INFORMATION_SCHEMA),並解析每個欄位的對應箭頭型別。 這些類型會以欄位中繼資料儲存在 Parquet 檔案中,以便以精確的欄位定義重建目標資料表。

    _SQL_TO_ARROW = {
        "bit": pa.bool_(),
        "tinyint": pa.uint8(),
        "smallint": pa.int16(),
        "int": pa.int32(),
        "bigint": pa.int64(),
        "float": pa.float64(),
        "real": pa.float32(),
        "smallmoney": pa.decimal128(10, 4),
        "money": pa.decimal128(19, 4),
        "date": pa.date32(),
        "datetime": pa.timestamp("us"),
        "datetime2": pa.timestamp("us"),
        "smalldatetime": pa.timestamp("s"),
        "uniqueidentifier": pa.string(),
        "xml": pa.string(),
        "image": pa.binary(),
        "binary": pa.binary(),
        "varbinary": pa.binary(),
        "timestamp": pa.binary(),
    }
    
    
    def _sql_type_str(data_type: str, max_length: int, precision: int, scale: int) -> str:
        """Build the exact SQL type string from INFORMATION_SCHEMA metadata."""
        dt = data_type.lower()
        if dt in ("char", "varchar", "nchar", "nvarchar", "binary", "varbinary"):
            length = "MAX" if max_length == -1 else str(max_length)
            return f"{dt.upper()}({length})"
        if dt in ("decimal", "numeric"):
            return f"{dt.upper()}({precision},{scale})"
        return dt.upper()
    
    
    def _arrow_type(sql_type: str, precision: int, scale: int) -> pa.DataType:
        sql_type = sql_type.lower()
        if sql_type in _SQL_TO_ARROW:
            return _SQL_TO_ARROW[sql_type]
        if sql_type in ("decimal", "numeric"):
            return pa.decimal128(precision, scale)
        if sql_type in ("char", "varchar", "nchar", "nvarchar", "text", "ntext", "sysname"):
            return pa.string()
        return pa.string()
    
  5. 加入結構內省和 DDL 生成函數。 _get_arrow_schema 使用參數化查詢對 INFORMATION_SCHEMA.COLUMNS 進行查詢,建立 Arrow 架構,並將原始 SQL 型別儲存為欄位元資料,以便目的資料表能以精確的欄位定義重建。 _create_table_ddl 會讀取該中繼資料以產生 DROP/CREATE TABLE DDL。 timestamp (rowversion)類型會被重新映射到VARBINARY(8),因為它是自動產生的,無法直接插入。

    def _get_arrow_schema(cursor, schema_name: str, table_name: str) -> pa.Schema:
        """Build an Arrow schema from INFORMATION_SCHEMA.COLUMNS.
    
        Stores the original SQL type as field metadata so the round-trip
        CREATE TABLE can reproduce exact column definitions.
        """
        cursor.execute(
            "SELECT COLUMN_NAME, DATA_TYPE, "
            "COALESCE(CHARACTER_MAXIMUM_LENGTH, 0), "
            "COALESCE(NUMERIC_PRECISION, 0), "
            "COALESCE(NUMERIC_SCALE, 0), "
            "IS_NULLABLE "
            "FROM INFORMATION_SCHEMA.COLUMNS "
            "WHERE TABLE_SCHEMA = ? AND TABLE_NAME = ? "
            "ORDER BY ORDINAL_POSITION",
            (schema_name, table_name),
        )
        rows = cursor.fetchall()
        if not rows:
            raise ValueError(f"No columns found for {schema_name}.{table_name}")
        fields = []
        for col_name, data_type, max_len, precision, scale, nullable in rows:
            arrow_t = _arrow_type(data_type, precision, scale)
            sql_t = _sql_type_str(data_type, max_len, precision, scale)
            fields.append(
                pa.field(
                    col_name, arrow_t,
                    nullable=(nullable == "YES"),
                    metadata={"sql_type": sql_t},
                )
            )
        return pa.schema(fields)
    
    
    def _create_table_ddl(target: str, schema: pa.Schema) -> str:
        """Build DROP/CREATE TABLE DDL from Arrow schema with SQL type metadata."""
        col_defs = []
        for f in schema:
            sql_t = f.metadata[b"sql_type"].decode()
            # timestamp/rowversion is auto-generated and not insertable
            if sql_t == "TIMESTAMP":
                sql_t = "VARBINARY(8)"
            null = "" if f.nullable else " NOT NULL"
            col_defs.append(f"[{f.name}] {sql_t}{null}")
        col_defs_str = ",\n    ".join(col_defs)
        return (
            f"IF OBJECT_ID('{target}', 'U') IS NOT NULL DROP TABLE {target};\n"
            f"CREATE TABLE {target} (\n    {col_defs_str}\n);"
        )
    
  6. 新增下載功能。 download_table 使用 cursor.arrow_batch(),在驅動程式的 C++ 層直接以 Arrow 記錄批次格式擷取資料,從而避免建立中間的 Python 物件。 每個批次都會轉換為來自 _get_arrow_schema 的中繼資料增強結構描述,讓原始 SQL 類型(例如 NVARCHAR(100))得以保留在 Parquet 檔案中。 此函式使用兩個獨立游標:一個用於讀取欄位元資料,另一個用於串流資料。

    def download_table(conn, schema_name: str, table_name: str, parquet_file: str) -> int:
        """Download a SQL table to a parquet file. Returns row count (0 if empty)."""
        _validate_ident(schema_name)
        _validate_ident(table_name)
        source = f"{schema_name}.[{table_name}]"
    
        with conn.cursor() as cursor:
            schema = _get_arrow_schema(cursor, schema_name, table_name)
    
        row_count = 0
        t0 = time.perf_counter()
    
        with conn.cursor() as cursor:
            cursor.execute(f"SELECT * FROM {source}")
            writer = None
            try:
                while True:
                    batch = cursor.arrow_batch(BATCH_SIZE)
                    if batch.num_rows == 0:
                        break
                    # Cast to the schema to preserve SQL type metadata in Parquet
                    arrays = [
                        batch.column(i).cast(schema.field(i).type)
                        for i in range(batch.num_columns)
                    ]
                    batch = pa.record_batch(arrays, schema=schema)
                    if writer is None:
                        writer = pq.ParquetWriter(parquet_file, schema)
                    writer.write_batch(batch)
                    row_count += batch.num_rows
            finally:
                if writer is not None:
                    writer.close()
    
        if row_count == 0:
            return 0
    
        elapsed = time.perf_counter() - t0
        rate = f"{int(row_count / elapsed):,} rows/sec" if elapsed > 0 else "n/a"
        print(
            f"{schema_name}.{table_name} -> {parquet_file}: {row_count:,} rows downloaded "
            f"in {elapsed:.2f}s ({rate})"
        )
        return row_count
    
  7. 添加擴充掛勾。 enrich_parquet 是一個佔位符,你可以在資料上傳前加入轉換、派生欄位或連接。 在這個快速入門中,它是不執行任何操作,會回傳不變的檔案路徑。

    def enrich_parquet(parquet_file: str) -> str:
        """Enrich a parquet file before upload. Returns the (possibly new) file path."""
        # TODO: add transformations, derived columns, or joins
        print(f"Enriching {parquet_file} (no-op)")
        return parquet_file
    
  8. 新增上傳功能。 upload_parquet 從 Parquet 檔案讀取 Arrow 架構,產生並執行 DROP/CREATE TABLE DDL 以準備目的地,然後將檔案的記錄批次串流成單一 cursor.bulkcopy_arrow() 呼叫,以實現高效能的批量插入。 由於 Parquet 批次本身就是 Apache Arrow 的記錄批次,因此此方法可直接載入這些批次,而無須先將每個值轉換為 Python 物件。 此 table_lock=True 選項透過最小化鎖爭用以提升吞吐量。 該方法會回傳複製的列數與時間,然後函式執行 a SELECT COUNT(*) ,若目標列數與上傳的列數不符,則會產生錯誤。

    def upload_parquet(conn, parquet_file: str, target: str) -> int:
        """Upload a parquet file into a SQL table via BCP. Returns row count."""
        # ── Create target table from parquet schema ──
        pf_schema = pq.read_schema(parquet_file)
        with conn.cursor() as cursor:
            cursor.execute(_create_table_ddl(target, pf_schema))
        conn.commit()
    
        # ── Bulk insert ──
        with pq.ParquetFile(parquet_file) as pf:
            with conn.cursor() as cursor:
                result = cursor.bulkcopy_arrow(
                    target, pf.iter_batches(batch_size=BATCH_SIZE),
                    batch_size=BATCH_SIZE, table_lock=True, timeout=3600,
                )
        uploaded = result["rows_copied"]
    
        # ── Verify ──
        with conn.cursor() as cursor:
            cursor.execute(f"SELECT COUNT(*) FROM {target}")
            count = cursor.fetchone()[0]
        if count != uploaded:
            raise ValueError(
                f"Row count mismatch for {target}: uploaded {uploaded:,}, destination has {count:,}"
            )
    
        print(
            f"{parquet_file} -> {target}: {uploaded:,} rows uploaded "
            f"in {result['elapsed_time']:.2f}s "
            f"({result['rows_per_second']:,.0f} rows/sec, {result['batch_count']} batches) "
            f"| destination rows: {count:,}"
        )
        return uploaded
    

    小提示

    將批次迭代器傳遞給單次 bulkcopy_arrow 呼叫,而非針對每個批次各呼叫一次該方法。 該方法會自行開啟連線,並在呼叫返回時將其關閉,因此以每個批次為單位的迴圈都會為每個批次登入一次,並取得一次資料表鎖定。

  9. 加入編排功能。 transfer_tables 將三個階段串連起來。 它會連接至來源資料庫,透過 INFORMATION_SCHEMA.TABLES 找出指定結構描述中的所有基礎資料表,將每個資料表下載為本機 Parquet 檔案,執行擴充 hook,然後連接至目標資料庫並上傳每個檔案。

    def transfer_tables(
        source_conn_str: str,
        dest_conn_str: str,
        source_schema: str,
        dest_schema: str,
    ) -> None:
        """Download all tables from source DB/schema to parquet, upload to dest DB/schema."""
        _validate_ident(source_schema)
        _validate_ident(dest_schema)
    
        parquet_dir = source_schema
        os.makedirs(parquet_dir, exist_ok=True)
    
        # ── Download from source ──
        with mssql_python.connect(source_conn_str) as src_conn:
            with src_conn.cursor() as cursor:
                cursor.execute(
                    "SELECT TABLE_NAME FROM INFORMATION_SCHEMA.TABLES "
                    "WHERE TABLE_SCHEMA = ? AND TABLE_TYPE = 'BASE TABLE' "
                    "ORDER BY TABLE_NAME",
                    (source_schema,),
                )
                tables = [row[0] for row in cursor.fetchall()]
    
            print(f"Found {len(tables)} {source_schema} tables: {', '.join(tables)}\n")
    
            parquet_files = []
            for table_name in tables:
                parquet_file = os.path.join(parquet_dir, f"{table_name}.parquet")
                row_count = download_table(src_conn, source_schema, table_name, parquet_file)
                if row_count == 0:
                    print(f"{source_schema}.{table_name}: empty, skipping")
                else:
                    parquet_files.append((table_name, parquet_file))
    
        # ── Enrich parquet files ──
        enriched = []
        for table_name, parquet_file in parquet_files:
            enriched.append((table_name, enrich_parquet(parquet_file)))
    
        # ── Upload to destination ──
        with mssql_python.connect(dest_conn_str) as dest_conn:
            for table_name, parquet_file in enriched:
                target = f"{dest_schema}.[{table_name}]"
                upload_parquet(dest_conn, parquet_file, target)
    
  10. 最後,加入 main 入口點。 它會載 .env 入檔案,呼叫 transfer_tables 來源與目的連接字串,並列印總經過時間。

    def main():
        load_dotenv()
        t_start = time.perf_counter()
    
        transfer_tables(
            source_conn_str=os.environ["SOURCE_CONNECTION_STRING"],
            dest_conn_str=os.environ["DEST_CONNECTION_STRING"],
            source_schema="SalesLT",
            dest_schema="dbo",
        )
    
        print(f"Total: {time.perf_counter() - t_start:.2f}s")
    
    
    if __name__ == "__main__":
        main()
    
  11. 儲存後關閉 main.py。

保存連接線

  1. 開啟 .gitignore 檔案,並為 .env 檔案新增排除項目。 您的檔案應該類似於此範例。 請務必儲存並在完成後將其關閉。

    # Python-generated files
    __pycache__/
    *.py[oc]
    build/
    dist/
    wheels/
    *.egg-info
    
    # Virtual environments
    .venv
    
    # Connection strings and secrets
    .env
    
  2. 在目前目錄中,建立名為 .env的新檔案。

  3. 在 .env 檔案中,新增您來源與目的地連接字串的項目。 把佔位符值替換成你實際的伺服器和資料庫名稱。

    SOURCE_CONNECTION_STRING="Server=<source_server_name>;Database=<source_database_name>;Encrypt=yes;TrustServerCertificate=no;Authentication=ActiveDirectoryInteractive"
    DEST_CONNECTION_STRING="Server=<dest_server_name>;Database=<dest_database_name>;Encrypt=yes;TrustServerCertificate=no;Authentication=ActiveDirectoryInteractive"
    

    小提示

    此處使用的連接字串很大程度上取決於您要連線的 SQL 資料庫類型。 如果您要連線到 Azure SQL 資料庫 或 Fabric 中的 SQL 資料庫,請使用 [連接字串] 索引標籤中的 ODBC 連接字串。您可能需要根據您的案例調整驗證類型。 如需連接字串及其語法的詳細資訊,請參閱 連接字串語法參考。

小提示

在 macOS 上,兩者都ActiveDirectoryInteractiveActiveDirectoryDefault適用於 Microsoft Entra 認證。 ActiveDirectoryInteractive 每次執行腳本時都會提示你登入。 為避免重複登入提示,請透過 Azure CLI 執行 az login,然後使用 ActiveDirectoryDefault,重複使用快取的憑證。

使用 uv run 執行腳本

  1. 在先前的終端機視窗中,或開啟至相同目錄的新終端機視窗中,執行下列命令。

     uv run main.py
    

    以下是指令碼完成時的預期輸出。

    Found 12 SalesLT tables: Address, Customer, CustomerAddress, ...
    
    SalesLT.Address → SalesLT/Address.parquet: 450 rows downloaded in 0.15s (3,000 rows/sec)
    ...
    SalesLT/Address.parquet → dbo.[Address]: 450 rows uploaded in 0.10s (4,500 rows/sec) | verified: 450
    ...
    Total: 2.35s
    
  2. 透過 VS Code 的 MSSQL 擴充功能連接目標資料庫,並確認資料表與資料已成功建立。

  3. 若要將指令碼部署到另一台電腦,請將資料夾以外的 .venv 所有檔案複製到另一部電腦。 虛擬環境會在第一次執行時重新建立。

程式碼運作方式

該應用程式可分三個階段進行完整的往返資料傳輸:

  1. 下載:連接原始資料庫,讀取欄位 INFORMATION_SCHEMA.COLUMNS中繼資料,建立 Apache Arrow 架構,然後將每個資料表下載到本地的 Parquet 檔案中。
  2. Enrich (可選):提供一個 hook (enrich_parquet),可在上傳前新增轉換、衍生欄位或連接。
  3. 上傳:批次讀取每個 Parquet 檔案,利用 Arrow 結構元資料產生的 DDL 在目標資料庫中重新建資料表,然後使用 cursor.bulkcopy_arrow() 進行高效能的批量插入。 由於原始碼已是 Arrow 格式,記錄批次會傳遞給驅動程式,而不會轉換成 Python 物件。

後續步驟