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

在這個快速入門中,你可以使用 mssql-python 驅動程式在資料庫間批量複製資料。 應用程式會使用 Apache Arrow 從來源資料庫結構下載資料表到本地 Parquet 檔案,然後使用高效能 bulkcopy 方法上傳至目標資料庫。 你可以利用此模式在 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. 在開發目錄中開啟命令提示字元。 如果你沒有,請建立一個新的目錄,例如 pythonscripts。 避免在 OneDrive 上放置資料夾,因為同步可能會干擾虛擬環境的管理。

  2. 建立新的專案 使用uv

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

新增依賴性

在相同的目錄中,安裝 mssql-pythonpython-dotenvpyarrow 套件。

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.5.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 擷取,使用 pyarrowpyarrow.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() 進行高效能的批量插入。 此 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 ──
        uploaded = 0
        t0 = time.perf_counter()
        with pq.ParquetFile(parquet_file) as pf:
            with conn.cursor() as cursor:
                for batch in pf.iter_batches(batch_size=BATCH_SIZE):
                    rows = zip(*(col.to_pylist() for col in batch.columns))
                    cursor.bulkcopy(
                        target, rows, batch_size=BATCH_SIZE,
                        table_lock=True, timeout=3600,
                    )
                    uploaded += batch.num_rows
        elapsed = time.perf_counter() - t0
    
        # ── 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:,}"
            )
    
        rate = f"{int(uploaded / elapsed):,} rows/sec" if elapsed > 0 else "n/a"
        print(
            f"{parquet_file} → {target}: {uploaded:,} rows uploaded "
            f"in {elapsed:.2f}s ({rate}) "
            f"| destination rows: {count:,}"
        )
        return uploaded
    
  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() 進行高效能的批量插入。

下一步

利用這些文章持續累積:

  • 批量複製 提供進階批量複製模式,包括欄位映射、批次大小及錯誤處理。
  • 資料載入與移動模式,用於從檔案、API 及其他資料庫載入資料的策略。
  • 效能調校以優化大型資料操作的吞吐量。