Usa mssql-python con pandas

La biblioteca pandas es la principal herramienta de análisis de datos de Python. Combinando pandas con el controlador mssql-python, puedes:

  • Carga los resultados de las consultas SQL directamente en DataFrames.
  • Escribe DataFrames de nuevo en Microsoft SQL de forma eficiente.
  • Realizar operaciones ETL.
  • Crear canalizaciones de datos.

Los ejemplos de este artículo consultan la Production.Product tabla y otras tablas en la base de datos de ejemplo de AdventureWorks. Los ejemplos que escriben datos utilizan tablas temporales para evitar modificar los datos de muestra.

Otras tablas referenciadas en ejemplos de análisis (Sales.SalesOrderHeader, Sales.SalesOrderDetail, Production.ProductSubcategory) forman parte de AdventureWorks. Sustituye tus propias tablas al adaptar estos patrones.

Lee datos en DataFrames

El controlador mssql-python devuelve las filas como objetos de Python, que conviertes en DataFrames de pandas leyendo los nombres de las columnas de cursor.description y los valores de las filas de fetchall(). Las funciones auxiliares de esta sección envuelven esa conversión en patrones reutilizables.

Consulta básica a DataFrame

Esta función ejecuta una consulta parametrizada y construye un DataFrame a partir del conjunto completo de resultados. Funciona bien para conjuntos de resultados que encajan cómodamente en la memoria.

import pandas as pd
import mssql_python

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

def query_to_dataframe(cursor, query: str, params: dict = None) -> pd.DataFrame:
    """Execute query and return results as DataFrame."""
    cursor.execute(query, params or {})
    
    # cursor.description is a list of tuples, one per column.
    # Each tuple's first element is the column name.
    columns = [col[0] for col in cursor.description]
    
    # Fetch all rows
    rows = cursor.fetchall()
    
    # Convert to DataFrame
    data = [tuple(row) for row in rows]
    return pd.DataFrame(data, columns=columns)

# Usage: %(cat)s is a parameterized placeholder. The driver safely substitutes
# the value from the dict, which prevents SQL injection.
df = query_to_dataframe(cursor, "SELECT * FROM Production.Product WHERE ProductSubcategoryID = %(cat)s", {"cat": 5})
print(df.head())

Note

Si su cadena de conexión usa Authentication=ActiveDirectoryDefault, el controlador usa DefaultAzureCredential, que intenta usar varios proveedores de credenciales en secuencia. La primera conexión puede ser lenta porque el SDK recorre la cadena hasta encontrar un proveedor que funcione. En producción, si sabes qué tipo de credencial utiliza tu entorno, especifícala directamente (por ejemplo, ActiveDirectoryMSI para identidad gestionada) para evitar el recorrido en cadena. Para más información, consulte Autenticación de Microsoft Entra.

Transmitir grandes conjuntos de datos

Para tablas con millones de filas, cargar todo a la vez puede agotar la memoria. El enfoque por fragmentos obtiene filas en lotes con fetchmany() y concatena los resultados, manteniendo el uso máximo de memoria proporcional a chunksize en lugar de al conjunto completo de resultados.

def query_to_dataframe_chunked(cursor, query: str, params: dict = None, 
                                chunksize: int = 10000) -> pd.DataFrame:
    """Load large query results in chunks for memory efficiency."""
    cursor.execute(query, params or {})
    columns = [col[0] for col in cursor.description]
    
    chunks = []
    while True:
        rows = cursor.fetchmany(chunksize)
        if not rows:
            break
        data = [tuple(row) for row in rows]
        chunks.append(pd.DataFrame(data, columns=columns))
    
    return pd.concat(chunks, ignore_index=True) if chunks else pd.DataFrame(columns=columns)

# Usage for large tables
df = query_to_dataframe_chunked(cursor, "SELECT * FROM Production.TransactionHistory", chunksize=50000)

Generador para grandes conjuntos de datos

Cuando necesites procesar datos de forma incremental sin tener el resultado completo en memoria, usa un generador. Cada uno yield produce un bloque de DataFrame que puedes procesar y descartar antes de recuperar el siguiente.

