mssql-pythonで一括コピーを使いましょう

mssql-pythonドライバーには、大量データを効率的にSQL Server、Azure SQL Database、Azure SQL Managed Instance、Microsoft Fabric内のSQLデータベースに挿入する一括コピー機能が含まれています。

cursor.bulkcopy()手法は、大規模データセットのロードに高性能な経路を提供します:

  • ネットワークの往復を最小限に抑えます。
  • オプションでロード中の制約チェックを回避できます。
  • 最適化されたTDSバルクインサートプロトコルを使用しています。
  • bcp.exeおよびSqlBulkCopyに匹敵するスループットを実現しています。

Rustベースの mssql_py_core ネイティブ拡張機能が一括コピー機能を支えています。 通常のカーソル execute() パイプラインの外で動作します。

基本的な使用方法

カーソル上で bulkcopy() を呼び出し、ターゲットテーブル名と行タプルまたは Row オブジェクトの反復を渡します:

Important

同じセッション内でターゲットテーブルを作成または変更した場合は、conn.commit()前にbulkcopy()に連絡してください。 バルクコピープロトコルはテーブルメタデータを読み取るために別の内部チャネルを使用するため、未コミットのDDL変更はデッドロックやタイムアウトを引き起こすことがあります。

import mssql_python

conn = mssql_python.connect(connection_string)
cursor = conn.cursor()

# Create a temp table for the demo
cursor.execute("""
    CREATE TABLE ##BulkDemo (
        ID INT,
        Name NVARCHAR(50),
        Amount MONEY
    )
""")
conn.commit()

data = [
    (1, "Alice", 50000.00),
    (2, "Bob", 60000.00),
    (3, "Carol", 55000.00),
]

result = cursor.bulkcopy("##BulkDemo", data)
print(f"Copied {result['rows_copied']} rows in {result['batch_count']} batch(es)")
print(f"Elapsed: {result['elapsed_time']}")

値を返す

bulkcopy() 辞書を返す:

Key タイプ 説明
rows_copied int コピー成功した行数。
batch_count int 処理されたバッチ数。
elapsed_time float 手術にかかる時間は数秒だった。

メソッドシグネチャ

cursor.bulkcopy(
    table_name,                    # str – target table (can include schema, e.g. "dbo.MyTable")
    data,                          # Iterable[Tuple | Row] – rows to insert
    batch_size=0,                  # int – rows per batch; 0 = server optimal
    timeout=30,                    # int – operation timeout in seconds
    column_mappings=None,          # List[str] | List[Tuple[int,str]] | None
    keep_identity=False,           # bool – preserve identity values from source
    check_constraints=False,       # bool – check constraints during load
    table_lock=False,              # bool – use table-level lock
    keep_nulls=False,              # bool – preserve NULLs instead of defaults
    fire_triggers=False,           # bool – fire INSERT triggers on target
    use_internal_transaction=False, # bool – use internal transaction per batch
)

列マッピング

既定では、bulkcopy() は列の順序に基づいてマップします。 各データ列は同じインデックスのテーブル列にマッピングされます。 この動作を上書きするために column_mappings パラメータを使います。

列名リスト

リスト内の各位置は、ソースデータインデックスに対応しています:

result = cursor.bulkcopy(
    "##BulkDemo",
    data,
    column_mappings=["ID", "Name", "Amount"],
)

高度なフォーマット:明示的インデックスマッピング

各タプルは (source_index, target_column_name)の形を取ります。 この形式を使って列をスキップまたは順序付け替えできます:

result = cursor.bulkcopy(
    "##BulkDemo",
    data,
    column_mappings=[(0, "ID"), (1, "Name"), (2, "Amount")],
)

ファイルからのロード

bulkcopy() にジェネレーターを渡すことで、CSVファイルや他のファイル形式からデータを読み込むことができます。

CSV ファイル

import csv
import io
import mssql_python

# In production, replace io.StringIO with open("data.csv", "r", ...)
csv_data = """ID,Name,Value
1,Widget,9.99
2,Gadget,24.50
3,Gizmo,4.75
"""

def csv_row_generator(file_obj):
    """Generator that yields tuples from a CSV file object."""
    reader = csv.reader(file_obj)
    next(reader)  # Skip header
    for row in reader:
        if row:  # skip blank lines
            yield (
                int(row[0]),      # ID
                row[1],           # Name
                float(row[2]),    # Value
            )

conn = mssql_python.connect(connection_string)
cursor = conn.cursor()

cursor.execute("""
    CREATE TABLE ##CSVImport (ID INT, Name NVARCHAR(100), Value FLOAT)
""")
conn.commit()
result = cursor.bulkcopy("##CSVImport", csv_row_generator(io.StringIO(csv_data)))
print(f"Imported {result['rows_copied']} rows from CSV")

バッチ処理付きの大ファイル

batch_sizeパラメータを設定して、ドライバーが1バッチに送る行数を制御します。 この方法は大きなファイルにうまく機能します:

import csv
import io
import mssql_python

# In production, replace io.StringIO with open("large_file.csv", "r", ...)
csv_data = "\n".join(
    ["ID,Name,Value"] + [f"{i},Item {i},{i * 1.5}" for i in range(1, 201)]
)

