Wähle ein Datenlade- und Bewegungsmuster mit mssql-python

Der Treiber mssql-python bietet mehrere Pfade zum Schreiben von Daten in Microsoft SQL. Jeder Weg passt zu unterschiedlichen Arbeitslasten. Dieser Leitfaden hilft Ihnen, basierend auf Ihrem Datenvolumen, Quellformat und Aktualisierungssemantik die richtige auszuwählen.

Entscheide nach Arbeitsbelastung

Arbeitsbelastung Empfohlener Pfad Warum?
CSV-Dateien in eine Tabelle laden CSV-Daten per Massenimport laden bulkcopy() mit einem Generator verarbeitet Dateien beliebiger Größe, ohne sie in den Speicher zu laden.
Füge eine einzelne Zeile aus dem Anwendungscode ein Einfügen einzelner Zeilen Geringer Overhead, einfache Fehlerbehandlung, funktioniert mit OUTPUT zur Rückgabe erzeugter Schlüssel.
Fügen Sie einen kleinen bis mittleren Batch aus dem Anwendungscode ein Stapelweise Einfügungen Reduziert Hin- und Rückfahrten im Vergleich zu einzelnen Einsätzen.
Lade hunderte Zeilen oder mehr von jeder Quelle Mehrfachkopie TDS Bulk Insert ist der effizienteste Weg für große Volumen.
Zeilen basierend auf einem Schlüssel einfügen oder aktualisieren Upsert mit MERGE MERGE behandelt INSERT, UPDATE, und DELETE in einer Anweisung.
Laden Sie einen DataFrame in eine Tabelle DataFrames laden Zeilen aus pandas oder Polars extrahieren und an bulkcopy() übergeben.
Daten in Parquet-Dateien bereitstellen Parkett-Inszenierung Nützlich für systemübergreifende ETL, bei denen ein Zwischendateiformat benötigt wird.

CSV-Daten per Massenimport laden

Das Laden von CSV-Daten ist die häufigste Eingabefrage bei Python-Datenbankarbeiten. csv.reader mit einem Generator verwenden, der bulkcopy() speist:

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

Das Generator-Muster hält den Speicherverbrauch unabhängig von der Dateigröße konstant. Für Spaltenabbildung und Identitätsbehandlung siehe Massenkopieroperationen.

Einfügen einzelner Zeilen

Verwenden Sie einzelne Inserts für Anwendungs-Schreibvorgänge, bei denen Sie jeweils einen Datensatz verarbeiten. Verwendung OUTPUT INSERTED zum Abrufen erzeugter Schlüssel:

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

Einzeleinsätze sind die richtige Wahl, wenn:

  • Du fügst pro Benutzeraktion (Formulareinreichung, API-Aufruf) eine Zeile ein.
  • Du musst jede Zeile einzeln validieren oder transformieren, bevor du sie einfügst.
  • Du brauchst sofort die eingefügte ID oder andere generierte Werte.

Stapel-Einfügungen

Verwenden Sie executemany(), wenn Sie eine mittlere Anzahl von Zeilen haben und nicht den Durchsatz von Massenkopiervorgängen benötigen:

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() jede Zeile als separate parametrisierte Anweisung sendet. Wenn der Durchsatz wichtiger ist als die Kontrolle über einzelne Zeilen, ist bulkcopy() effizienter, weil es das TDS-Bulk-Insert-Protokoll verwendet. Der Schwellenwert hängt von der Zeilenbreite und der Netzwerklatenz ab, liegt aber typischerweise im niedrigen Hunderterbereich.

Massenkopieren

Wenn der Durchsatz wichtiger ist als die Steuerung auf Zeilenebene, verwenden Sie bulkcopy(). Es verwendet das TDS-Bulk-Insert-Protokoll, das deutlich effizienter ist als zeilenweise Einfügungen:

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

Performance-Tipps für Massenkopien

  • Verwenden Sie Generatoren für große Datensätze, um die Speichernutzung konstant zu halten.
  • Set batch_size um zu steuern, wie viele Zeilen pro TDS-Batch gesendet werden. Fang mit 5.000 an und passe je nach Reihenbreite an.
  • Verwenden Sie Tabellensperren für exklusive Ladungen: cursor.bulkcopy("dbo.ProductImport", rows, table_lock=True).
  • Deaktiviere die Indizes vor dem Laden und baue sie danach neu auf. Diese Sequenz vermeidet während des Ladens den Wartungsaufwand für Indizes.

Für Spaltenabbildungen, Identitätsspalten, NULL-Behandlung und paralleles Laden siehe Massenkopieroperationen.

Upsert mit MERGE

MERGE ist die Anweisung in Microsoft SQL für bedingte INSERT, UPDATE und DELETE in einem einzigen Vorgang. Es unterstützt das Muster „einfügen, wenn neu; aktualisieren, wenn vorhanden“, das Python-Entwickler typischerweise benötigen.

