Kies een datalaad- en bewegingspatroon met mssql-python

De mssql-python driver biedt meerdere paden om data in Microsoft SQL te schrijven. Elk pad past bij verschillende werklasten. Deze gids helpt je de juiste te kiezen op basis van je datavolume, bronformaat en update-semantiek.

Bepaal op basis van werklast

Werklast Aanbevolen pad Waarom
Laad CSV-bestanden in een tabel CSV-gegevens laden met bulksgewijs kopiëren bulkcopy() met een generator verwerkt bestanden van elke grootte zonder ze in het geheugen te laden.
Voeg een enkele rij in uit applicatiecode Invoegingen van één rij Lage overhead, eenvoudige foutafhandeling, werkt met OUTPUT voor het teruggeven van gegenereerde sleutels.
Voeg een kleine tot matige batch toe uit applicatiecode Batchgewijze invoegingen Vermindert het aantal heen-en-weerbewegingen vergeleken met afzonderlijke inserts.
Laad honderden rijen of meer van elke bron bulksgewijs kopiëren TDS bulk-insert is de meest efficiënte methode voor grote volumes.
Rijen invoegen of bijwerken op basis van een sleutel Upsert met MERGE MERGE verwerkt INSERT, UPDATE, en DELETE in één stelling.
Laad een DataFrame in een tabel DataFrames laden bulkcopy_arrow()leest de Arrow-data van de DataFrame zonder voor elke waarde een Python-object te bouwen.
Laad Apache Arrow-gegevens in een tabel Laad pijlgegevens bulkcopy_arrow()leest direct het geheugen van Arrow, zonder Python-tuples te bouwen.
Gegevens voorbereiden via Parquet-bestanden Parquet-staging Nuttig voor cross-system ETL waarbij een tussentijds bestandsformaat nodig is. Parquet is al Arrow-gegevens, dus het wordt geladen zonder conversie van rijen.

Laad CSV-gegevens door deze bulksgewijs te kopiëren

Het laden van CSV-gegevens is de meest voorkomende ingest-vraag bij Python-databasewerk. Gebruik csv.reader met een generator die bulkcopy() voedt:

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

Het generatorpatroon houdt het geheugengebruik constant ongeacht de bestandsgrootte. Voor kolommapping en identiteitsbehandeling, zie Bulk copy operations.

Invoegingen van één rij

Gebruik single inserts voor applicatieniveau schrijfopdrachten waarbij je één record tegelijk verwerkt. Gebruik OUTPUT INSERTED om gegenereerde sleutels op te halen:

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

Afzonderlijke inzetstukken zijn de beste keuze als:

  • Je voegt per gebruikersactie één rij in (formulierindiening, API-aanroep).
  • Je moet elke rij afzonderlijk valideren of transformeren voordat je het invoegt.
  • Je hebt de ingevoegde ID of andere gegenereerde waarden direct nodig.

Invoegingen in batches

Gebruik executemany() wanneer je een matig aantal rijen hebt en de doorvoersnelheid van bulkkopiëren niet nodig hebt:

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() stuurt elke rij als een aparte geparametriseerde instructie. Wanneer doorvoersnelheid belangrijker is dan controle per rij, bulkcopy() is het efficiënter omdat het het TDS bulk insert-protocol gebruikt. Het omslagpunt hangt af van de rijbreedte en netwerklatentie, maar ligt doorgaans bij enkele honderden rijen.

Bulksgewijs kopiëren

Wanneer doorvoer belangrijker is dan controle per rij, gebruik bulkcopy(). Het gebruikt het TDS bulk insert-protocol, dat rijen streamt in plaats van één instructie per rij te sturen:

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

Tips voor betere prestaties bij bulkkopiëren

  • Gebruik generatoren voor grote datasets om het geheugengebruik constant te houden.
  • Gebruik bulkcopy_arrow() wanneer de bron kolomvormig is, zoals een DataFrame of een Parquet-bestand. Het slaat de conversie naar Python row tuples over.
  • Set batch_size om te bepalen hoeveel rijen per TDS-batch worden verzonden. Begin met 5.000 en pas aan op basis van de rijbreedte.
  • Gebruik tabelvergrendelingen voor exclusieve ladingen: cursor.bulkcopy("dbo.ProductImport", rows, table_lock=True).
  • Schakel indexen uit voor het laden, en bouw daarna opnieuw op. Deze reeks voorkomt de indexonderhoudsoverhead tijdens de belasting.

Voor kolommappingen, identiteitskolommen, NULL-afhandeling en parallel laden, zie Bulk copy operations.

Upsert met MERGE

MERGEis de Microsoft SQL-instructie voor conditionele INSERT, UPDATE, en DELETE in één enkele operatie. Het behandelt het "insert if new, update if exists"-patroon dat Python-ontwikkelaars vaak nodig hebben.

Enkelrijige upsert

Voor één rij gebruik je MERGE met een USING-clausule die aliassen voor parameters definieert:

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 met een opsteltafel

Voor bulk-upserts plaats je de gegevens eerst in een tijdelijke tabel en gebruik je vervolgens MERGE om die gegevens van daaruit bij te werken. Gebruik insert-or-update als standaardpatroon voor DataFrame-upserts en batchupdates:

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

Dit voorbeeld toont het standaard invoeg- of updatepatroon:

  • INSERT rijen uit de bron die niet in het doel bestaan (WHEN NOT MATCHED BY TARGET).
  • UPDATE rijen die in beide voorkomen (WHEN MATCHED).
  • OUTPUT-clausule geeft aan welke actie er op elke rij is ondernomen, wat nuttig is voor auditsporen.

