Elige un patrón de carga y movimiento de datos con mssql-python

El mssql-python controlador proporciona múltiples vías para escribir datos en Microsoft SQL. Cada opción se adapta a distintas cargas de trabajo. Esta guía te ayuda a elegir la adecuada según el volumen de datos, el formato de la fuente y la semántica de la actualización.

Decide por carga de trabajo

Carga de trabajo Ruta de acceso recomendada Por qué
Cargar archivos CSV en una tabla Cargar datos CSV mediante copia masiva bulkcopy() con un generador gestiona archivos de cualquier tamaño sin cargarlos en la memoria.
Insertar una sola fila del código de la aplicación Inserciones de una sola fila Baja sobrecarga, manejo de errores sencillo, funciona con OUTPUT para devolver las claves generadas.
Insertar un lote pequeño o moderado desde el código de la aplicación Inserciones por lotes Reduce los viajes de ida y vuelta en comparación con las inserciones individuales.
Carga cientos de filas o más desde cualquier fuente Copia en bloque La inserción masiva mediante TDS es la vía más eficiente para grandes volúmenes.
Insertar o actualizar filas basándose en una clave Upsert con MERGE MERGE maneja INSERT, UPDATE, y DELETE en una sola afirmación.
Carga un DataFrame en una tabla Cargar DataFrames bulkcopy_arrow()lee los datos de Arrow del DataFrame sin construir un objeto Python para cada valor.
Carga los datos de Apache Arrow en una tabla Cargar datos de Arrow bulkcopy_arrow() lee la memoria de Arrow directamente, sin construir tuplas de Python.
Almacena los datos temporalmente en archivos Parquet Montaje de parquet Útil para ETL entre sistemas donde se necesita un formato de archivo intermedio. Parquet ya está en formato Arrow, así que se carga sin conversión por filas.

Carga datos CSV con copia masiva

Cargar datos CSV es la pregunta de ingesta más común para el trabajo con bases de datos en Python. Usa csv.reader con un generador 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()

El patrón generador mantiene el uso de memoria constante independientemente del tamaño del archivo. Para el mapeo de columnas y la gestión de identidades, véase Operaciones de copia masiva.

Inserciones de una sola fila

Utiliza inserciones individuales para escrituras en el nivel de aplicación en las que se procesa un registro cada vez. Uso OUTPUT INSERTED para recuperar las claves generadas:

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

Los insertos individuales son la opción adecuada cuando:

  • Insertas una fila por acción del usuario (envío de formularios, llamada a la API).
  • Necesitas validar o transformar cada fila individualmente antes de insertarla.
  • Necesitas el ID insertado u otros valores generados inmediatamente.

Inserciones por lotes

Usa executemany() cuando tengas un número moderado de filas y no necesites el rendimiento de la copia en bloque:

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() envía cada fila como una sentencia parametrizada separada. Cuando el rendimiento importa más que el control por fila, bulkcopy() es más eficiente porque utiliza el protocolo TDS de inserción masiva. El punto de cruce depende del ancho de las filas y de la latencia de la red, pero suele situarse en unos pocos cientos de filas.

Copia masiva

Cuando el rendimiento importa más que el control por fila, utiliza bulkcopy(). Utiliza el protocolo TDS bulk insert, que transmite filas en lugar de enviar una sentencia por fila:

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

Consejos de rendimiento para copia masiva

  • Utiliza generadores para conjuntos de datos grandes para mantener constante el uso de memoria.
  • Uso bulkcopy_arrow() cuando la fuente es columnar, como un DataFrame o un archivo Parquet. Omite la conversión a tuplas de fila de Python.
  • Establezca batch_size para controlar cuántas filas se envían en cada lote TDS. Empieza con 5.000 y ajusta según el ancho de la fila.
  • Usa cerraduras de mesa para cargas exclusivas: cursor.bulkcopy("dbo.ProductImport", rows, table_lock=True).
  • Desactiva los índices antes de cargarlos y luego reconstruye después. Esta secuencia evita la sobrecarga de mantenimiento del índice durante la carga.

Para mapeos de columnas, columnas de identidad, manejo de NULL y carga paralela, véase Operaciones de copia masiva.

Upsert con MERGE

MERGEes la sentencia de Microsoft SQL para condiciones INSERT, UPDATE, y DELETE en una sola operación. Gestiona el patrón de "insertar si es nuevo, actualizar si existe" que los desarrolladores de Python suelen necesitar.

Upsert de una sola fila

Para una sola fila, usa MERGE con una USING cláusula que defina alias 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 masivo con una tabla de preparación

Para los upserts masivos, prepara primero los datos en una tabla temporal y, a continuación, utiliza MERGE para actualizarlos a partir de ella. Utiliza insert-or-update como patrón predeterminado para los upserts de DataFrame y las actualizaciones por lotes:

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 ejemplo demuestra el patrón predeterminado de insertar o actualizar:

  • INSERT filas de la fuente que no existen en el destino (WHEN NOT MATCHED BY TARGET).
  • UPDATE filas que existen en ambos (WHEN MATCHED).
  • La cláusula OUTPUT informa qué acción se realizó en cada fila, lo cual es útil para las auditorías.

