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 Estrai righe da pandas o Polars e passale a bulkcopy().
Caricare temporaneamente i dati tramite file Parquet Scenografia in parquet Utile per ETL cross-system dove è necessario un formato file intermedio.

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 è significativamente più efficiente rispetto agli inserti riga 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.
  • 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

Estrai le righe da un pandas o da un DataFrame Polar e caricale usando bulkcopy():

pandas

Converti un DataFrame pandas in tuple e passa a 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

Converti un DataFrame Polars in tuple usando il metodo .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()

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

Scenografia in parquet

Usa Parquet come formato intermedio quando si migra dati tra sistemi o quando la tua pipeline ETL produce già file 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()

Per file Parquet di grandi dimensioni, leggi in gruppi di righe per mantenere costante l'uso della memoria:

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

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