Usa copia di massa con mssql-python

Il driver mssql-python include una funzione di copia in massa che inserisce in modo efficiente grandi quantità di dati in SQL Server, database SQL di Azure, Istanza gestita di SQL di Azure e database SQL in Microsoft Fabric.

Il cursor.bulkcopy() metodo fornisce un percorso ad alte prestazioni per il caricamento di grandi dataset:

  • Minimizza i viaggi di andata e ritorno in rete.
  • Opzionalmente bypassa il controllo dei vincoli durante il caricamento.
  • Utilizza il protocollo ottimizzato TDS bulk insert.
  • Raggiunge una portata paragonabile a bcp.exe e SqlBulkCopy.

L'estensione nativa basata mssql_py_core su Rust alimenta la funzione di copia in massa. Viene eseguito al di fuori della normale pipeline del cursore execute().

Utilizzo di base

Chiama bulkcopy() su un cursore, passando il nome della tabella di destinazione e un iterabile di tuple di righe o oggetti Row:

Importante

Se crei o modifichi la tabella di destinazione nella stessa sessione, chiama conn.commit() prima di bulkcopy(). Il protocollo bulk copy utilizza un canale interno separato per leggere i metadati della tabella, quindi una modifica DDL non commessa può causare un deadlock o un timeout.

import mssql_python

conn = mssql_python.connect(connection_string)
cursor = conn.cursor()

# Create a temp table for the demo
cursor.execute("""
    CREATE TABLE ##BulkDemo (
        ID INT,
        Name NVARCHAR(50),
        Amount MONEY
    )
""")
conn.commit()

data = [
    (1, "Alice", 50000.00),
    (2, "Bob", 60000.00),
    (3, "Carol", 55000.00),
]

result = cursor.bulkcopy("##BulkDemo", data)
print(f"Copied {result['rows_copied']} rows in {result['batch_count']} batch(es)")
print(f"Elapsed: {result['elapsed_time']}")

Valore restituito

bulkcopy() restituisce un dizionario:

Key TIPO Descrizione
rows_copied int Numero di righe copiate con successo.
batch_count int Numero di lotti processati.
elapsed_time float Tempo impiegato per l'operazione in secondi.

Firma del metodo

cursor.bulkcopy(
    table_name,                    # str - target table (can include schema, e.g. "dbo.MyTable")
    data,                          # Iterable[Tuple | Row] - rows to insert
    batch_size=0,                  # int - rows per batch; 0 = server optimal
    timeout=30,                    # int - operation timeout in seconds
    column_mappings=None,          # List[str] | List[Tuple[int,str]] | None
    keep_identity=False,           # bool - preserve identity values from source
    check_constraints=False,       # bool - check constraints during load
    table_lock=False,              # bool - use table-level lock
    keep_nulls=False,              # bool - preserve NULLs instead of defaults
    fire_triggers=False,           # bool - fire INSERT triggers on target
    use_internal_transaction=False, # bool - use internal transaction per batch
)

Mappatura delle colonne

Per impostazione predefinita, bulkcopy() mappa le colonne in base alla posizione ordinale. Ogni colonna di dati corrisponde alla colonna della tabella con lo stesso indice. Usa il column_mappings parametro per sovrascrivere questo comportamento.

Elenco dei nomi delle colonne

Ogni posizione nell'elenco corrisponde all'indice dei dati di origine:

result = cursor.bulkcopy(
    "##BulkDemo",
    data,
    column_mappings=["ID", "Name", "Amount"],
)

Formato avanzato: mappatura esplicita degli indici

Ogni tupla assume la forma (source_index, target_column_name). Usa questo formato per saltare o riordinare le colonne:

result = cursor.bulkcopy(
    "##BulkDemo",
    data,
    column_mappings=[(0, "ID"), (1, "Name"), (2, "Amount")],
)

Caricamento dai file

Puoi caricare dati da file CSV e altri formati passando un generatore a bulkcopy().

File CSV

import csv
import io
import mssql_python

