mssql-pythonでデータロードと移動パターンを選択します

mssql-pythonドライバーは、Microsoft SQLにデータを書き込むための複数の方法を提供しています。 それぞれの道は異なる課題に対応しています。 このガイドは、あなたのデータ量、ソースフォーマット、最新の意味論に基づいて最適なものを選ぶのに役立ちます。

仕事量で決めてください

Workload 推奨パス なぜでしょうか
CSVファイルをテーブルに読み込む CSVデータをまとめて読み込む bulkcopy() ジェネレーターを使うことで、メモリにロードせずに任意のサイズのファイルを処理できます。
アプリケーションコードから1行を挿入します 単一行の挿入 低オーバーヘッドで単純なエラー処理が可能で、生成キーの返却にはOUTPUTと連携します。
アプリケーションコードから小規模から中規模のバッチを挿入します 一括挿入 単一挿入と比べてラウンドトリップ数を削減します。
どのソースからでも何百行以上も読み込んでください 一括コピー TDSバルクインサートは大量輸送において最も効率的な経路です。
キーに基づいて行を挿入または更新する MERGE を使用したアップサート MERGE 1つの文で INSERT、 UPDATE、 DELETE を処理します。
データフレームをテーブルに読み込む データフレームのロード bulkcopy_arrow()すべての値に対してPythonオブジェクトを構築せずにDataFrameの矢印データを読み取ります。
Apache Arrowデータをテーブルに読み込む ロードアローデータ bulkcopy_arrow()Pythonタプルを構築せずにArrowメモリを直接読み込みます。
Parquet ファイルを使用してデータをステージングする Parquet ステージング 中間ファイル形式が必要なクロスシステムETLに有用です。 ParquetはすでにArrowのデータなので、行変換なしで読み込みます。

一括コピーを使用してCSVデータを読み込む

CSVデータの読み込みはPythonデータベース作業で最も一般的な取り込み質問です。 発電機の給電csv.readerbulkcopy()を組み合わせて使用:

import csv
import mssql_python

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

# Create a target table
cursor.execute("""
    IF NOT EXISTS (SELECT * FROM sys.tables WHERE name = 'ProductImport')
    CREATE TABLE dbo.ProductImport (
        Name nvarchar(100),
        ProductNumber nvarchar(25),
        ListPrice decimal(10,2)
    )
""")
conn.commit()

def csv_rows(path):
    with open(path, newline="", encoding="utf-8") as f:
        reader = csv.reader(f)
        next(reader)  # Skip header
        for row in reader:
            yield (row[0], row[1], float(row[2]))

result = cursor.bulkcopy(
    "dbo.ProductImport",
    csv_rows("products.csv"),
    batch_size=5000
)
print(f"Loaded {result['rows_copied']} rows")
conn.commit()

ジェネレーターパターンはファイルサイズに関係なくメモリ使用量を一定に保ちます。 列の写像と同一性の処理については、 バルクコピー操作を参照してください。

単一行の挿入

アプリケーションレベルの書き込みにはシングルインサートを使い、一度に1つのレコードを処理します。 生成された鍵を取得するには OUTPUT INSERTED を使う:

cursor.execute("""
    INSERT INTO dbo.ProductImport (Name, ProductNumber, ListPrice)
    OUTPUT INSERTED.Name
    VALUES (%(name)s, %(product_number)s, %(list_price)s)
""", {"name": "Widget", "product_number": "WG-1000", "list_price": 19.99})

inserted_name = cursor.fetchval()
conn.commit()

シングルインサートは以下の場合に適切な選択です:

  • ユーザーアクション(フォーム提出、API呼び出し)ごとに1行ずつ挿入します。
  • 挿入前に各行を個別に検証または変換する必要があります。
  • 挿入されたIDや他の生成された値がすぐに必要です。

バッチ挿入

行数が適度で、まとめコピーのスループットが不要なときに executemany() を使います:

rows = [
    {"name": "Widget A", "product_number": "WG-1001", "list_price": 19.99},
    {"name": "Widget B", "product_number": "WG-1002", "list_price": 24.99},
    {"name": "Widget C", "product_number": "WG-1003", "list_price": 29.99},
]

cursor.executemany(
    "INSERT INTO dbo.ProductImport (Name, ProductNumber, ListPrice) VALUES (%(name)s, %(product_number)s, %(list_price)s)",
    rows
)
conn.commit()

executemany() 各行を別々のパラメータ化された文として送信します。 スループットが行単位制御よりも重要な場合、 bulkcopy() はTDSのバルクインサートプロトコルを使用しているためより効率的です。 クロスオーバーは行幅やネットワークのレイテンシに依存しますが、通常は数百行前半です。