def query_to_dataframe_generator(cursor, query: str, params: dict = None,
                                  chunksize: int = 10000):
    """Yield DataFrame chunks for processing without loading all data."""
    cursor.execute(query, params or {})
    columns = [col[0] for col in cursor.description]
    
    while True:
        rows = cursor.fetchmany(chunksize)
        if not rows:
            break
        data = [tuple(row) for row in rows]
        yield pd.DataFrame(data, columns=columns)

# Process chunks without loading entire dataset
huge_query = """
    SELECT * FROM Production.TransactionHistory
    UNION ALL SELECT * FROM Production.TransactionHistory
    UNION ALL SELECT * FROM Production.TransactionHistory
"""
for chunk_df in query_to_dataframe_generator(cursor, huge_query):
    # Process each chunk, then discard it before the next fetch
    print(f"Processing chunk of {len(chunk_df)} rows")
    total_cost = chunk_df["ActualCost"].sum()
    print(f"Chunk total cost: {total_cost}")

Escribe DataFrames en Microsoft SQL

Entrecomille los identificadores para evitar inyección SQL

Los nombres de tablas y columnas no pueden pasarse como parámetros de consulta en SQL. Cuando construyas sentencias SQL con identificadores dinámicos, envuelve cada nombre entre corchetes y evita cualquier carácter incrustado ] para evitar la inyección SQL.

def quote_id(identifier: str) -> str:
    """Quote a Microsoft SQL identifier to prevent SQL injection.
    
    Wraps the name in square brackets and escapes any embedded ] characters.
    Raises ValueError if the identifier is empty or contains null bytes.
    """
    if not identifier or "\x00" in identifier:
        raise ValueError(f"Invalid identifier: {identifier!r}")
    escaped = identifier.replace("]", "]]")
    return f"[{escaped}]"

Las funciones auxiliares de esta sección utilizan quote_id() para todos los nombres de tablas y columnas en el SQL generado.

Insertar filas de DataFrame

El enfoque más sencillo itera sobre las filas del DataFrame y genera un INSERT por fila. El enfoque sencillo funciona para DataFrames pequeños, pero es lento para volúmenes grandes porque cada fila requiere un viaje de ida y vuelta separado al servidor.

def dataframe_to_sql(cursor, conn, df: pd.DataFrame, table: str, 
                     if_exists: str = "append") -> int:
    """Write DataFrame to Microsoft SQL table."""
    if if_exists == "replace":
        cursor.execute(f"TRUNCATE TABLE {quote_id(table)}")
    
    columns = df.columns.tolist()
    placeholders = ", ".join([f"%({col})s" for col in columns])
    col_list = ", ".join([quote_id(col) for col in columns])
    
    query = f"INSERT INTO {quote_id(table)} ({col_list}) VALUES ({placeholders})"
    
    rows_inserted = 0
    for _, row in df.iterrows():
        params = {col: (None if pd.isna(val) else val) for col, val in row.items()}
        cursor.execute(query, params)
        rows_inserted += 1
    
    conn.commit()
    return rows_inserted

# Usage
cursor.execute("""
    CREATE TABLE #Products (
        Name NVARCHAR(100),
        ListPrice DECIMAL(10,2),
        ProductSubcategoryID INT
    )
""")
df = pd.DataFrame({
    "Name": ["Product A", "Product B"],
    "ListPrice": [29.99, 49.99],
    "ProductSubcategoryID": [1, 2]
})
rows = dataframe_to_sql(cursor, conn, df, "#Products")
print(f"Inserted {rows} rows")

Para DataFrames grandes, utiliza el método bulkcopy() del controlador, que envía filas de forma masiva a través del protocolo TDS (Tabular Data Stream), el protocolo nativo de comunicación que usa Microsoft SQL. Este enfoque es más rápido que los insertos fila por fila porque minimiza los viajes de ida y vuelta.

def dataframe_to_sql_bulk(conn, df: pd.DataFrame, table: str) -> int:
    """Bulk insert DataFrame using BCP for better performance."""
    # Convert DataFrame to list of tuples, handling NaN
    rows = []
    for _, row in df.iterrows():
        row_data = tuple(None if pd.isna(v) else v for v in row)
        rows.append(row_data)
    
    cursor = conn.cursor()
    result = cursor.bulkcopy(table, rows)
    conn.commit()
    return result["rows_copied"]