# In production, replace io.StringIO with open("data.csv", "r", ...)
csv_data = """ID,Name,Value
1,Widget,9.99
2,Gadget,24.50
3,Gizmo,4.75
"""

def csv_row_generator(file_obj):
    """Generator that yields tuples from a CSV file object."""
    reader = csv.reader(file_obj)
    next(reader)  # Skip header
    for row in reader:
        if row:  # skip blank lines
            yield (
                int(row[0]),      # ID
                row[1],           # Name
                float(row[2]),    # Value
            )

conn = mssql_python.connect(connection_string)
cursor = conn.cursor()

cursor.execute("""
    CREATE TABLE ##CSVImport (ID INT, Name NVARCHAR(100), Value FLOAT)
""")
conn.commit()
result = cursor.bulkcopy("##CSVImport", csv_row_generator(io.StringIO(csv_data)))
print(f"Imported {result['rows_copied']} rows from CSV")

File di grandi dimensioni con elaborazione in batch

Imposta il parametro batch_size per controllare quante righe il driver invia in ogni batch. Questo approccio funziona bene per file di grandi dimensioni:

import csv
import io
import mssql_python

# In production, replace io.StringIO with open("large_file.csv", "r", ...)
csv_data = "\n".join(
    ["ID,Name,Value"] + [f"{i},Item {i},{i * 1.5}" for i in range(1, 201)]
)

def csv_rows(file_obj):
    reader = csv.reader(file_obj)
    next(reader)  # Skip header
    for row in reader:
        if row:
            yield (int(row[0]), row[1], float(row[2]))

conn = mssql_python.connect(connection_string)
cursor = conn.cursor()

cursor.execute("""
    CREATE TABLE ##LargeCSV (ID INT, Name NVARCHAR(100), Value FLOAT)
""")
conn.commit()
result = cursor.bulkcopy(
    "##LargeCSV",
    csv_rows(io.StringIO(csv_data)),
    batch_size=50,
)
print(f"Imported {result['rows_copied']} rows in {result['batch_count']} batches")

Carica i DataFrame di pandas

Un DataFrame è organizzato per colonne, quindi il percorso più veloce è bulkcopy_arrow(), che accetta la tabella Arrow che pandas sa già produrre. bulkcopy() prende tuple di righe, quindi devi prima convertire le colonne in oggetti Python.

Converti la tabella Arrow nei tipi delle colonne di destinazione prima di caricarla. pyarrow deduce float64 per una colonna numerica, che il conducente non può mappare a denaro, decimale o numerico:

import pandas as pd
import pyarrow as pa
import mssql_python

df = pd.DataFrame({
    'ID': [1, 2, 3],
    'Name': ['Alice', 'Bob', 'Carol'],
    'Amount': [50000.0, 60000.0, 55000.0],
})

conn = mssql_python.connect(connection_string)
cursor = conn.cursor()

cursor.execute("""
    CREATE TABLE ##PandasDemo (ID INT, Name NVARCHAR(50), Amount MONEY)
""")
conn.commit()

target = pa.schema([
    pa.field('ID', pa.int32()),
    pa.field('Name', pa.string()),
    pa.field('Amount', pa.decimal128(19, 4)),   # MONEY
])

table = pa.Table.from_pandas(df, preserve_index=False).cast(target)
result = cursor.bulkcopy_arrow("##PandasDemo", table)

Senza il lancio, il carico fallisce con ValueError: Cannot map Arrow column 'Amount' (Float64) to SQL column 'Amount' (Money). Crea il cast con Table.cast() invece di passare lo schema a Table.from_pandas(), che non può convertire direttamente una colonna in virgola mobile in decimal128. NaN i valori diventano SQL NULL su questo percorso, quindi non è necessario sostituirli prima.

Se invece hai bisogno del percorso riga-tuple, itertuples() ottieni già tuple quando passi name=None:

data = list(df.itertuples(index=False, name=None))
result = cursor.bulkcopy("##PandasDemo", data)

Carica i dati di Apache Arrow

Usalo cursor.bulkcopy_arrow() per caricare i dati di Apache Arrow. Questo metodo legge direttamente dalla memoria Arrow, quindi non è necessario costruire tuple di righe Python prima di chiamare il metodo.

