Escolha um padrão de carregamento e movimento de dados com mssql-python

O mssql-python driver fornece múltiplos caminhos para escrever dados no Microsoft SQL. Cada opção adequa-se a diferentes cargas de trabalho. Este guia ajuda-o a escolher o correto com base no volume de dados, formato de origem e semântica de atualização.

Decide por carga de trabalho

Carga de trabalho Caminho recomendado Porquê
Carregar ficheiros CSV numa tabela Carregar dados CSV através de cópia em massa bulkcopy() com um gerador lida com ficheiros de qualquer tamanho sem os carregar na memória.
Inserir uma única linha a partir do código da aplicação Inserções de uma única linha Baixo overhead, gestão de erros simples, funciona com o OUTPUT para devolver as chaves geradas.
Insira um lote pequeno a médio a partir do código da aplicação Inserções em lote Reduz as idas e voltas em comparação com inserções únicas.
Carregue centenas de linhas ou mais de qualquer fonte Cópia em massa A inserção em massa via TDS é a forma mais eficiente para grandes volumes.
Inserir ou atualizar linhas com base numa chave Upsert com MERGE MERGE trata INSERT, UPDATE, e DELETE numa afirmação.
Carregar um DataFrame numa tabela Carregar DataFrames Extrair filas de pandas ou polares e alimentar para bulkcopy().
Armazenar os dados temporariamente em ficheiros Parquet Encenação em parquet Útil para ETL entre sistemas onde é necessário um formato de ficheiro intermédio.

Carregar dados CSV com cópia em bloco

Carregar dados CSV é a questão de ingesta mais comum no trabalho de bases de dados em Python. Utilize csv.reader com um gerador a alimentar 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()

O padrão gerador mantém o uso de memória constante independentemente do tamanho do ficheiro. Para mapeamento de colunas e gestão de identidade, veja Operações de cópia em massa.

Inserções de uma única fila

Utilize operações de inserção individuais para operações de escrita ao nível da aplicação, quando processa um registo de cada vez. Uso OUTPUT INSERTED para recuperar chaves geradas:

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

Os insertos individuais são a escolha certa quando:

  • Insere uma linha por ação do utilizador (submissão de formulário, chamada à API).
  • Tens de validar ou transformar cada linha individualmente antes de inserir.
  • Precisa imediatamente do ID inserido ou de outros valores gerados.

Inserções em lote

Utilize executemany() quando tiver um número moderado de linhas e não precisar da taxa de transferência da cópia em 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() envia cada linha como uma instrução parametrizada separada. Quando a taxa de transferência é mais importante do que o controlo individual por linha, bulkcopy() é mais eficiente porque utiliza o protocolo TDS de inserção em massa. O ponto de transição depende da largura das linhas e da latência da rede, mas situa-se tipicamente nas poucas centenas de linhas.

Cópia de grandes volumes

Quando o throughput importa mais do que o controlo por linha, use bulkcopy(). Utiliza o protocolo TDS de inserção em massa, que é significativamente mais eficiente do que inserções linha a linha:

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

Dicas de desempenho para cópia em massa

  • Use geradores para grandes conjuntos de dados para manter o uso de memória constante.
  • Defina batch_size para controlar quantas linhas são enviadas em cada lote TDS. Começa com 5.000 e ajusta com base na largura das filas.
  • Utilize bloqueios de tabela para carregamentos exclusivos: cursor.bulkcopy("dbo.ProductImport", rows, table_lock=True).
  • Desative os índices antes de carregar e, em seguida, reconstrua-os. Esta sequência evita a sobrecarga de manutenção do índice durante a carga.

Para mapeamentos de colunas, colunas de identidade, processamento de NULL e carregamento paralelo, consulte operações de cópia em massa.

Upsert com MERGE

MERGE é a instrução do Microsoft SQL para INSERT, UPDATE e DELETE condicionais numa só operação. Lida com o padrão "inserir se for novo, atualizar se existir" que os programadores Python normalmente precisam.

Upsert de uma só fila

Para uma única linha, use MERGE com uma USING cláusula que defina aliases de parâmetros:

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 em lote com uma tabela de staging

Para upserts em lote, carregue primeiro os dados para uma tabela temporária e, em seguida, use MERGE para atualizar a partir dessa tabela. Utilize insert-or-update como o padrão predefinido para operações de upsert em DataFrames e atualizações em lote:

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

Este exemplo demonstra o padrão padrão de inserção ou atualização:

  • INSERT linhas da fonte que não existem no alvo (WHEN NOT MATCHED BY TARGET).
  • UPDATE linhas que existem em ambos (WHEN MATCHED).
  • A cláusula OUTPUT reporta que ação foi tomada em cada linha, o que é útil para registos de auditoria.

Atenção

Adicione WHEN NOT MATCHED BY SOURCE THEN DELETE apenas quando os dados de staging forem um snapshot completo e autoritativo do alvo. Se o lote contiver apenas linhas alteradas, essa cláusula elimina linhas que foram intencionalmente omitidas do feed de origem.

Se precisar de uma reconciliação completa, efetue a extensão MERGE apenas depois de confirmar que a fonte é fidedigna para a tabela de destino:

WHEN NOT MATCHED BY SOURCE THEN
    DELETE

Em ambientes partilhados, utilize um nome único de tabela temporária global por execução ou uma tabela permanente de preparação para evitar colisões entre tarefas concorrentes.

Quando usar instruções separadas UPDATE e INSERT em vez disso

MERGE é poderoso, mas tem casos limite. Considere usar instruções separadas quando:

  • Não precisas DELETE de lógica. Ter um elemento separado UPDATE seguido de INSERT WHERE NOT EXISTS é mais legível e mais simples de depurar.
  • A MERGE afirmação é suficientemente complexa para que o comportamento de bloqueio seja difícil de prever. As instruções separadas permitem um controlo explícito sobre a granularidade dos bloqueios.
  • Estás a atualizar uma tabela com elevada concorrência onde MERGE a escalada de bloqueios pode causar bloqueio.
# 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()

Carregar DataFrames

Extrai linhas de um pandas ou Polars DataFrame e carrega-as usando bulkcopy():

pandas

Converta um DataFrame do pandas em tuplas e passe-o para 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

Converter um DataFrame do Polars em tuplos usando o método .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()

Para conhecer as formas completas de carregamento do DataFrame, consulte a integração com pandas e a integração com Polars.

Encenação em parquet

Use o Parquet como formato intermédio ao migrar dados entre sistemas ou quando o seu pipeline ETL já produz ficheiros 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()

Para ficheiros Parquet grandes, leia por grupos de linhas para manter a utilização de memória constante:

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

Validar dados carregados

Após o carregamento, verifique a contagem do número de linhas e faça uma verificação pontual dos dados:

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

Para cargas em produção, não confie na transação da conexão invocante para proteger uma chamada a bulkcopy(). bulkcopy() abre a sua própria ligação interna e confirma de forma independente as linhas que foram copiadas, por isso uma conn.rollback() na tua ligação principal não as pode desfazer. Duas abordagens dão-lhe atomicidade:

  • Define use_internal_transaction=True para envolver cada lote numa transação própria. Um lote que falha a meio do processo anula esse lote, em vez de o deixar parcialmente carregado.
  • Para validar os dados antes de os mover para produção, copie-os em massa para uma tabela de preparação, valide-os e, em seguida, mova as linhas para a tabela de destino utilizando um INSERT ... SELECT dentro de uma transação na ligação principal. Uma vez que isso INSERT é executado na tua conexão, conn.rollback() reverte-o se a validação falhar.
# 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