Scegli un pattern di caricamento e movimento dati con mssql-python

Il mssql-python driver offre molteplici percorsi per la scrittura di dati in Microsoft SQL. Ogni percorso si adatta a carichi di lavoro diversi. Questa guida ti aiuta a scegliere quella giusta in base al volume dei dati, al formato della fonte e alla semantica degli aggiornamenti.

Decidi in base al carico di lavoro

Carico di lavoro Percorso consigliato Perché
Carica file CSV in una tabella Carica dati CSV con copia in blocco bulkcopy() con un generatore gestisce file di qualsiasi dimensione senza caricarli in memoria.
Inserisci una singola riga dal codice dell'applicazione Inserti a riga singola Basso overhead, gestione degli errori semplice, funziona con OUTPUT per restituire le chiavi generate.
Inserisci un lotto piccolo o moderato dal codice applicativo Inserimenti in batch Riduce i viaggi di andata e ritorno rispetto agli inserti singoli.
Carica centinaia di righe o più da qualsiasi fonte Copia in blocco L'inserimento in massa TDS è il percorso più efficiente per volumi grandi.
Inserisci o aggiorna righe basandoti su una chiave Upsert con MERGE MERGE gestisce INSERT, UPDATE, e DELETE in una sola affermazione.
Carica un DataFrame in una tabella Carica DataFrame bulkcopy_arrow()legge i dati Arrow del DataFrame senza costruire un oggetto Python per ogni valore.
Carica i dati di Apache Arrow in una tabella Carica dati Arrow bulkcopy_arrow()legge direttamente la memoria Arrow, senza costruire tuple Python.
Caricare temporaneamente i dati tramite file Parquet Scenografia in parquet Utile per ETL cross-system dove è necessario un formato file intermedio. Parquet è già in formato Arrow, quindi viene caricato senza conversione delle righe.

Carica dati CSV mediante copia in blocco

Il caricamento dei dati CSV è la domanda di ingest più comune per il lavoro su database Python. Usa csv.reader con un generatore che alimenta bulkcopy():

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

Il pattern generatore mantiene costante l'uso della memoria indipendentemente dalla dimensione del file. Per la mappatura delle colonne e la gestione delle identità, vedi Operazioni di copia in massa.

Inserimenti di una singola riga

Usa inserti singoli per le scritture a livello applicativo, dove elabori un record alla volta. Utilizzare OUTPUT INSERTED per recuperare le chiavi generate:

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

Gli inserti singoli sono la scelta giusta quando:

  • Inserisci una riga per ogni azione utente (invio del modulo, chiamata API).
  • Devi validare o trasformare ogni riga singolarmente prima di inserirla.
  • Ti serve subito l'ID inserito o altri valori generati.

Inserimenti in batch

Usa executemany() quando hai un numero moderato di righe e non hai bisogno della velocità di throughput delle copie in massa:

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() invia ogni riga come un'istruzione parametrizzata separata. Quando la velocità di trasmissione conta più del controllo per riga, bulkcopy() è più efficiente perché utilizza il protocollo TDS bulk insert. Il crossover dipende dalla larghezza delle righe e dalla latenza di rete, ma di solito si trova nelle poche centinaia di righe.

Copia in blocco

Quando la velocità di rendimento conta più del controllo per riga, si usa bulkcopy(). Utilizza il protocollo TDS bulk insert, che trasmette le righe invece di inviare una sola istruzione per riga:

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

Suggerimenti per le prestazioni della copia in blocco

  • Usa generatori per grandi dataset per mantenere costante l'uso della memoria.
  • Uso bulkcopy_arrow() quando la sorgente è colonnare, come un DataFrame o un file Parquet. Ignora la conversione in tuple di righe Python.
  • Set batch_size per controllare quante righe vengono inviate per ogni batch TDS. Inizia con 5.000 e aggiusta in base alla larghezza delle file.
  • Usa serrature da tavolo per carichi esclusivi: cursor.bulkcopy("dbo.ProductImport", rows, table_lock=True).
  • Disabilita gli indici prima di caricare, poi ricostruisci dopo. Questa sequenza evita il sovraccarico di manutenzione dell'indice durante il carico.

Per le mappature di colonne, colonne identità, gestione NULL e caricamento parallelo, vedi Operazioni di copia in massa.

Upsert con MERGE

MERGE è l'istruzione di SQL di Microsoft per eseguire in modo condizionale INSERT, UPDATE e DELETE in un'unica operazione. Gestisce il pattern "inserisci se nuovo, aggiorna se esiste" di cui gli sviluppatori Python hanno comunemente bisogno.

Upsert a riga singola

Per una singola riga, si usa MERGE con una USING clausola che definisce gli alias dei parametri:

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

Upsert in blocco con tavolo di preparazione

Per le operazioni di upsert in blocco, carica prima i dati in una tabella temporanea, quindi usa MERGE per aggiornare a partire da essa. Usa insert-or-update come modello predefinito per le operazioni di upsert dei DataFrame e gli aggiornamenti batch:

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