# Usage
cursor.execute("CREATE TABLE ##PandasProducts (Name NVARCHAR(50), ListPrice DECIMAL(10,2), ProductSubcategoryID INT)")
conn.commit()

df = pd.DataFrame({
    "Name": ["Product A", "Product B", "Product C"],
    "ListPrice": [29.99, 49.99, 19.99],
    "ProductSubcategoryID": [1, 2, 1]
})

rows = dataframe_to_sql_bulk(conn, df, "##PandasProducts")

Actualizar filas existentes desde DataFrame

Para actualizar las filas que ya existen en la tabla, itere sobre el DataFrame y ejecute instrucciones parametrizadas UPDATE. El key_column indica qué fila se debe actualizar.

def update_from_dataframe(cursor, conn, df: pd.DataFrame, table: str,
                          key_column: str) -> int:
    """Update existing rows based on key column."""
    columns = [col for col in df.columns if col != key_column]
    set_clause = ", ".join([f"{quote_id(col)} = %({col})s" for col in columns])
    
    query = f"UPDATE {quote_id(table)} SET {set_clause} WHERE {quote_id(key_column)} = %({key_column})s"
    
    rows_updated = 0
    for _, row in df.iterrows():
        params = {col: (None if pd.isna(val) else val) for col, val in row.items()}
        cursor.execute(query, params)
        rows_updated += cursor.rowcount
    
    conn.commit()
    return rows_updated

# Usage
cursor.execute("""
    CREATE TABLE #ProductPrices (
        ProductID INT PRIMARY KEY,
        ListPrice DECIMAL(10,2)
    );
    INSERT INTO #ProductPrices VALUES (1, 29.99), (2, 49.99), (3, 19.99);
""")
conn.commit()

df_updates = pd.DataFrame({
    "ProductID": [1, 2, 3],
    "ListPrice": [31.99, 52.99, 21.99]
})
updated = update_from_dataframe(cursor, conn, df_updates, "#ProductPrices", "ProductID")

Patrón de upsert (fusión)

Cuando algunas filas puedan ser nuevas y otras ya existan, usa una instrucción SQL MERGE para insertar o actualizar en una sola operación. MERGE compara cada fila entrante con la tabla objetivo usando las columnas clave. Si se encuentra una coincidencia, se actualiza; de lo contrario, se inserta. MERGE evita verificar su existencia por separado.

def upsert_from_dataframe(cursor, conn, df: pd.DataFrame, table: str,
                          key_columns: list[str]) -> int:
    """Insert or update rows based on key columns. Returns total rows affected."""
    all_columns = df.columns.tolist()
    value_columns = [c for c in all_columns if c not in key_columns]
    
    total_affected = 0
    
    for _, row in df.iterrows():
        params = {col: (None if pd.isna(val) else val) for col, val in row.items()}
        
        # Build MERGE statement with quoted identifiers
        key_match = " AND ".join([f"t.{quote_id(k)} = s.{quote_id(k)}" for k in key_columns])
        update_set = ", ".join([f"{quote_id(c)} = s.{quote_id(c)}" for c in value_columns])
        all_cols = ", ".join([quote_id(c) for c in all_columns])
        all_vals = ", ".join([f"%({c})s" for c in all_columns])
        
        cursor.execute(f"""
            MERGE {quote_id(table)} AS t
            USING (SELECT {', '.join([f'%({c})s AS {quote_id(c)}' for c in all_columns])}) AS s
            ON {key_match}
            WHEN MATCHED THEN UPDATE SET {update_set}
            WHEN NOT MATCHED THEN INSERT ({all_cols}) VALUES ({all_vals});
        """, params)
        
        total_affected += cursor.rowcount
    
    conn.commit()
    return total_affected

Patrones de análisis de datos

Los siguientes ejemplos muestran tareas de análisis comunes que combinan consultas de Microsoft SQL con transformaciones pandas.

Agregar consultas a DataFrame