L'argomento source accetta un pyarrow.Table, a pyarrow.RecordBatch, un pyarrow.RecordBatchReader, o qualsiasi oggetto che esponga l'interfaccia dati Arrow C. Gli argomenti rimanenti sono gli stessi di bulkcopy().

import mssql_python
import pyarrow as pa

conn = mssql_python.connect(connection_string)

# bulkcopy_arrow() opens its own connection, so commit the table creation first.
conn.autocommit = True
cursor = conn.cursor()

cursor.execute("""
    CREATE TABLE ##ArrowDemo (ID INT, Name NVARCHAR(50), Amount FLOAT)
""")

table = pa.table({
    "ID": pa.array([1, 2, 3], type=pa.int32()),
    "Name": pa.array(["Alice", "Bob", "Carol"], type=pa.string()),
    "Amount": pa.array([50000.0, 60000.0, 55000.0], type=pa.float64()),
})

result = cursor.bulkcopy_arrow("##ArrowDemo", table)
print(f"Copied {result['rows_copied']} rows")

Ogni tipo di colonna freccia deve essere compatibile con il tipo di colonna SQL di destinazione. Il writer non converte tra famiglie di tipi, quindi il passaggio di una colonna float64 a una colonna money genera ValueError prima che venga scritta qualsiasi riga. Usa decimal128 per colonne di denaro, decimali e numeriche .

Il passaggio di una sorgente Arrow a bulkcopy() genera TypeError e reindirizza a bulkcopy_arrow().

Per maggiori informazioni sul supporto Arrow, incluso come trasmettere un set di risultati da una tabella all'altra, vedi integrazione con Apache Arrow.

Gestire i valori NULL

Passa None in qualsiasi posizione di colonna per inserire un valore SQL NULL :

cursor.execute("""
    CREATE TABLE ##NullDemo (ID INT, Name NVARCHAR(50), Amount MONEY)
""")
conn.commit()

data = [
    (1, "Alice", 50000.00),
    (2, "Bob", None),       # NULL Amount
    (3, None, 55000.00),    # NULL Name
]

cursor.bulkcopy("##NullDemo", data)

Colonne di identità

Per inserire valori identità espliciti, impostare keep_identity=True:

cursor.execute("""
    CREATE TABLE ##IdentDemo (ID INT, Name NVARCHAR(50), Amount MONEY)
""")
conn.commit()

data = [
    (100, "Alice", 50000.00),
    (200, "Bob", 60000.00),
]

cursor.bulkcopy("##IdentDemo", data, keep_identity=True)

Quando keep_identity=False (il valore predefinito), ometti la colonna identità dai tuoi dati e usala column_mappings per indirizzare le colonne non identità.

Opzioni di copia in blocco

Parametro Impostazione predefinita Descrizione
batch_size 0 Righe per blocco. 0 Permette al server di scegliere la dimensione ottimale.
timeout 30 Timeout dell'operazione in secondi. Si applica all'operazione di copia di massa stessa, non alla connessione interna.
keep_identity False Preservare i valori ID dei dati di origine.
check_constraints False Controllare i vincoli della tabella durante il caricamento.
table_lock False Acquisire un blocco a livello di tabella invece di blocchi a livello di riga.
keep_nulls False Preservare i valori NULL invece di inserire valori predefiniti delle colonne.
fire_triggers False Attiva INSERT trigger sulla tabella di destinazione.
use_internal_transaction False Racchiudi ogni batch in una transazione interna.

Note

bulkcopy() apre una connessione interna separata al server. Quella connessione interna eredita il timeout della query associata al cursore: imposta Connection.timeout su un valore positivo prima di creare il cursore e lo stesso valore definisce anche il tempo massimo per il tentativo di connessione per la copia in blocco. Se il timeout della query del cursore è 0, la connessione interna utilizza il timeout predefinito di 15 secondi di connessione. Un cursore prende il valore quando viene creato, quindi cambiare Connection.timeout dopo non influisce su un cursore esistente o una copia in blocco in volo. Aumenta il timeout della query prima di creare il cursore per endpoint lenti, soggetti a limitazione o ad alta latenza (ad esempio, tramite VPN o tra diverse aree geografiche).