Questo esempio dimostra il pattern predefinito di inserimento o aggiornamento:

  • INSERT righe dalla fonte che non esistono nel target (WHEN NOT MATCHED BY TARGET).
  • UPDATE righe presenti in entrambi (WHEN MATCHED).
  • La clausola OUTPUT riporta quale azione è stata compiuta su ciascuna riga, il che è utile per le audit trail.

Attenzione

Aggiungi WHEN NOT MATCHED BY SOURCE THEN DELETE solo quando i dati di staging sono un'istantanea completa e autorevole della destinazione. Se il batch contiene solo righe modificate, quella clausola elimina le righe che sono state intenzionalmente omesse dal feed sorgente.

Se hai bisogno di una riconciliazione completa, estendi MERGE solo dopo aver confermato che la fonte è autorevole per la tabella di destinazione:

WHEN NOT MATCHED BY SOURCE THEN
    DELETE

Negli ambienti condivisi, utilizzare un nome univoco globale per la tabella temporanea a ogni esecuzione oppure una tabella di staging permanente per evitare collisioni tra processi concorrenti.

Quando usare invece istruzioni separate UPDATE e INSERT

MERGE è potente ma ha casi limite. Considera di usare affermazioni separate quando:

  • Non hai bisogno di DELETE logica. Un separato UPDATE seguito da INSERT WHERE NOT EXISTS è più leggibile e semplice da debug.
  • L'affermazione MERGE è abbastanza complessa da rendere difficile prevedere il comportamento di blocco. Le istruzioni distinte consentono di controllare esplicitamente la granularità del blocco.
  • Stai aggiornando una tabella con elevata concorrenza in cui MERGE l'escalation dei lock potrebbe causare blocchi.
# 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()

Carica DataFrame

Un DataFrame è a colonne, quindi caricalo con bulkcopy_arrow() invece di trasformarlo in tuple di righe per bulkcopy().

Abbina i tipi di frecce alle colonne di destinazione prima del caricamento. pyarrow deduce float64 per una colonna numerica, che il conducente non può mappare in denaro, decimale o numerico.

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

Usa Table.cast() invece di passare lo schema a Table.from_pandas(), che non può convertire direttamente una colonna float in decimal128 .

Polars

Polars supporta l'interfaccia dati C di Arrow, quindi puoi passare il DataFrame stesso. Getta prima le colonne per lo stesso motivo:

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

Puoi anche impostare i tipi quando leggi il file, con pl.read_csv("products.csv", schema_overrides={"ListPrice": pl.Decimal(19, 4)}).

Passare direttamente il DataFrame fornisce i relativi buffer al driver senza effettuare una copia. df.to_arrow() funziona anch'esso, ma Polars ricodifica le colonne stringhe durante quella conversione, copiando così tutti i dati della stringa.

bulkcopy_arrow() accetta un pyarrow.Table, un RecordBatch, un RecordBatchReader, o qualsiasi oggetto che implementi l'interfaccia dati Arrow C tramite __arrow_c_stream__ o __arrow_c_array__. Passando uno qualsiasi di questi a bulkcopy() viene generato TypeError.

Per i pattern completi di caricamento DataFrame, vedi integrazione con pandas e integrazione con Polars.

Carica i dati della freccia

Quando il sorgente è già in formato Apache Arrow, cursor.bulkcopy_arrow() caricalo senza prima costruire tuple 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")

Il metodo accetta anche un pyarrow.RecordBatch o un pyarrow.RecordBatchReader, quindi puoi trasmettere un insieme di risultati direttamente cursor.arrow_reader() in un'altra tabella.

Ogni tipo di colonna Arrow deve essere compatibile con il tipo di colonna SQL di destinazione, e lo scrittore non converte tra famiglie di tipi. Per maggiori informazioni, vedi integrazione con Apache Arrow.

Scenografia in parquet

Usa Parquet come formato intermedio quando migri dati tra sistemi o quando la tua pipeline ETL produce già file Parquet. Un file Parquet legge i dati Arrow, quindi passalo direttamente a bulkcopy_arrow():

import pyarrow.parquet as pq

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

Per i file Parquet di grandi dimensioni, iterare i gruppi di righe per mantenere costante l'uso della memoria. Ogni lotto è un RecordBatch, che bulkcopy_arrow() accetta direttamente:

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

Per trasmettere l'intero file in una sola chiamata, avvolgi i batch in un 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()

Valida i dati caricati

Dopo il caricamento, verifica il conteggio delle righe e controlla a campione i dati:

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

Per i carichi di produzione, non affidarti alla transazione della connessione chiamante per proteggere una bulkcopy() chiamata. bulkcopy() apre una propria connessione interna e conferma in modo indipendente le righe copiate, quindi un conn.rollback() nella connessione principale non può annullarne gli effetti. Due approcci ti danno atomicità:

  • Imposta use_internal_transaction=True per racchiudere ogni lotto in una transazione separata. Un lotto che fallisce a metà processo annulla quel lotto invece di lasciarlo a metà carica.
  • Per convalidare i dati prima di caricarli nella tabella finale, esegui una copia in blocco in una tabella di staging, convalida i dati, quindi sposta le righe nella tabella di destinazione usando un INSERT ... SELECT all'interno di una transazione nella connessione principale. Poiché questo INSERT funziona sulla tua connessione, conn.rollback() annulla se la validazione fallisce.
# 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