Caution

Añade WHEN NOT MATCHED BY SOURCE THEN DELETE solo cuando los datos de preparación sean una instantánea completa y fidedigna del destino. Si el lote solo contiene filas modificadas, esa cláusula elimina filas que se omitieron intencionadamente del feed de origen.

Si necesitas una conciliación completa, amplía el MERGE solo después de confirmar que la fuente es fidedigna para la tabla de destino:

WHEN NOT MATCHED BY SOURCE THEN
    DELETE

En entornos compartidos, utiliza un nombre global único para la tabla temporal en cada ejecución o una tabla permanente de preparación para evitar colisiones entre trabajos concurrentes.

Cuándo usar sentencias separadas UPDATE y INSERT en su lugar

MERGE es potente pero tiene casos límite. Considera usar sentencias separadas cuando:

  • No necesitas DELETE lógica. Un separado UPDATE seguido de INSERT WHERE NOT EXISTS es más legible y sencillo de depurar.
  • La MERGE afirmación es lo suficientemente compleja como para que el comportamiento de bloqueo sea difícil de predecir. Las instrucciones separadas proporcionan un control explícito de la granularidad del bloqueo.
  • Estás actualizando una tabla con alta concurrencia en la que la escalada de bloqueos de MERGE podría provocar bloqueos.
# 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()

Cargar DataFrames

Un DataFrame tiene estructura columnar, así que cárgalo con bulkcopy_arrow() en lugar de aplanarlo en tuplas de filas para bulkcopy().

Asocia los tipos de Arrow con las columnas de destino antes de realizar la carga. pyarrow infiere float64 para una columna numérica, que el conductor no puede asignar a dinero, decimal o 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()

Usa Table.cast() en lugar de pasar el esquema a Table.from_pandas(), que no puede convertir una columna flotante a decimal128 directamente.

Polars

Polars implementa la interfaz de datos C de Arrow, por lo que puedes pasar el propio DataFrame. Lanza las columnas primero por la misma razón:

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

También puedes establecer los tipos al leer el archivo, con pl.read_csv("products.csv", schema_overrides={"ListPrice": pl.Decimal(19, 4)}).

Al pasar el DataFrame, sus búferes se transfieren directamente al controlador sin realizar una copia. df.to_arrow() también funciona, pero Polars vuelve a codificar columnas de cadena durante esa conversión, lo que copia todos los datos de la cadena.

bulkcopy_arrow() acepta un pyarrow.Table, un RecordBatch, un RecordBatchReader o cualquier objeto que implemente la interfaz de datos C de Arrow a través de __arrow_c_stream__ o __arrow_c_array__. Pasar cualquiera de ellos a bulkcopy() genera TypeError.

Para consultar todos los patrones de carga de DataFrames, consulte la integración de pandas y la integración de Polars.

Cargar datos de Arrow

Cuando los datos de origen ya están en formato Apache Arrow, cursor.bulkcopy_arrow() los carga sin necesidad de crear primero tuplas de 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")

El método también acepta un pyarrow.RecordBatch o un pyarrow.RecordBatchReader, por lo que puedes transferir en flujo un conjunto de resultados desde cursor.arrow_reader() directamente a otra tabla.

Cada tipo de columna Arrow debe ser compatible con su tipo de columna SQL de destino, y el escritor no convierte entre familias de tipos. Para más información, véase integración con Apache Arrow.

Montaje de parquet

Utiliza Parquet como formato intermedio al migrar datos entre sistemas o cuando tu pipeline ETL ya produce archivos Parquet. Un archivo Parquet lee datos de Arrow, así que pásalos directamente a bulkcopy_arrow():

import pyarrow.parquet as pq

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

Para archivos grandes de Parquet, iterar grupos de filas para mantener el uso de memoria constante. Cada lote es un RecordBatch, que bulkcopy_arrow() acepta directamente:

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 todo el archivo en una sola llamada, envuelva los lotes en 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()

Validar datos cargados

Después de la carga, verifica el número de filas y realiza comprobaciones puntuales de los datos:

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 en producción, no confíes en que la transacción de la conexión que realiza la llamada proteja una llamada a bulkcopy(). bulkcopy() abre su propia conexión interna y confirma de forma independiente las filas copiadas, así que un conn.rollback() en tu conexión principal no puede deshacerlas. Dos enfoques te dan atomicidad:

  • Configura use_internal_transaction=True para envolver cada lote en su propia transacción. Un lote que falla a mitad de camino se revierte en su totalidad, en lugar de quedar a medio cargar.
  • Para validar los datos antes de promocionarlos, cópialos de forma masiva a una tabla de paso, valídalos y, a continuación, traslada las filas a la tabla de destino utilizando un INSERT ... SELECT dentro de una transacción en tu conexión principal. Dado que eso INSERT se ejecuta en tu conexión, conn.rollback() lo deshace si la validación falla.
# 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