クイック スタート: Python 用 mssql-python ドライバーを使用した一括コピー

このクイック スタートでは、 mssql-python ドライバーを使用して、データベース間でデータを一括コピーします。 アプリケーションは、Apache Arrow を使用してソース データベース スキーマからローカル Parquet ファイルにテーブルをダウンロードし、高パフォーマンスの bulkcopy_arrow メソッドを使用してコピー先データベースにアップロードします。 このパターンを使用して、Fabric の SQL Server、Azure SQL Database、SQL Database の間でデータを移行、レプリケート、または変換できます。

mssql-python ドライバーでは、Windows マシンへの外部依存関係は必要ありません。 ドライバーは、1 つの pip インストールで必要なすべてのものをインストールします。これにより、アップグレードとテストの時間がない他のスクリプトを中断することなく、新しいスクリプトに最新バージョンのドライバーを使用できます。

mssql-python のドキュメント | mssql-python ソース コード | パッケージ (PyPI) | Uv

前提条件

  • Python 3.10 以降

  • Python をまだお持ちでない場合は、python.org から Python ランタイムと pip パッケージ マネージャーをインストールします。

  • 自分の環境を使いたくないですか? コンテナやローカル開発に従って、再現可能な開発コンテナやGitHub Codespaces環境を作成しましょう。

  • Visual Studio Code と次の拡張機能:

  • macOS および Linux でのパスワードレス認証用の Azure Command-Line インターフェイス (CLI)。

  • uvがまだない場合は、インストール手順に従います。

  • AdventureWorksLTサンプルスキーマと有効な接続文字列を持つソースデータベースです。

  • 有効な接続文字列を持つ宛先データベース。 ユーザーには、テーブルを作成して書き込むアクセス許可が必要です。 もし2つ目のデータベースがなければ、同じデータベースを使い、宛先のテーブルには別のスキーマを使うことができます。

オペレーティング システム固有の 1 回限りの前提条件をインストールします。 Windowsユーザーはこのステップをスキップできます。 プラットフォームの詳細は 「Install 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. ファイルの内容を確認します。 この例のようになります。 mssql-pythonの Python のバージョンと依存関係では、>=を使用して最小バージョンを定義します。 正確なバージョンを使用する場合は、バージョン番号の前の >= を ==に変更します。 その後、各パッケージの解決済みバージョンが uv.lock に格納されます。 lockfile を使用すると、プロジェクトに取り組む開発者が一貫したパッケージ バージョンを使用できるようになります。 また、エンド ユーザーにパッケージを配布するときに、まったく同じパッケージ バージョンのセットが使用されるようになります。 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の上部に、インポートと定数を追加します。 このスクリプトは、データベース接続とArrowフェッチにmssql_python、列データ処理とParquetファイルI/Oにpyarrowとpyarrow.parquet、python-dotenvファイルからの接続文字列の読み込みに.env、そしてインジェクションを防ぐために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 への型マッピングを追加します。 このディクショナリは、PARQUET に書き込むときにデータの忠実性が維持されるように、SQL Server 列の型を Apache Arrow に変換します。 ヘルパー関数は、NVARCHAR(100)メタデータから正確なSQL型文字列(例:DECIMAL(18,2)やINFORMATION_SCHEMA)を構築し、各列に対応するArrow型を解決します。 これらの型は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 は、そのメタデータを読み取って DDL DROP/CREATE TABLE 生成します。 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ファイル内に保持されます。 この関数では、列メタデータを読み取るカーソルとデータをストリーミングするカーソルの 2 つを使用します。

    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 は、変換、派生列、またはデータへの結合をアップロード前に追加できるプレースホルダーです。 このクイック スタートでは、ファイル パスを変更せずに返す no-op です。

    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スキーマを読み込み、DDL DROP/CREATE TABLE 生成・実行して宛先を準備し、ファイルのレコードバッチを単一の cursor.bulkcopy_arrow() 呼び出しにストリーミングして高性能な一括挿入を行います。 ParquetのバッチはすでにApache Arrowのレコードバッチであるため、この方法はすべての値をPythonオブジェクトに変換せずにロードします。 table_lock=True オプションを使用すると、ロックの競合を最小限に抑えることでスループットが向上します。 このメソッドはコピーした行数とタイミングを返し、関数は 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
    

    ヒント

    バッチイテレーターを1バッチごとに1回呼ぶのではなく、単一の bulkcopy_arrow コールに渡します。 メソッドは自身の接続を開き、呼び出しが戻ると閉じるため、バッチごとにループがログインし、各バッチごとにテーブルロックを取得します。

  9. オーケストレーション関数を追加します。 transfer_tables 3 つのフェーズを結び付けます。 ソースデータベースに接続し、 INFORMATION_SCHEMA.TABLESを通じて与えられたスキーマ内のすべてのベーステーブルを発見し、それぞれをローカルのParquetファイルにダウンロードし、エンリッチメントフックを実行し、その後宛先データベースに接続して各ファイルをアップロードします。

    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 データベースの種類によって大きく異なります。 Fabric で Azure SQL Database または SQL データベースに接続する場合は、[接続文字列] タブから ODBC 接続文字列を使用します。シナリオによっては、認証の種類の調整が必要になる場合があります。 接続文字列とその構文の詳細については、 接続文字列の構文リファレンスを参照してください。

ヒント

macOS では、 ActiveDirectoryInteractive と ActiveDirectoryDefault の両方が Microsoft Entra 認証で機能します。 ActiveDirectoryInteractive スクリプトを実行するたびにサインインするように求められます。 繰り返しサインインのプロンプトを避けるために、を実行して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 フォルダーを除くすべてのファイルを他のコンピューターにコピーします。 仮想環境は、最初の実行で再作成されます。

コードのしくみ

アプリケーションは、次の 3 つのフェーズで完全なラウンドトリップ データ転送を実行します。

  1. ダウンロード: ソース データベースに接続し、 INFORMATION_SCHEMA.COLUMNSから列メタデータを読み取り、Apache Arrow スキーマをビルドしてから、各テーブルをローカル Parquet ファイルにダウンロードします。
  2. エンリッチ (省略可能): アップロードする前に変換、派生列、または結合を追加できるフック (enrich_parquet) を提供します。
  3. アップロード: 各 Parquet ファイルをバッチで読み取り、方向スキーマ メタデータから生成された DDL を使用してコピー先データベース内のテーブルを再作成した後、 cursor.bulkcopy_arrow() を使用して高パフォーマンスの一括挿入を行います。 ソースがすでにArrow形式であるため、レコードのバッチはPythonオブジェクトに変換されることなくドライバーに渡されます。

次のステップ