一括コピー

スループットが行ごとの制御よりも重要な場合は bulkcopy()を使いましょう。 TDSのバルクインサートプロトコルを使用しており、1行ごとに1つの文を送信するのではなく、行をストリーミングします。

rows = [
    ("Widget A", "WG-1001", 19.99),
    ("Widget B", "WG-1002", 24.99),
    ("Widget C", "WG-1003", 29.99),
]

result = cursor.bulkcopy("dbo.ProductImport", rows, batch_size=5000)
print(f"Loaded {result['rows_copied']} rows")
conn.commit()

一括コピーのパフォーマンスに関するヒント

  • 大規模なデータセットにはメモリ使用量を一定に保つためにジェネレーターを使いましょう。
  • bulkcopy_arrow() DataFrameやParquetファイルのように、ソースがカラム状の場合に利用してください。 Pythonの行タプルへの変換をスキップします。
  • batch_size を設定して、TDS バッチあたりに送信する行数を制御します。 5,000から始めて、行幅に応じて調整してください。
  • 排他的な荷物にはテーブルロックを使いましょう:cursor.bulkcopy("dbo.ProductImport", rows, table_lock=True)
  • 読み込む前にインデックスを無効にし、その後再構築してください。 この手順により、負荷中のインデックスメンテナンスのオーバーヘッドを回避できます。

列のマッピング、識別列、NULL処理、並列読み込みについては、 バルクコピー操作を参照してください。

MERGE を使用したアップサート

MERGEは、1回の操作で条件付きのINSERT、UPDATE、およびDELETEを実行するためのMicrosoft SQLのステートメントです。 これはPython開発者がよく必要とする「新しい場合は挿入、存在すれば更新」というパターンを処理します。

単一行アップサート

単一行の場合、パラメータエイリアスを定義する MERGE 節付きのUSINGを用います。

cursor.execute("""
    MERGE dbo.ProductImport AS target
    USING (SELECT %(name)s AS Name, %(product_number)s AS ProductNumber, %(list_price)s AS ListPrice) AS source
    ON target.ProductNumber = source.ProductNumber
    WHEN MATCHED THEN
        UPDATE SET
            Name = source.Name,
            ListPrice = source.ListPrice
    WHEN NOT MATCHED THEN
        INSERT (Name, ProductNumber, ListPrice)
        VALUES (source.Name, source.ProductNumber, source.ListPrice);
""", {"name": "Widget A", "product_number": "WG-1001", "list_price": 24.99})
conn.commit()

ステージング テーブル付きの一括アップサート

バルクアップサートの場合は、まずデータを一時テーブルにステージ化し、そこから更新するために MERGE を使いましょう。 DataFrameアップサートおよびバッチ更新のデフォルトパターンとしてinsert-or-updateを使用してください:

import csv
import mssql_python

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

# Step 1: Create a global temp table for staging
# Note: bulkcopy() requires global temp tables (##), not session temp tables (#)
cursor.execute("""
    IF OBJECT_ID('tempdb..##ProductImportStage') IS NOT NULL
        DROP TABLE ##ProductImportStage;
    CREATE TABLE ##ProductImportStage (
        Name nvarchar(100),
        ProductNumber nvarchar(25),
        ListPrice decimal(10,2)
    )
""")
cursor.commit()

# Step 2: Bulk load into the staging table
def csv_rows(path):
    with open(path, newline="", encoding="utf-8") as f:
        reader = csv.reader(f)
        next(reader)
        for row in reader:
            yield (row[0], row[1], float(row[2]))

cursor.bulkcopy("##ProductImportStage", csv_rows("products_update.csv"), batch_size=5000)

# Step 3: MERGE from staging into the target table
cursor.execute("""
    MERGE dbo.ProductImport AS target
    USING ##ProductImportStage AS source
    ON target.ProductNumber = source.ProductNumber
    WHEN MATCHED THEN
        UPDATE SET
            Name = source.Name,
            ListPrice = source.ListPrice
    WHEN NOT MATCHED BY TARGET THEN
        INSERT (Name, ProductNumber, ListPrice)
        VALUES (source.Name, source.ProductNumber, source.ListPrice)
    OUTPUT $action, INSERTED.ProductNumber, DELETED.ProductNumber;
""")

# Step 4: Read the OUTPUT to see what changed
for row in cursor.fetchall():
    print(f"{row[0]}: inserted={row[1]}, deleted={row[2]}")