def get_sales_summary(cursor) -> pd.DataFrame:
    """Get sales summary by category."""
    return query_to_dataframe(cursor, """
        SELECT 
            pc.Name AS CategoryName,
            COUNT(*) AS ProductCount,
            AVG(p.ListPrice) AS AvgPrice,
            MIN(p.ListPrice) AS MinPrice,
            MAX(p.ListPrice) AS MaxPrice
        FROM Production.Product p
        JOIN Production.ProductSubcategory pc ON p.ProductSubcategoryID = pc.ProductSubcategoryID
        GROUP BY pc.Name
        ORDER BY ProductCount DESC
    """)

df = get_sales_summary(cursor)
print(df.to_string())

Datos de serie temporal

Utiliza la indexación de fechas y el remuestreo de Pandas para trabajar con datos de series temporales de Microsoft SQL. Para permitir operaciones como promedios móviles y remuestreo, establece la columna de fecha como índice DataFrame.

def get_daily_sales(cursor, start_date: str, end_date: str) -> pd.DataFrame:
    """Get daily sales time series."""
    df = query_to_dataframe(cursor, """
        SELECT 
            CAST(OrderDate AS DATE) AS Date,
            COUNT(*) AS OrderCount,
            SUM(TotalDue) AS Revenue
        FROM Sales.SalesOrderHeader
        WHERE OrderDate BETWEEN %(start)s AND %(end)s
        GROUP BY CAST(OrderDate AS DATE)
        ORDER BY Date
    """, {"start": start_date, "end": end_date})
    
    # Set date as index for time series operations
    df["Date"] = pd.to_datetime(df["Date"])
    df.set_index("Date", inplace=True)
    
    return df

# Usage
sales_df = get_daily_sales(cursor, "2024-01-01", "2024-12-31")

# Resample to weekly
weekly = sales_df.resample("W").sum()

# Calculate rolling average
sales_df["RollingAvg"] = sales_df["Revenue"].rolling(window=7).mean()

Tablas dinámicas a partir de datos SQL

Las tablas dinámicas reorganizan los datos de filas en una matriz. Para reorganizarlo por dimensiones como año, mes y categoría, extrae los datos en bruto de Microsoft SQL y luego usa pivot_table().

def get_sales_pivot(cursor) -> pd.DataFrame:
    """Get sales data and create pivot table."""
    df = query_to_dataframe(cursor, """
        SELECT 
            YEAR(soh.OrderDate) AS Year,
            MONTH(soh.OrderDate) AS Month,
            pc.Name AS CategoryName,
            SUM(sod.OrderQty * sod.UnitPrice) AS Revenue
        FROM Sales.SalesOrderHeader soh
        JOIN Sales.SalesOrderDetail sod ON soh.SalesOrderID = sod.SalesOrderID
        JOIN Production.Product p ON sod.ProductID = p.ProductID
        JOIN Production.ProductSubcategory pc ON p.ProductSubcategoryID = pc.ProductSubcategoryID
        GROUP BY YEAR(soh.OrderDate), MONTH(soh.OrderDate), pc.Name
    """)
    
    # Create pivot table
    pivot = df.pivot_table(
        values="Revenue",
        index=["Year", "Month"],
        columns="CategoryName",
        aggfunc="sum",
        fill_value=0
    )
    
    return pivot

pivot_df = get_sales_pivot(cursor)
print(pivot_df)

Patrones ETL

Para construir pipelines de extracción, transformación y carga, combina consultas SQL de Microsoft con transformaciones pandas para construir pipelines de extracción, transformación y carga. El conductor se encarga de la extracción y la carga mientras que los pandas se encargan del paso de transformación.

Extracción, transformación y carga

Este ejemplo extrae datos activos de clientes, aplica reglas de negocio a segmentar clientes y carga los resultados en una tabla de destinos.