def csv_rows(file_obj):
    reader = csv.reader(file_obj)
    next(reader)  # Skip header
    for row in reader:
        if row:
            yield (int(row[0]), row[1], float(row[2]))

conn = mssql_python.connect(connection_string)
cursor = conn.cursor()

cursor.execute("""
    CREATE TABLE ##LargeCSV (ID INT, Name NVARCHAR(100), Value FLOAT)
""")
conn.commit()
result = cursor.bulkcopy(
    "##LargeCSV",
    csv_rows(io.StringIO(csv_data)),
    batch_size=50,
)
print(f"Imported {result['rows_copied']} rows in {result['batch_count']} batches")

pandas DataFramesをロード

pandasのデータフレームをタプルのリストに変換してから bulkcopy()に渡します:

import pandas as pd
import mssql_python

df = pd.DataFrame({
    'ID': [1, 2, 3],
    'Name': ['Alice', 'Bob', 'Carol'],
    'Amount': [50000.0, 60000.0, 55000.0],
})

conn = mssql_python.connect(connection_string)
cursor = conn.cursor()

cursor.execute("""
    CREATE TABLE ##PandasDemo (ID INT, Name NVARCHAR(50), Amount MONEY)
""")
conn.commit()

data = [tuple(row) for row in df.itertuples(index=False, name=None)]
result = cursor.bulkcopy("##PandasDemo", data)

NULL値の処理

SQL None値を挿入するために、任意の列位置にNULLパスします:

cursor.execute("""
    CREATE TABLE ##NullDemo (ID INT, Name NVARCHAR(50), Amount MONEY)
""")
conn.commit()

data = [
    (1, "Alice", 50000.00),
    (2, "Bob", None),       # NULL Amount
    (3, None, 55000.00),    # NULL Name
]

cursor.bulkcopy("##NullDemo", data)

識別列

明示的な単位元値を挿入するには、 keep_identity=Trueを設定します:

cursor.execute("""
    CREATE TABLE ##IdentDemo (ID INT, Name NVARCHAR(50), Amount MONEY)
""")
conn.commit()

data = [
    (100, "Alice", 50000.00),
    (200, "Bob", 60000.00),
]

cursor.bulkcopy("##IdentDemo", data, keep_identity=True)

keep_identity=False(デフォルト)になったら、データから識別列を省略し、column_mappingsで非識別列をターゲットにしてください。

一括コピーオプション

パラメーター Default 説明
batch_size 0 バッチあたりの行数 0 サーバーに最適なサイズを選ばせます。
timeout 30 数秒でタイムアウト。
keep_identity False 元のデータから識別値を保持します。
check_constraints False ロード中にテーブルの制約を確認してください。
table_lock False 行レベルのロックではなく、テーブルレベルのロックを取得しましょう。
keep_nulls False 列のデフォルトを挿入する代わりにNULL値を保持しましょう。
fire_triggers False ターゲットテーブルでトリガー INSERT 発射。
use_internal_transaction False 各バッチを内部トランザクションで囲みます。

エラーを処理する

bulkcopy() ロードが失敗すると例外が発生します。エラーを検出するために try/except ブロックで呼び出しをラップします。 ただし、 bulkcopy() は独自の内部接続で動作し、コピーした行を独立してコミットするため、メイン接続の conn.rollback() では元に戻せません。 バッチをアトミックにするには、各バッチを独自のトランザクションでラップし、失敗した場合は自動的にロールバックする use_internal_transaction=True を設定します。

import mssql_python

conn = mssql_python.connect(connection_string)
cursor = conn.cursor()

cursor.execute("""
    CREATE TABLE ##ImportDemo (ID INT, Name NVARCHAR(50), Value FLOAT)
""")
conn.commit()

data = [
    (1, "Alice", 50000.00),
    (2, "Bob", 60000.00),
    (3, "Carol", 55000.00),
]

try:
    result = cursor.bulkcopy("##ImportDemo", data, use_internal_transaction=True)
    print(f"Successfully copied {result['rows_copied']} rows")
except (mssql_python.DatabaseError, ValueError) as e:
    # bulkcopy() commits on its own connection, so there's nothing to roll back
    # here. With use_internal_transaction=True, a failed batch is already rolled
    # back on the bulk copy connection.
    print(f"Bulk copy failed: {e}")

独自の検証ロジックを通したうえでデータをロードするには、まずステージング テーブルに一括コピーし、次にメイン接続上のトランザクション内で INSERT ... SELECT を使用してその行をターゲット テーブルに移します。 その INSERT は接続で実行されるため、検証に失敗した場合は conn.rollback() でそれを元に戻します。

Authentication

バルクコピーは独自のトークンを必要とする別の内部チャネルを使用します。 ドライバは、サポートされる認証方法のトークン取得を自動的に処理します。

マネージドアイデンティティ(ActiveDirectoryMSI)

システム割り当てまたはユーザー割り当てのマネージデンティティには Authentication=ActiveDirectoryMSI を使用してください。 この認証方法は、Azure VMS、App Service、Functions、AKSなどのAzureホストサービスに推奨されています。