Waarschuwing

Voeg alleen toe WHEN NOT MATCHED BY SOURCE THEN DELETE wanneer de stagingdata een gezaghebbende volledige snapshot van het doel is. Als de batch alleen gewijzigde rijen bevat, verwijdert die clausule rijen die opzettelijk uit de bronfeed zijn weggelaten.

Als je volledige reconciliatie nodig hebt, breid MERGE dan pas uit nadat je hebt bevestigd dat de bron de gezaghebbende bron is voor de doeltabel:

WHEN NOT MATCHED BY SOURCE THEN
    DELETE

In gedeelde omgevingen gebruik per run een unieke globale tijdelijke tabelnaam of een permanente stagingtabel om botsingen tussen gelijktijdige taken te voorkomen.

Wanneer gebruik je in plaats daarvan afzonderlijke UPDATE- en INSERT-instructies

MERGE is krachtig, maar kent uitzonderingssituaties. Overweeg het gebruik van afzonderlijke statements wanneer:

  • Je hebt geen logica nodig DELETE . Een apart UPDATE gevolg door INSERT WHERE NOT EXISTS is leesbaarder en eenvoudiger te debuggen.
  • De MERGE uitspraak is complex genoeg dat het vergrendelingsgedrag moeilijk te voorspellen is. Aparte statements geven je expliciete controle over de vergrendelingsgranulariteit.
  • Je bent een tabel met hoge gelijktijdigheid aan het bijwerken waarbij MERGE lock-escalatie blokkades kan veroorzaken.
# 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 inladen

Een DataFrame is kolomvormig, dus laad het met bulkcopy_arrow() in plaats van het af te vlakken in rij-tuples voor bulkcopy().

Stem de Arrow-typen af op de bestemmingskolommen voordat je de gegevens laadt. pyarrowAfleidt float64 voor een numerieke kolom, die de bestuurder niet kan toewijzen aan geld,decimaal of numeriek.

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

Gebruik Table.cast() in plaats van het schema aan Table.from_pandas()te geven, wat een float-kolom niet direct kan omzetten naar decimal128 .

Polars

Polars implementeert de Arrow C data-interface, zodat je het DataFrame zelf kunt doorgeven. Stort eerst de kolommen om dezelfde reden:

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

Je kunt ook de types instellen wanneer je het bestand leest, met pl.read_csv("products.csv", schema_overrides={"ListPrice": pl.Decimal(19, 4)}).

Door de DataFrame rechtstreeks door te geven, worden de buffers zonder te kopiëren aan de driver overgedragen. df.to_arrow() werkt ook, maar Polars codeert tijdens die conversie de stringkolommen opnieuw, waardoor alle stringgegevens worden gekopieerd.

bulkcopy_arrow() accepteert een pyarrow.Table, een RecordBatch, een RecordBatchReader of elk object dat de Arrow C-data-interface implementeert via __arrow_c_stream__ of __arrow_c_array__. Als een van deze aan bulkcopy() wordt doorgegeven, wordt TypeError gegenereerd.

Voor volledige laadpatronen van DataFrame, zie pandas-integratie en Polars-integratie.

Laad pijlgegevens

Wanneer de bron al in Apache Arrow-formaat is, cursor.bulkcopy_arrow() laad je deze zonder eerst Python-tuples te bouwen.

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

De methode accepteert ook een pyarrow.RecordBatch of een pyarrow.RecordBatchReader, zodat je een resultaatset rechtstreeks naar cursor.arrow_reader() een andere tabel kunt streamen.

Elk Arrow-kolomtype moet compatibel zijn met het bestemmings-SQL-kolomtype, en de schrijver converteert niet tussen typefamilies. Voor meer informatie, zie Apache Arrow-integratie.

Parquet-staging

Gebruik Parquet als tussenformaat bij het migreren van data tussen systemen of wanneer je ETL-pijplijn al Parquet-bestanden produceert. Een Parquet-bestand leest Arrow-data in, dus geef het direct door naar bulkcopy_arrow():

import pyarrow.parquet as pq

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

Doorloop bij grote Parquet-bestanden de rijgroepen om het geheugengebruik constant te houden. Elke batch is een RecordBatch, die bulkcopy_arrow() direct accepteert:

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

Om het hele bestand in één aanroep te streamen, wikkel je de batches in een 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()

Valideer geladen gegevens

Controleer na het laden het aantal rijen en controleer de gegevens steekproefs:

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

Vertrouw voor productieworkloads niet op de transactie van de aanroepende verbinding om een bulkcopy()-aanroep te beveiligen. bulkcopy() opent zijn eigen interne verbinding en legt de gekopieerde rijen onafhankelijk vast, dus een conn.rollback() op je primaire verbinding kan ze niet ongedaan maken. Twee benaderingen geven je atomiciteit:

  • Stel use_internal_transaction=True in dat elke batch in een eigen transactie wordt verpakt. Een batch die halverwege mislukt, wordt teruggedraaid in plaats van gedeeltelijk geladen achter te blijven.
  • Om gegevens te valideren voordat je die doorzet, kopieer je de gegevens in bulk naar een stagingtabel, valideer je ze en verplaats je de rijen vervolgens naar de doeltabel door binnen een transactie via je hoofdverbinding een INSERT ... SELECT te gebruiken. Omdat INSERT via je verbinding wordt uitgevoerd, maakt conn.rollback() dit ongedaan als de validatie mislukt.
# 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