def etl_pipeline(source_cursor, dest_cursor, dest_conn):
    """Simple ETL pipeline with pandas."""
    
    # Extract
    df = query_to_dataframe(source_cursor, """
        SELECT 
            c.CustomerID,
            COUNT(soh.SalesOrderID) AS OrderCount,
            SUM(soh.TotalDue) AS TotalSpent
        FROM Sales.Customer c
        JOIN Sales.SalesOrderHeader soh ON c.CustomerID = soh.CustomerID
        WHERE soh.OrderDate > DATEADD(YEAR, -1, GETDATE())
        GROUP BY c.CustomerID
    """)
    
    # Transform
    df["CustomerSegment"] = pd.cut(
        df["TotalSpent"],
        bins=[0, 100, 500, 1000, float("inf")],
        labels=["Bronze", "Silver", "Gold", "Platinum"]
    )
    df["AvgOrderValue"] = df["TotalSpent"] / df["OrderCount"].replace(0, 1)
    df["IsHighValue"] = df["TotalSpent"] > 500
    
    # Load
    dataframe_to_sql_bulk(dest_conn, df[["CustomerID", "CustomerSegment", "AvgOrderValue", "IsHighValue"]], 
                          "#CustomerAnalytics")
    
    return len(df)

Patrón incremental de carga

Para canalizaciones de datos en curso, solo carga los registros que cambiaron desde la última ejecución. Este enfoque consulta la tabla de destino para obtener la marca de tiempo máxima y luego recupera solo los registros más recientes de la fuente.

def incremental_load(cursor, conn, source_table: str, dest_table: str,
                     timestamp_col: str) -> int:
    """Load only new/changed records based on timestamp."""
    
    # Get last loaded timestamp
    cursor.execute(f"SELECT MAX({quote_id(timestamp_col)}) FROM {quote_id(dest_table)}")
    last_loaded = cursor.fetchval()
    
    # Build query for new records
    if last_loaded:
        df = query_to_dataframe(cursor, f"""
            SELECT * FROM {quote_id(source_table)}
            WHERE {quote_id(timestamp_col)} > %(last)s
        """, {"last": last_loaded})
    else:
        df = query_to_dataframe(cursor, f"SELECT * FROM {quote_id(source_table)}")
    
    if df.empty:
        return 0
    
    # Load new records
    return dataframe_to_sql_bulk(conn, df, dest_table)

Consejos de rendimiento

Uso del tipo de datos adecuado

Pandas por defecto utiliza tipos de 64 bits para números, desperdiciando memoria cuando los tipos más pequeños son suficientes. Reducir enteros y flotadores, y convertir columnas de cadenas de baja cardinalidad en categorías puede reducir significativamente el uso de memoria.

def optimize_dataframe_types(df: pd.DataFrame) -> pd.DataFrame:
    """Optimize DataFrame memory usage."""
    for col in df.columns:
        col_type = df[col].dtype
        
        if col_type == "int64":
            # Downcast integers
            df[col] = pd.to_numeric(df[col], downcast="integer")
        elif col_type == "float64":
            # Downcast floats
            df[col] = pd.to_numeric(df[col], downcast="float")
        elif col_type == "object":
            # Convert to category if low cardinality
            num_unique = df[col].nunique()
            if num_unique / len(df) < 0.5:
                df[col] = df[col].astype("category")
    
    return df

Usa SQL para el trabajo pesado

Microsoft SQL es más rápido para agregaciones, filtrado y uniones que extraer todos tus datos en bruto por cable y procesarlos localmente en Python. Deja que Microsoft SQL haga el trabajo duro siempre que sea posible, mueve solo los datos que necesitas por la red y usa pandas para análisis y transformaciones que sean más cómodos en Python.

# Avoid: pulling all rows over the wire to aggregate locally in pandas
df_all = query_to_dataframe(cursor, "SELECT * FROM Production.Product")  # transfers entire table
summary = df_all.groupby("Color").agg({"ListPrice": "sum"})  # aggregation that SQL can do faster

# Better: push the aggregation into SQL and transfer only the summary
df = query_to_dataframe(cursor, """
    SELECT Color, SUM(ListPrice) AS TotalPrice
    FROM Production.Product
    WHERE Color IS NOT NULL
    GROUP BY Color
""")

Escribe por lotes

Para los DataFrames que sean demasiado grandes para una sola inserción masiva, divide el trabajo en lotes y realiza un seguimiento del progreso.

def batch_insert(cursor, conn, df: pd.DataFrame, table: str, batch_size: int = 1000):
    """Insert in batches with progress tracking."""
    total = len(df)
    
    for i in range(0, total, batch_size):
        batch = df.iloc[i:i + batch_size]
        dataframe_to_sql(cursor, conn, batch, table)
        print(f"Inserted {min(i + batch_size, total)}/{total}")