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 Batch-Einfügungen Reduziert die Anzahl der Roundtrips im Vergleich zu einzelnen Einfügungen.
Lade hunderte Zeilen oder mehr von jeder Quelle Mehrfachkopie TDS-Bulk-Einfügungen sind der effizienteste Weg für große Volumes.
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 bulkcopy_arrow()liest die Arrow-Daten des DataFrame, ohne für jeden Wert ein Python-Objekt zu bauen.
Apache-Pfeil-Daten in eine Tabelle laden Pfeildaten laden bulkcopy_arrow() liest direkt aus dem Arrow-Speicher, ohne Python-Tupel zu erstellen.
Bereiten Sie Daten über Parquet-Dateien vor Parquet-Staging Nützlich für systemübergreifende ETL, bei denen ein Zwischendateiformat benötigt wird. Parquet sind bereits Arrow-Daten und werden daher ohne Zeilenumwandlung geladen.

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.

Batch-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. Die Schwelle hängt von der Zeilenbreite und der Netzwerklatenz ab, liegt jedoch typischerweise im niedrigen dreistelligen Zeilenbereich.

Massenkopieren

Wenn der Durchsatz wichtiger ist als die Steuerung auf Zeilenebene, verwenden Sie bulkcopy(). Es verwendet das TDS-Bulk-Insert-Protokoll, das Zeilen streamt, anstatt pro Zeile eine Anweisung zu senden:

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.
  • Verwenden Sie bulkcopy_arrow(), wenn die Quelle spaltenbasiert ist, z. B. ein DataFrame oder eine Parquet-Datei. Es überspringt die Umwandlung in Python-Zeilentupel.
  • 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 Ladevorgänge: 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.

Einzeiliges „Upsert“

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

Bulk-„Upsert“ mit einer Stagingtabelle

Bei Bulk-„Upserts“ sollten Sie die Daten zunächst in eine temporäre Tabelle zwischenspeichern und anschließend MERGE verwenden, um die Aktualisierung von dort aus durchzuführen. 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 erst hinzu, wenn die zwischengespeicherten Daten einen verbindlichen vollständigen Snapshot der Zieltabelle darstellen. 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 Abgleichung benötigen, erweitern Sie den MERGE erst, nachdem Sie sich vergewissert 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. Separate Anweisungen ermöglichen Ihnen eine explizite Kontrolle über die Sperrgranularität.
  • Sie aktualisieren eine Tabelle mit hoher Parallelität, bei der eine MERGE Lock-Eskalation zu Blockierungen führen könnte.
# 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

Ein DataFrame ist spaltenbasiert, also lade es mit bulkcopy_arrow(), anstatt es für bulkcopy() in Zeilen-Tupel umzuwandeln.

Ordne die Pfeiltypen den Zielspalten zu, bevor du lädst. pyarrow leitet auf float64 eine numerische Spalte, die der Fahrer nicht auf Geld,Dezimal- oder Zahlenzahl abbilden kann.

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

Verwenden Sie Table.cast(), anstatt das Schema an Table.from_pandas() zu übergeben, da Table.from_pandas() eine Float-Spalte nicht direkt in decimal128 konvertieren kann.

Polars

Polars implementiert die Arrow C Datenschnittstelle, sodass du den DataFrame selbst weitergeben kannst. Gieße die Säulen zuerst aus demselben Grund:

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

Du kannst auch die Typen beim Lesen der Datei festlegen, mit pl.read_csv("products.csv", schema_overrides={"ListPrice": pl.Decimal(19, 4)}).

Die direkte Übergabe des DataFrame an den Treiber übergibt dessen Puffer, ohne eine Kopie zu erstellen. df.to_arrow() funktioniert auch, aber Polars kodiert während dieser Umwandlung die String-Spalten neu, wodurch alle Stringdaten kopiert werden.

bulkcopy_arrow() akzeptiert ein pyarrow.Table, , RecordBatch, ein RecordBatchReader, oder jedes Objekt, das die Datenschnittstelle von Arrow C über __arrow_c_stream__ oder __arrow_c_array__implementiert. Wenn eines davon an bulkcopy() übergeben wird, wird TypeError ausgelöst.

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

Pfeildaten laden

Wenn die Quelldaten bereits im Apache-Arrow-Format vorliegen, lädt cursor.bulkcopy_arrow() sie, ohne zuvor Python-Tupel zu erstellen.

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

Die Methode akzeptiert außerdem ein pyarrow.RecordBatch oder ein pyarrow.RecordBatchReader, sodass Sie eine Ergebnismenge direkt aus cursor.arrow_reader() in eine andere Tabelle streamen können.

Jeder Arrow-Spaltentyp muss mit seinem Ziel-SQL-Spaltentyp kompatibel sein, und der Schreiber konvertiert nicht zwischen den Typfamilien. Weitere Informationen finden Sie unter Apache Arrow Integration.

Parquet-Staging

Verwenden Sie Parquet als Zwischenformat, wenn Sie Daten zwischen Systemen migrieren oder wenn Ihre ETL-Pipeline bereits Parquet-Dateien erzeugt. Eine Parquet-Datei liest in die Arrow-Daten, also übergebe sie direkt an bulkcopy_arrow():

import pyarrow.parquet as pq

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

Für große Parquet-Dateien iterieren Sie Zeilengruppen, um den Speicherverbrauch konstant zu halten. Jeder Batch ist ein RecordBatch, der bulkcopy_arrow() direkt akzeptiert:

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

Um die gesamte Datei in einem einzigen Aufruf zu streamen, packen Sie die Chargen in einem 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()

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

Verlassen Sie sich bei Produktivladungen nicht auf die Transaktion der aufrufenden Verbindung, um einen bulkcopy()-Aufruf zu schützen. 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:

  • Legen Sie use_internal_transaction=True fest, um jeden Batch in eine eigene Transaktion zu verpacken. Eine Charge, die teilweise fehlschlägt, rollt diese Charge zurück, anstatt sie halb geladen zu lassen.
  • Um Daten vor der Hochstufung zu validieren, kopieren Sie sie im bulk-Verfahren in eine Stagingtabelle, validieren Sie sie und verschieben Sie die Zeilen anschließend mithilfe eines INSERT ... SELECT innerhalb einer Transaktion über Ihre Haupt-Verbindung in die Zieltabelle. 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