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 Insertos en lote 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 Extrae las filas de pandas o Polars y pásalas a bulkcopy().
Almacenar datos temporalmente mediante archivos Parquet Montaje de parquet Útil para ETL entre sistemas donde se necesita un formato de archivo intermedio.

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 insertos individuales para escrituras a nivel de aplicación, donde procesas un registro a la 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 inflexión depende de la anchura de la fila y de la latencia de la red, pero normalmente se sitúa en torno a unos pocos cientos de filas.

Copia masiva

Cuando el rendimiento importa más que el control por fila, utiliza bulkcopy(). Utiliza el protocolo de inserción masiva TDS, que es significativamente más eficiente que las inserciones fila a 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.
  • 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.

Inserción o actualización 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 ensayo

Para operaciones upsert en bloque, carga primero los datos en una tabla temporal y luego usa MERGE para actualizar a partir de ella. Utiliza insert-or-update como patrón predeterminado para los upserts 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 staging sean una instantánea completa y autorizada del objetivo. 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 MERGE solo después de confirmar que la fuente es la fuente autorizada 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 de alta concurrencia donde MERGE la escalada de bloqueos podría causar 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

Extrae filas de un DataFrame de pandas o de Polars y cárgalas con bulkcopy():

pandas

Convierte un DataFrame de pandas en tuplas y pásalo 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

Convierte un DataFrame Polars en tuplas usando el 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 consultar todos los patrones de carga de DataFrames, consulte la integración de pandas y la integración de Polars.

Montaje de parquet

Utiliza Parquet como formato intermedio al migrar datos entre sistemas o cuando tu pipeline ETL ya produce archivos 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 archivos Parquet de gran tamaño, leer por grupos de filas para mantener constante el uso de 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()

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 de trabajo de producción, no confíes en la transacción de la conexión invocadora para proteger 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 revierte ese lote en lugar de dejarlo medio cargado.
  • Para validar los datos antes de pasarlos, cópialos en bloque a una tabla de preparación, valídalos y, a continuación, mueve las filas a la tabla de destino mediante un INSERT ... SELECT dentro de una transacción en la 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