conn.commit()

この例はデフォルトの挿入または更新パターンを示しています:

  • INSERT ターゲット(WHEN NOT MATCHED BY TARGET)には存在しないソースからの行。
  • UPDATE 両方 (WHEN MATCHED) に存在する行。
  • OUTPUT 節は各行で行われたアクションを報告し、監査トレイルに役立ちます。

Caution

ステージングデータがターゲットの権威ある完全なスナップショットである場合にのみ、 WHEN NOT MATCHED BY SOURCE THEN DELETE を追加してください。 バッチに変更された行のみが含まれている場合、その節は意図的にソースフィードから省略された行を削除します。

完全な照合が必要な場合は、ターゲットテーブルのソースが権威あるものであることを確認した後にのみ MERGE を延長してください。

WHEN NOT MATCHED BY SOURCE THEN
    DELETE

共有環境では、実行ごとに固有のグローバル一時テーブル名や恒久的なステージングテーブルを使用して、同時ジョブ間の衝突を防ぎます。

代わりに別々の UPDATE 文と INSERT 文を使うべき時

MERGE 強力ですが、例外もあります。 以下の場合に別々の文を使うことを考えてみてください:

  • DELETEロジックは必要ありません。 別 UPDATE に続く INSERT WHERE NOT EXISTS の方が読みやすく、デバッグも簡単です。
  • MERGE文は非常に複雑で、ロック挙動の予測が難しいです。 別々のステートメントはロックの細かさを明確にコントロールできます。
  • 高並行処理テーブルを更新しており、MERGE ロックのエスカレーション ブロックを引き起こす可能性があります。
# Simpler alternative: UPDATE then INSERT
cursor.execute("""
    UPDATE dbo.ProductImport
    SET Name = %(name)s, ListPrice = %(list_price)s
    WHERE ProductNumber = %(product_number)s
""", {"name": "Widget A", "list_price": 24.99, "product_number": "WG-1001"})

if cursor.rowcount == 0:
    cursor.execute("""
        INSERT INTO dbo.ProductImport (Name, ProductNumber, ListPrice)
        VALUES (%(name)s, %(product_number)s, %(list_price)s)
    """, {"name": "Widget A", "product_number": "WG-1001", "list_price": 24.99})

conn.commit()

データフレームのロード

DataFrame は列指向なので、bulkcopy() 用に行タプルへ平坦化するのではなく、bulkcopy_arrow() で読み込みます。

読み込む前に、Arrow 型を宛先列に対応付けてください。 pyarrow 数値列の float64 を推論しますが、運転手は それをお金小数点数値にマッピングできません。

pandas

import pandas as pd
import pyarrow as pa

df = pd.read_csv("products.csv")

target = pa.schema([
    pa.field("Name", pa.string()),
    pa.field("ProductNumber", pa.string()),
    pa.field("ListPrice", pa.decimal128(19, 4)),   # MONEY
])

table = pa.Table.from_pandas(
    df[["Name", "ProductNumber", "ListPrice"]], preserve_index=False
).cast(target)

cursor.bulkcopy_arrow("dbo.ProductImport", table)
conn.commit()

スキーマをTable.from_pandas()に渡す代わりにTable.cast()を使うべきです。はフロートカラムを直接の decimal128 に変換できません。

Polars

PolarsはArrow Cのデータインターフェースを実装しているので、DataFrame自体をパスできます。 同じ理由でまず柱を鋳造してください:

import polars as pl

df = pl.read_csv("products.csv")

cursor.bulkcopy_arrow(
    "dbo.ProductImport",
    df.select([
        "Name",
        "ProductNumber",
        pl.col("ListPrice").cast(pl.Decimal(19, 4)),   # MONEY
    ]),
)
conn.commit()

ファイルを読み取る際に pl.read_csv("products.csv", schema_overrides={"ListPrice": pl.Decimal(19, 4)})で型を設定することもできます。

DataFrameを渡すことで、そのバッファはコピーなしで直接ドライバーに渡されます。 df.to_arrow() こちらも動作しますが、Polarsはその変換中に文字列の列を再エンコードし、すべての文字列データをコピーします。

bulkcopy_arrow() pyarrow.TableRecordBatchRecordBatchReader、または__arrow_c_stream__または__arrow_c_array__を通じてArrow Cのデータインターフェースを実装する任意のオブジェクトを受け入れます。 これらのいずれかを bulkcopy() に渡すと、TypeError が発生します。

完全なデータフレームの読み込みパターンについては、 pandas積分 および Polars積分を参照してください。