Gestire gli errori

bulkcopy() solleva un'eccezione se il carico fallisce, quindi avvolgi la chiamata in un try/except blocco per rilevare gli errori. Tieni presente che bulkcopy() viene eseguito su una propria connessione interna e conferma le righe copiate in modo indipendente, quindi un conn.rollback() sulla connessione principale non può annullarle. Per rendere atomico un batch, imposta use_internal_transaction=True, che racchiude ogni batch nella propria transazione, che viene automaticamente annullata se il batch non riesce:

import mssql_python

conn = mssql_python.connect(connection_string)
cursor = conn.cursor()

cursor.execute("""
    CREATE TABLE ##ImportDemo (ID INT, Name NVARCHAR(50), Value FLOAT)
""")
conn.commit()

data = [
    (1, "Alice", 50000.00),
    (2, "Bob", 60000.00),
    (3, "Carol", 55000.00),
]

try:
    result = cursor.bulkcopy("##ImportDemo", data, use_internal_transaction=True)
    print(f"Successfully copied {result['rows_copied']} rows")
except (mssql_python.DatabaseError, ValueError) as e:
    # bulkcopy() commits on its own connection, so there's nothing to roll back
    # here. With use_internal_transaction=True, a failed batch is already rolled
    # back on the bulk copy connection.
    print(f"Bulk copy failed: {e}")

Per subordinare un caricamento alla tua logica di convalida, esegui una copia in blocco in una tabella di staging, quindi sposta le righe nella tabella di destinazione con un'INSERT ... SELECT all'interno di una transazione nella connessione principale. Questo INSERT funziona sulla tua connessione, quindi conn.rollback() annulla se la validazione fallisce.

Authentication

La copia di massa utilizza un canale interno separato che richiede un proprio token. Il driver gestisce automaticamente l'acquisizione dei token per i metodi di autenticazione supportati.

Identità gestita (ActiveDirectoryMSI)

Utilizzare Authentication=ActiveDirectoryMSI per identità gestita assegnata dal sistema o dall'utente. Questo metodo di autenticazione è consigliato per servizi ospitati su Azure come Azure VM, App Service, Functions e AKS.

import mssql_python

# System-assigned managed identity
conn = mssql_python.connect(
    "Server=<server>.database.windows.net;"
    "Database=<database>;"
    "Authentication=ActiveDirectoryMSI;"
    "Encrypt=yes"
)
cursor = conn.cursor()

cursor.execute("CREATE TABLE ##MsiDemo (ID INT, Name NVARCHAR(50))")
conn.commit()

result = cursor.bulkcopy("##MsiDemo", [(1, "Alice"), (2, "Bob")])
print(f"Copied {result['rows_copied']} rows")

Per un'identità gestita assegnata dall'utente, passa l'ID client nella stringa di connessione:

conn = mssql_python.connect(
    "Server=<server>.database.windows.net;"
    "Database=<database>;"
    "Authentication=ActiveDirectoryMSI;"
    "UID=<client-id>;"
    "Encrypt=yes"
)

Principale di servizio (ActiveDirectoryServicePrincipal)

Utilizzare Authentication=ActiveDirectoryServicePrincipal per l'autenticazione del principale servizio (credenziali client).

conn = mssql_python.connect(
    "Server=<server>.database.windows.net;"
    "Database=<database>;"
    "Authentication=ActiveDirectoryServicePrincipal;"
    "UID=<application-client-id>;"
    "PWD=<client-secret>;"
    "Encrypt=yes"
)
cursor = conn.cursor()

cursor.execute("CREATE TABLE ##SpDemo (ID INT, Value FLOAT)")
conn.commit()

result = cursor.bulkcopy("##SpDemo", [(1, 1.5), (2, 2.5)])
print(f"Copied {result['rows_copied']} rows")

Catena di credenziali predefinita (ActiveDirectoryDefault)

