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

O mssql-python driver oferece múltiplos caminhos para gravar dados no Microsoft SQL. Cada caminho se encaixa em diferentes cargas de trabalho. Este guia ajuda você a escolher a opção certa com base no volume de dados, formato de origem e semântica de atualização.

Decida com base na carga de trabalho

Carga de Trabalho Caminho recomendado Por que
Carregar arquivos CSV em uma tabela Carregar dados CSV com cópia em massa bulkcopy() com um gerador lida com arquivos de qualquer tamanho sem precisar carregá-los na memória.
Insira uma única linha do código da aplicação Inserções de linha única Baixo overhead, tratamento direto de erros, funciona com OUTPUT para retornar chaves geradas.
Insira um lote pequeno a moderado a partir do código da aplicação Inserções em lote Reduz o número de comunicações de ida e volta 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.
Insira ou atualize linhas com base em uma chave Upsert com MERGE MERGE trata INSERT, UPDATE, e DELETE em uma única afirmação.
Carregar um DataFrame em uma tabela Carregar DataFrames bulkcopy_arrow()lê os dados Arrow do DataFrame sem construir um objeto Python para cada valor.
Carregar os dados do Apache Arrow em uma tabela Carregar dados Arrow bulkcopy_arrow()lê a memória do Arrow diretamente, sem construir tuplas em Python.
Dados de estágio através de arquivos Parquet Preparação de Parquet Útil para ETL entre sistemas onde é necessário um formato de arquivo intermediário. Parquet já está em formato Arrow, portanto é carregado sem conversão linha por linha.

Carregar dados CSV com cópia em massa

Carregar dados CSV é a pergunta de ingestão mais comum para trabalhos com bancos de dados em Python. Use csv.reader com um gerador que 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()

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

Inserções de uma única linha

Use inserções individuais para operações de gravação no nível da aplicação, quando você processa um registro por vez. Use 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()

Insertos individuais são a escolha certa quando:

  • Você insere uma linha por ação do usuário (envio de formulário, chamada de API).
  • Você precisa validar ou transformar cada linha individualmente antes de inserir.
  • Você precisa do ID inserido ou outros valores gerados imediatamente.

Inserções em lote

Use executemany() quando você tiver um número moderado de linhas e não precisa 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 importa mais do que o controle por linha, bulkcopy() é mais eficiente porque utiliza o protocolo TDS de inserção em massa. O ponto de transição depende da largura da linha e da latência da rede, mas normalmente fica na casa das poucas centenas de linhas.

Cópia em lote

Quando a taxa de transferência importa mais do que o controle por linha, use bulkcopy(). Ele usa o protocolo de inserção em massa TDS, que transmite linhas em fluxo em vez de enviar uma instrução por 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 conjuntos de dados grandes e mantenha o consumo de memória constante.
  • Use bulkcopy_arrow() quando a fonte for colunar, como um arquivo DataFrame ou Parquet. Ele ignora a conversão para tuplas de linha do Python.
  • Defina batch_size para controlar quantas linhas são enviadas em cada lote TDS. Comece com 5.000 e ajuste de acordo com a largura da linha.
  • Use bloqueios de tabela para cargas exclusivas: cursor.bulkcopy("dbo.ProductImport", rows, table_lock=True).
  • Desative os índices antes do carregamento e reconstrua-os depois. Essa sequência evita a sobrecarga de manutenção do índice durante a carga.

Para mapeamentos de coluna, colunas de identidade, tratamento 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 em uma única operação. Ele lida com o padrão "inserir se for novo, atualizar se existir" que os desenvolvedores Python normalmente precisam.

Execução de upsert de linha única

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

Execução de upsert em lote com tabela de preparação

Para execuções de upserts em massa, coloque os dados em uma tabela temporária primeiro e depois use MERGE para atualizar a partir dela. Use inserir ou atualizar como o padrão para operações de execução de upsert e atualizações em lote em DataFrames:

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 destino (WHEN NOT MATCHED BY TARGET).
  • UPDATE linhas presentes em ambos (WHEN MATCHED).
  • A cláusula OUTPUT informa qual ação foi tomada em cada linha, o que é útil para trilhas de auditoria.

Cuidado