Arrowデータを読み込む

ソースがすでにApache Arrow形式であれば、cursor.bulkcopy_arrow()Pythonタプルを先に構築せずにロードします。

from decimal import Decimal

import pyarrow as pa

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

table = pa.table({
    "Name": pa.array(["Widget", "Gadget"], type=pa.string()),
    "ProductNumber": pa.array(["WI-1000", "GA-2000"], type=pa.string()),
    "ListPrice": pa.array([Decimal("29.99"), Decimal("49.99")], type=pa.decimal128(10, 2)),
})

result = cursor.bulkcopy_arrow("dbo.ProductImport", table, batch_size=5000)
print(f"Copied {result['rows_copied']} rows")

このメソッドは pyarrow.RecordBatchpyarrow.RecordBatchReaderも受け入れているので、 cursor.arrow_reader() の結果セットを別のテーブルに直接ストリーミングできます。

各Arrow列のタイプは宛先のSQLカラムタイプと互換性があり、書き手はタイプファミリー間で変換しません。 詳細については、 Apache Arrowの統合を参照してください。

Parquet ステージング

システム間のデータ移行やETLパイプラインがすでにParquetファイルを生成している場合に、中間フォーマットとしてParquetを使用してください。 ParquetファイルはArrowのデータを読み取るので、直接 bulkcopy_arrow()に渡してください:

import pyarrow.parquet as pq

cursor.bulkcopy_arrow("dbo.ProductImport", pq.read_table("products.parquet"))
conn.commit()

大きなParquetファイルでは、メモリ使用量を一定に保つために行グループを反復します。 各バッチは RecordBatchであり、 bulkcopy_arrow() は直接以下を受け付けます:

import pyarrow.parquet as pq

parquet_file = pq.ParquetFile("products.parquet")

for batch in parquet_file.iter_batches(batch_size=10000):
    cursor.bulkcopy_arrow("dbo.ProductImport", batch)

conn.commit()

ファイルを1回の呼び出しでストリーミングするには、バッチを RecordBatchReaderでラップします:

import pyarrow as pa
import pyarrow.parquet as pq

parquet_file = pq.ParquetFile("products.parquet")
reader = pa.RecordBatchReader.from_batches(
    parquet_file.schema_arrow, parquet_file.iter_batches(batch_size=10000)
)

cursor.bulkcopy_arrow("dbo.ProductImport", reader)
conn.commit()

読み込まれたデータの検証

読み込み後、行数とスポットチェックデータを確認してください:

cursor.execute("SELECT COUNT(*) FROM dbo.ProductImport")
count = cursor.fetchval()
print(f"Total rows: {count}")

cursor.execute("""
    SELECT TOP 5 Name, ProductNumber, ListPrice
    FROM dbo.ProductImport
    ORDER BY Name
""")
for row in cursor:
    print(f"  {row.Name} ({row.ProductNumber}): ${row.ListPrice:.2f}")

本番環境のロードでは、 bulkcopy() 通話を保護するために呼び出し接続のトランザクションに頼らないでください。 bulkcopy() 独自の内部接続を開き、コピーした行を独立してコミットするため、メイン接続の conn.rollback() では元に戻せません。 原子性を得る方法は二つあります。

  • use_internal_transaction=True を、各バッチをそれぞれ独立したトランザクションで処理するように設定してください。 途中で失敗したバッチは、中途半端に読み込まれた状態のままにするのではなく、そのバッチ全体がロールバックされます。
  • プロモート前にデータを検証するには、ステージングテーブルに一括コピーし、検証後、メイン接続のトランザクション内の INSERT ... SELECT を使って行をターゲットテーブルに移動させます。 そのINSERTは接続上で実行されるため、検証が失敗した場合はconn.rollback()がそれを取り消します。
# Stage the data. bulkcopy() runs on its own connection, so these rows
# persist regardless of the transaction below.
cursor.bulkcopy("dbo.ProductImport_Stage", rows, batch_size=5000)

try:
    cursor.execute("SELECT COUNT(*) FROM dbo.ProductImport_Stage")
    count = cursor.fetchval()

    if count < expected_count:
        raise ValueError(f"Expected {expected_count} rows, got {count}")

    # This INSERT runs on your connection, so it's covered by the transaction.
    cursor.execute("""
        INSERT INTO dbo.ProductImport (Name, ProductNumber, ListPrice)
        SELECT Name, ProductNumber, ListPrice FROM dbo.ProductImport_Stage
    """)
    conn.commit()
except Exception:
    conn.rollback()
    raise