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

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

仕事量で決めてください

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

一括コピーを使用して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行の挿入

アプリケーションレベルの書き込みにはシングルインサートを使い、一度に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のバルクインサートプロトコルを使用しており、これは行ごとの挿入よりもはるかに効率的です。

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()

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

  • 大規模なデータセットにはメモリ使用量を一定に保つためにジェネレーターを使いましょう。
  • 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()

データフレームのロード

pandasやPolars DataFrameから行を抽出し、以下の bulkcopy()で読み込みます:

pandas

pandasのデータフレームをタプルに変換し、 bulkcopy()に渡します:

import pandas as pd

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

# Convert DataFrame rows to tuples
rows = list(df[["Name", "ProductNumber", "ListPrice"]].itertuples(index=False, name=None))

cursor.bulkcopy("dbo.ProductImport", rows, batch_size=5000)
conn.commit()

Polars

Polarsデータフレームを .rows() メソッドを用いてタプルに変換します:

import polars as pl

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

# Convert Polars DataFrame to list of tuples
rows = df.select(["Name", "ProductNumber", "ListPrice"]).rows()

cursor.bulkcopy("dbo.ProductImport", rows, batch_size=5000)
conn.commit()

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

Parquet ステージング

システム間のデータ移行やETLパイプラインですでにParquetファイルを生成している場合に、中間フォーマットとしてParquetを使用してください:

import pyarrow.parquet as pq

# Read Parquet file
table = pq.read_table("products.parquet")

# Convert to rows for bulkcopy
rows = [tuple(row) for row in zip(*[col.to_pylist() for col in table.columns])]

cursor.bulkcopy("dbo.ProductImport", rows, batch_size=5000)
conn.commit()

大きなParquetファイルの場合、メモリ使用量を一定に保つために行グループを読み込みます:

import pyarrow.parquet as pq

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

for batch in parquet_file.iter_batches(batch_size=10000):
    rows = [tuple(row) for row in zip(*[col.to_pylist() for col in batch.columns])]
    cursor.bulkcopy("dbo.ProductImport", rows, batch_size=10000)

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