import mssql_python

# System-assigned managed identity
conn = mssql_python.connect(
    "Server=<server>.database.windows.net;"
    "Database=<database>;"
    "Authentication=ActiveDirectoryMSI;"
    "Encrypt=yes"
)
cursor = conn.cursor()

cursor.execute("CREATE TABLE ##MsiDemo (ID INT, Name NVARCHAR(50))")
conn.commit()

result = cursor.bulkcopy("##MsiDemo", [(1, "Alice"), (2, "Bob")])
print(f"Copied {result['rows_copied']} rows")

ユーザー割り当て管理IDの場合、クライアントIDを接続文字列で渡します:

conn = mssql_python.connect(
    "Server=<server>.database.windows.net;"
    "Database=<database>;"
    "Authentication=ActiveDirectoryMSI;"
    "UID=<client-id>;"
    "Encrypt=yes"
)

Service principal(ActiveDirectoryServicePrincipal)

サービスプリンシパル(クライアント認証)認証には Authentication=ActiveDirectoryServicePrincipal を使いましょう。

conn = mssql_python.connect(
    "Server=<server>.database.windows.net;"
    "Database=<database>;"
    "Authentication=ActiveDirectoryServicePrincipal;"
    "UID=<application-client-id>;"
    "PWD=<client-secret>;"
    "Encrypt=yes"
)
cursor = conn.cursor()

cursor.execute("CREATE TABLE ##SpDemo (ID INT, Value FLOAT)")
conn.commit()

result = cursor.bulkcopy("##SpDemo", [(1, 1.5), (2, 2.5)])
print(f"Copied {result['rows_copied']} rows")

デフォルトの認証情報チェーン(ActiveDirectoryDefault)

ActiveDirectoryDefault 環境変数、ワークロードアイデンティティ、マネージデントIDなど、複数の認証プロバイダーを順番に試します。 ローカル開発とAzureホストサービスの両方でコード変更なしで動作します。

認証の詳細については、Microsoft Entra認証をご覧ください。

パフォーマンスに関するヒント

以下の技術は、大量コピーのスループットを最大化するのに役立ちます。

大規模なデータセットにはジェネレーターを使いましょう

ジェネレーターは、 bulkcopy() 任意の反復可能なものを受け入れるため、メモリ使用を最小限に抑えます:

def data_generator(count):
    """Generate rows without loading all into memory."""
    for i in range(count):
        yield (i, f"Item {i}", i * 1.5)

cursor = conn.cursor()
cursor.execute("""
    CREATE TABLE ##LargeDemo (ID INT, Name NVARCHAR(50), Value FLOAT)
""")
conn.commit()
result = cursor.bulkcopy("##LargeDemo", data_generator(1000))

読み込みを高速化するには、テーブルロックを使用する

同時リーダーがない場合は、 table_lock=True 設定して大きな初期負荷時のロックオーバーヘッドを減らしてください。

result = cursor.bulkcopy(
    "##LargeDemo",
    data,
    table_lock=True,
    batch_size=100000,
)

ロード時にインデックスを無効にする

一括読み込み前に非クラスタインデックスを一時的に無効にし、その後再構築してパフォーマンス向上を図る:

cursor = conn.cursor()

cursor.execute("""
    CREATE TABLE ##IndexDemo (ID INT, Name NVARCHAR(50), Value FLOAT)
""")
cursor.execute("CREATE NONCLUSTERED INDEX IX_Name ON ##IndexDemo(Name)")
conn.commit()

cursor.execute("ALTER INDEX IX_Name ON ##IndexDemo DISABLE")
conn.commit()

result = cursor.bulkcopy("##IndexDemo", data)
conn.commit()

cursor.execute("ALTER INDEX IX_Name ON ##IndexDemo REBUILD")
conn.commit()

テーブルを並列にロードする

各テーブルごとに別々の接続を開いて、ロードを同時に実行します。

import concurrent.futures

def load_table(table_name, rows):
    conn = mssql_python.connect(connection_string)
    cursor = conn.cursor()
    cursor.execute(f"CREATE TABLE {table_name} (ID INT, Name NVARCHAR(50), Value FLOAT)")
    conn.commit()
    result = cursor.bulkcopy(table_name, rows)
    conn.commit()
    conn.close()
    return result["rows_copied"]

data = [(i, f"Item {i}", i * 1.5) for i in range(100)]

with concurrent.futures.ThreadPoolExecutor(max_workers=3) as executor:
    futures = [
        executor.submit(load_table, "##Load1", data),
        executor.submit(load_table, "##Load2", data),
        executor.submit(load_table, "##Load3", data),
    ]
    for future in concurrent.futures.as_completed(futures):
        print(f"Loaded {future.result()} rows")

代替案との比較

以下の表は、バルクコピーと他のデータ挿入方法を比較しています。

Method 利用シーン パフォーマンス
cursor.bulkcopy() 大規模なデータセット(1,000行以上)を扱います。 最 速
cursor.executemany() パラメータ付きのメディアデータセット。 Moderate
cursor.execute() ループ内で 小さなデータセットで単純な論理。 最も遅い