ActiveDirectoryDefault Prova più fornitori di credenziali in sequenza, come variabili ambientali, identità del carico di lavoro, identità gestita e altro ancora. Funziona sia per lo sviluppo locale che per i servizi ospitati su Azure senza modifiche al codice.

Per maggiori informazioni sull'autenticazione, vedi Microsoft Entra autentication.

Suggerimenti per le prestazioni

Le seguenti tecniche ti aiutano a massimizzare la velocità di copia in massa.

Inizia da un'origine colonnare

bulkcopy()prende un iterabile di tuple righe, quindi ogni valore deve esistere come oggetto Python prima che la copia inizi. Quando i dati sono già colonnari, bulkcopy_arrow() legge direttamente i buffer Arrow e salta quel passaggio. Un DataFrame pandas o Polars, un file Parquet e il risultato di cursor.arrow() sono tutte origini Arrow. Per maggiori informazioni, vedi Carica dati di Apache Arrow.

Usa generatori per grandi dataset

I generatori minimizzano l'uso della memoria perché bulkcopy() accettano qualsiasi iterabile:

def data_generator(count):
    """Generate rows without loading all into memory."""
    for i in range(count):
        yield (i, f"Item {i}", i * 1.5)

cursor = conn.cursor()
cursor.execute("""
    CREATE TABLE ##LargeDemo (ID INT, Name NVARCHAR(50), Value FLOAT)
""")
conn.commit()
result = cursor.bulkcopy("##LargeDemo", data_generator(1000))

Usa serrature da tavolo per carichi più veloci

Quando non hai lettori concorrenti, imposta table_lock=True per ridurre il sovraccarico di blocco durante carichi iniziali elevati.

result = cursor.bulkcopy(
    "##LargeDemo",
    data,
    table_lock=True,
    batch_size=100000,
)

Disabilita gli indici durante il caricamento

Disabilita temporaneamente gli indici non clusterizzati prima del caricamento in massa e ricostruiscili successivamente per migliorare le prestazioni:

cursor = conn.cursor()

cursor.execute("""
    CREATE TABLE ##IndexDemo (ID INT, Name NVARCHAR(50), Value FLOAT)
""")
cursor.execute("CREATE NONCLUSTERED INDEX IX_Name ON ##IndexDemo(Name)")
conn.commit()

cursor.execute("ALTER INDEX IX_Name ON ##IndexDemo DISABLE")
conn.commit()

result = cursor.bulkcopy("##IndexDemo", data)
conn.commit()

cursor.execute("ALTER INDEX IX_Name ON ##IndexDemo REBUILD")
conn.commit()

Tabelle di carico in parallelo

Apri una connessione separata per ogni tabella ed esegui i carichi contemporaneamente.

import concurrent.futures

def load_table(table_name, rows):
    conn = mssql_python.connect(connection_string)
    cursor = conn.cursor()
    cursor.execute(f"CREATE TABLE {table_name} (ID INT, Name NVARCHAR(50), Value FLOAT)")
    conn.commit()
    result = cursor.bulkcopy(table_name, rows)
    conn.commit()
    conn.close()
    return result["rows_copied"]

data = [(i, f"Item {i}", i * 1.5) for i in range(100)]

with concurrent.futures.ThreadPoolExecutor(max_workers=3) as executor:
    futures = [
        executor.submit(load_table, "##Load1", data),
        executor.submit(load_table, "##Load2", data),
        executor.submit(load_table, "##Load3", data),
    ]
    for future in concurrent.futures.as_completed(futures):
        print(f"Loaded {future.result()} rows")

Confronto con le alternative

La tabella seguente confronta il copio di massa con altri metodi di inserimento dati.

metodo Caso di utilizzo Prestazioni
cursor.bulkcopy_arrow() Set di dati di grandi dimensioni che sono già colonnari. Il più veloce
cursor.bulkcopy() Grandi dataset (più di 1.000 righe) provenienti da fonti orientate alle righe. Veloce
cursor.executemany() Dataset medi con parametri. Moderate
cursor.execute() in un ciclo Piccoli dataset con una logica semplice. Più lento