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をロード

DataFrameはカラム形式なので、最速パスは bulkcopy_arrow()であり、pandasがすでに生成方法を知っているArrowテーブルを消費します。 bulkcopy()行タプルを取るので、まず列をPythonオブジェクトにフラット化する必要があります。

読み込む前に、Arrowテーブルを宛先の列のデータ型にキャストしてください。 pyarrow 数値列に対して float64 を推論しますが、ドライバーはそれを お金小数点数値に対応できません。

import pandas as pd
import pyarrow as pa
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()

target = pa.schema([
    pa.field('ID', pa.int32()),
    pa.field('Name', pa.string()),
    pa.field('Amount', pa.decimal128(19, 4)),   # MONEY
])

table = pa.Table.from_pandas(df, preserve_index=False).cast(target)
result = cursor.bulkcopy_arrow("##PandasDemo", table)

型変換がないと、ValueError: Cannot map Arrow column 'Amount' (Float64) to SQL column 'Amount' (Money) で読み込みに失敗します。 スキーマをTable.cast()に渡すのではなく、Table.from_pandas()でキャストを構築してください。Table.cast()では、float 型の列を直接decimal128に変換できません。 NaN この経路では値がSQL NULL になるので、最初に置き換える必要はありません。

代わりに行タプルのパスが必要な場合は、name=None を渡すと itertuples() はすでにタプルを返します:

data = list(df.itertuples(index=False, name=None))
result = cursor.bulkcopy("##PandasDemo", data)

Apache Arrow データのロード

cursor.bulkcopy_arrow()を使ってApache Arrowのデータを読み込みます。 このメソッドはArrowのメモリから直接読み込むため、呼び出す前にPythonの行タプルを作らない。

source引数は、Arrow Cのデータインターフェースを公開するpyarrow.Tablepyarrow.RecordBatchpyarrow.RecordBatchReader、または任意のオブジェクトを受け入れます。 残りの議論も bulkcopy()と同じです。

import mssql_python
import pyarrow as pa

conn = mssql_python.connect(connection_string)

# bulkcopy_arrow() opens its own connection, so commit the table creation first.
conn.autocommit = True
cursor = conn.cursor()

cursor.execute("""
    CREATE TABLE ##ArrowDemo (ID INT, Name NVARCHAR(50), Amount FLOAT)
""")

table = pa.table({
    "ID": pa.array([1, 2, 3], type=pa.int32()),
    "Name": pa.array(["Alice", "Bob", "Carol"], type=pa.string()),
    "Amount": pa.array([50000.0, 60000.0, 55000.0], type=pa.float64()),
})

result = cursor.bulkcopy_arrow("##ArrowDemo", table)
print(f"Copied {result['rows_copied']} rows")

各Arrow列タイプは、宛先のSQL列タイプと互換性がある必要があります。 書き手は型ファミリー間の変換を行わないため、 float64 列を マネー 列に渡すと、行が書かれる前に ValueError が上がります。 金額小数decimal128の列にはを使用します。

Arrow ソースを bulkcopy() に渡すと、TypeError が発生し、bulkcopy_arrow() に誘導されます。

Arrowのサポートについての詳細や、あるテーブルから別のテーブルへ結果セットをストリーミングする方法については、 Apache Arrowの統合を参照してください。

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)

ID 列

明示的な単位元値を挿入するには、 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 各バッチを内部トランザクションで囲みます。

Note

bulkcopy() サーバーへの別の内部接続を開きます。 その内部接続はカーソルのクエリタイムアウトを引き継ぎます。カーソルを作成する前に Connection.timeout を正の値に設定し、同じ値が一括コピー接続の試みを制限します。 カーソルのクエリタイムアウトが 0の場合、内部接続はデフォルトの15秒間の接続タイムアウトを使用します。 カーソルは作成時に値を取りますので、その後 Connection.timeout を変更しても既存のカーソルや進行中の大量コピーには影響しません。 遅い、スロットルがかかる、または高遅延のエンドポイント(例えばVPNや地域間)には、カーソルを作成する前にクエリタイムアウトを上げてください。

エラーを処理する

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()は行タプルの反復可能を取ります。したがって、コピーが始まる前にすべての値はPythonオブジェクトとして存在しなければなりません。 データがすでにカラム状の場合、 bulkcopy_arrow() はArrowバッファを直接読み取り、そのステップをスキップします。 pandasやPolars DataFrame、Parquetファイル、そしてその cursor.arrow() 結果はすべてArrowソースです。 詳細については、「 Apache Arrow データを読み込み」をご覧ください。

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

ジェネレーターは、 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_arrow() すでにカラム形式の大規模なデータセット。 最 速
cursor.bulkcopy() 行指向ソースからの大規模データセット(1,000行以上)。 速い
cursor.executemany() パラメータ付きのメディアデータセット。 Moderate
cursor.execute() ループ内で 小さなデータセットで単純な論理。 最も遅い