Adicione WHEN NOT MATCHED BY SOURCE THEN DELETE apenas quando os dados de preparação forem um snapshot completo e autoritativo do destino. Se o lote contém apenas linhas alteradas, essa cláusula exclui linhas que foram intencionalmente omitidas do feed de origem.

Se precisar de reconciliação completa, estenda a MERGE apenas depois de confirmar que a fonte é autoritativa para a tabela de destino:

WHEN NOT MATCHED BY SOURCE THEN
    DELETE

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

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

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

  • Você não precisa de lógica DELETE. Um UPDATE separado, seguido de INSERT WHERE NOT EXISTS, é mais legível e mais simples de depurar.
  • A MERGE instrução é complexa o suficiente para que o comportamento de bloqueio seja difícil de prever. Declarações separadas permitem controle explícito sobre a granularidade do bloqueio.
  • Você está atualizando uma tabela de alta concorrência na qual um MERGE escalonamento de bloqueio pode causar bloqueios.
# 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

Um DataFrame é colunar, então carregue-o usando bulkcopy_arrow() em vez de convertê-lo em tuplas de linha para bulkcopy().

Associe os tipos do Arrow às colunas de destino antes de carregar. pyarrow infere float64 para uma coluna numérica, que o motorista não pode mapear para dinheiro, decimal ou numérico.

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

Use Table.cast() , em vez de passar o esquema para Table.from_pandas(), que não pode converter uma coluna float em decimal128 diretamente.

Polars

Polars implementa a interface de dados Arrow C, então você pode passar o próprio DataFrame. Molde as colunas primeiro pelo mesmo 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()

Você também pode definir os tipos ao ler o arquivo, com pl.read_csv("products.csv", schema_overrides={"ListPrice": pl.Decimal(19, 4)}).

Ao passar o DataFrame diretamente, seus buffers são passados ao driver sem fazer uma cópia. df.to_arrow() também funciona, mas Polars recodifica as colunas de texto durante essa conversão, o que copia todos os dados de texto.

bulkcopy_arrow() aceita um pyarrow.Table, a RecordBatch, a RecordBatchReader, ou qualquer objeto que implemente a interface de dados Arrow C através de __arrow_c_stream__ ou __arrow_c_array__. Fornecer qualquer um destes para bulkcopy() gera TypeError.

Para padrões completos de carregamento do DataFrame, veja integração com pandas e integração com Polars.

Carregar dados do Arrow

Quando a origem já está no formato Apache Arrow, cursor.bulkcopy_arrow() a carrega sem primeiro criar tuplas em 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")

O método também aceita um pyarrow.RecordBatch ou um pyarrow.RecordBatchReader, para que você possa transferir um conjunto de resultados de cursor.arrow_reader() diretamente para outra tabela.

Cada tipo de coluna Arrow deve ser compatível com seu tipo de coluna SQL de destino, e o escritor não converte entre famílias de tipos. Para mais informações, veja integração com o Apache Arrow.

Preparação do Parquet

Use o Parquet como formato intermediário ao migrar dados entre sistemas ou quando seu pipeline ETL já produz arquivos Parquet. Um arquivo Parquet lê dados do Arrow, então passe diretamente para bulkcopy_arrow():

import pyarrow.parquet as pq

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

Para arquivos grandes de Parquet, itere grupos de linhas para manter o uso de memória constante. Cada lote é um RecordBatch, que bulkcopy_arrow() aceita diretamente:

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

Para transmitir o arquivo inteiro em uma única chamada, envolva os lotes em um 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()

Validar dados carregados

Após o carregamento, verifique a contagem de linhas e verifique os dados pontuais:

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 de trabalho em produção, não dependa da transação da conexão chamadora para proteger uma chamada bulkcopy(). bulkcopy() abre sua própria conexão interna e confirma as linhas copiadas de forma independente, de forma que um conn.rollback() na conexão principal não possa desfazê-las. Duas abordagens garantem atomicidade:

  • Configure use_internal_transaction=True para envolver cada lote em sua própria transação. Um lote que falha durante o processamento é revertido, em vez de permanecer parcialmente carregado.
  • Para validar os dados antes de promovê-los, copie-os em massa para uma tabela de preparação, valide-os e, em seguida, mova as linhas para a tabela de destino usando um(a) INSERT ... SELECT dentro de uma transação na conexão principal. Como isso INSERT é executado na sua conexão, conn.rollback() o desfaz caso a validação falhe.
# 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