Upsert für eine einzelne Zeile

Für eine einzelne Zeile verwenden Sie MERGE mit einer USING Klausel, die Parameter-Aliasse definiert:

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

Aufbauen Sie mit einem Aufbautisch

Für Bulk-Upserts laden Sie die Daten zunächst in eine temporäre Tabelle und verwenden anschließend MERGE, um daraus zu aktualisieren. Verwenden Sie insert-or-update als Standardmuster für DataFrame-Upserts und Batch-Updates:

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

Dieses Beispiel zeigt das Standard-Einfüge- oder Update-Muster:

  • INSERT Zeilen aus der Quelle, die im Ziel nicht existieren (WHEN NOT MATCHED BY TARGET).
  • UPDATE Zeilen, die in beiden vorhanden sind (WHEN MATCHED).
  • Die OUTPUT-Klausel berichtet, welche Aktion in jeder Zeile durchgeführt wurde, was für Audit-Trails nützlich ist.

Caution

Fügen Sie WHEN NOT MATCHED BY SOURCE THEN DELETE nur hinzu, wenn die Staging-Daten eine maßgebliche vollständige Momentaufnahme des Ziels sind. Wenn der Batch nur geänderte Zeilen enthält, löscht diese Klausel Zeilen, die absichtlich aus dem Quellfeed weggelassen wurden.

Wenn Sie eine vollständige Abstimmung benötigen, erweitern Sie MERGE erst, nachdem Sie bestätigt haben, dass die Quelle für die Zieltabelle maßgeblich ist:

WHEN NOT MATCHED BY SOURCE THEN
    DELETE

In gemeinsamen Umgebungen verwenden Sie pro Ausführung einen eindeutigen globalen temporären Tabellennamen oder eine permanente Staging-Tabelle, um Kollisionen zwischen gleichzeitigen Jobs zu vermeiden.

Wann stattdessen separate UPDATE- und INSERT-Anweisungen verwendet werden sollten

MERGE ist mächtig, hat aber Randfälle. Erwägen Sie, separate Aussagen zu verwenden, wenn:

  • Du brauchst keine DELETE Logik. Ein separater UPDATE gefolgt von INSERT WHERE NOT EXISTS ist besser lesbar und leicht zu debuggen.
  • Die Aussage MERGE ist komplex genug, dass das Sperrverhalten schwer vorherzusagen ist. Getrennte Anweisungen ermöglichen eine explizite Kontrolle über die Sperrgranularität.
  • Du aktualisierst eine Tabelle mit hoher Parallelität, bei der MERGE eine Sperreskalation zu Blockierungen führen kann.
# 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()

DataFrames laden

Extrahieren Sie Zeilen aus einem Pandas- oder Polars-DataFrame und laden Sie sie mit bulkcopy():

pandas

Einen pandas DataFrame in Tupel umwandeln und an bulkcopy() übergeben:

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

Konvertieren eines Polars-DataFrame in Tupel mithilfe der Methode .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()

Für vollständige DataFrame-Lademuster siehe pandas-Integration und Polars-integration.

Parkett-Inszenierung

Verwenden Sie Parquet als Zwischenformat, wenn Sie Daten zwischen Systemen migrieren oder wenn Ihre ETL-Pipeline bereits Parquet-Dateien produziert:

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

Für große Parquet-Dateien lesen Sie in Zeilengruppen, um den Speicherverbrauch konstant zu halten:

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

Validiere geladene Daten

Überprüfen Sie nach dem Laden die Anzahl der Zeilen und prüfen Sie die Daten stichprobenartig:

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}")

Bei Produktionslast solltest du dich nicht darauf verlassen, dass die Transaktion der aufrufenden Verbindung einen bulkcopy()-Aufruf absichert. bulkcopy() öffnet eine eigene interne Verbindung und schreibt die kopierten Zeilen unabhängig fest, sodass ein conn.rollback() auf Ihrer Hauptverbindung sie nicht rückgängig machen kann. Zwei Ansätze ergeben Atomizität:

  • Setzen Sie use_internal_transaction=True so, dass jeder Batch in einer eigenen Transaktion ausgeführt wird. Eine Charge, die teilweise fehlschlägt, rollt diese Charge zurück, anstatt sie halb geladen zu lassen.
  • Um die Daten vor dem Promoten zu validieren, kopiere in Massen in eine Staging-Tabelle, validiere sie und verschiebe dann die Zeilen in die Zieltabelle, indem du eine INSERT ... SELECT Transaktion innerhalb einer Transaktion auf deiner Hauptverbindung verwendest. Weil das INSERT auf deiner Verbindung läuft, macht conn.rollback() es rückgängig, wenn die Validierung fehlschlägt.
# 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