Usar mssql-python con Polars

Polars es una biblioteca para DataFrames de alto rendimiento escrita en Rust que ofrece una alternativa rápida y eficiente en memoria a pandas. Polars, junto con el controlador mssql-python, te permite:

  • Carga los resultados de consultas SQL directamente en Polars DataFrames.
  • Usa Apache Arrow para transferir datos sin copia desde Microsoft SQL.
  • Vuelve a escribir los DataFrames de Polars en Microsoft SQL eficientemente.
  • Crea canalizaciones de datos de alto rendimiento con evaluación diferida.

Los ejemplos de este artículo consultan la AdventureWorks base de datos de ejemplo. Si aún no lo tienes, consulta las bases de datos de ejemplo de AdventureWorks.

Lee datos en Polars DataFrames

Puedes cargar datos SQL de Microsoft en Polar de dos maneras: conversión fila por fila mediante métodos estándar de cursor, o transferencia sin copia a través de Apache Arrow. "Copia cero" significa que los datos permanecen en un único búfer de memoria que el controlador, Arrow y Polars leen directamente, por lo que ninguna fila se duplica en objetos Python intermedios. Utiliza el enfoque Arrow para la mayoría de las cargas de trabajo debido a esta eficiencia.

Consulta básica a DataFrame

Este enfoque obtiene todas las filas con el cursor estándar y construye manualmente un DataFrame Polars. Acepta consultas parametrizadas para sustitución segura de valores. Funciona sin PyArrow pero es más lento para conjuntos de resultados grandes porque cada valor pasa por Python.

import polars as pl
import mssql_python

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


def query_to_polars(cursor, query: str, params: dict = None) -> pl.DataFrame:
    """Execute query and return results as Polars DataFrame."""
    cursor.execute(query, params or {})

    columns = [col[0] for col in cursor.description]
    rows = cursor.fetchall()

    data = {col: [row[i] for row in rows] for i, col in enumerate(columns)}
    return pl.DataFrame(data)

# Usage: %(color)s is a parameterized placeholder. The driver safely substitutes
# the value from the dict, which prevents SQL injection.
df = query_to_polars(cursor, "SELECT TOP 5 Name, ListPrice FROM Production.Product WHERE Color = %(color)s", {"color": "Black"})
print(df)

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.

La forma más eficiente de cargar datos SQL de Microsoft en Polar es a través de Apache Arrow. El método arrow() del controlador mssql-python devuelve un pyarrow.Table que Polars puede consumir sin sobrecarga de copia.

def query_to_polars_arrow(cursor, query: str, params: dict = None) -> pl.DataFrame:
    """Execute query and load results through Arrow for best performance."""
    cursor.execute(query, params or {})
    arrow_table = cursor.arrow()
    return pl.from_arrow(arrow_table)

# Usage
df = query_to_polars_arrow(cursor, "SELECT ProductID, Name, ListPrice FROM Production.Product")
print(df)

Transmita grandes conjuntos de datos mediante lotes de Arrow

Para conjuntos de datos que no caben en memoria, úsalo arrow_reader() para procesar datos en lotes de streaming. Cada lote es un pyarrow.RecordBatch que Polars puede consumir de forma independiente, por lo que el uso de memoria se mantiene proporcional a batch_size en lugar de al conjunto completo de resultados.

def process_large_query(cursor, query: str, params: dict = None, batch_size: int = 50000) -> pl.DataFrame:
    """Process large query results as streaming Arrow batches."""
    cursor.execute(query, params or {})
    reader = cursor.arrow_reader(batch_size=batch_size)

    results = []
    for batch in reader:
        chunk_df = pl.from_arrow(batch)
        # Process each chunk
        results.append(chunk_df)

    return pl.concat(results) if results else pl.DataFrame()

# Usage
df = process_large_query(cursor, "SELECT * FROM Production.TransactionHistory")

Usa LazyFrames para la ejecución diferida

Polars LazyFrames te permite construir una cadena de operaciones (filtro, grupo, ordenación) sin ejecutarlas inmediatamente. Polars optimiza toda la cadena antes de ejecutarla, lo que puede ser más rápido que aplicar cada paso individualmente.

def query_to_lazy(cursor, query: str, params: dict = None) -> pl.LazyFrame:
    """Execute query and return a Polars LazyFrame."""
    cursor.execute(query, params or {})
    arrow_table = cursor.arrow()
    return pl.from_arrow(arrow_table).lazy()

# Build a query plan without executing immediately
lf = query_to_lazy(cursor, "SELECT SalesOrderID, CustomerID, TotalDue, OrderDate FROM Sales.SalesOrderHeader")
result = (
    lf.filter(pl.col("TotalDue") > 100)
    .group_by("CustomerID")
    .agg([
        pl.col("TotalDue").sum().alias("TotalSpent"),
        pl.col("SalesOrderID").count().alias("OrderCount")
    ])
    .sort("TotalSpent", descending=True)
    .collect()  # Execute the optimized plan
)
print(result)

Escribe DataFrames Polars 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 fila por fila itera sobre el DataFrame con iter_rows(named=True) y ejecuta uno INSERT por fila. Este enfoque es sencillo pero lento para volúmenes grandes porque cada fila requiere un viaje de ida y vuelta al servidor.

def polars_to_sql(cursor, conn, df: pl.DataFrame, table: str) -> int:
    """Write Polars DataFrame to Microsoft SQL table."""
    columns = df.columns
    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.iter_rows(named=True):
        params = {k: (None if v is None else v) for k, v in row.items()}
        cursor.execute(query, params)
        rows_inserted += 1

    conn.commit()
    return rows_inserted

# Usage
cursor.execute("CREATE TABLE #PolarsInsert (Name NVARCHAR(100), Price DECIMAL(10,2), CategoryID INT)")
df = pl.DataFrame({
    "Name": ["Product A", "Product B"],
    "Price": [29.99, 49.99],
    "CategoryID": [1, 2]
})
rows = polars_to_sql(cursor, conn, df, "#PolarsInsert")
print(f"Inserted {rows} rows")

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

def polars_to_sql_bulk(conn, df: pl.DataFrame, table: str) -> int:
    """Bulk insert Polars DataFrame using BCP for best performance."""
    rows = [tuple(None if v is None else v for v in row) for row in df.iter_rows()]

    cursor = conn.cursor()
    result = cursor.bulkcopy(table, rows)
    conn.commit()
    return result["rows_copied"]

# Usage
cursor.execute("CREATE TABLE ##PolarsBulk (Name NVARCHAR(50), Price FLOAT, CategoryID INT)")
conn.commit()
df = pl.DataFrame({
    "Name": ["Product A", "Product B", "Product C"],
    "Price": [29.99, 49.99, 19.99],
    "CategoryID": [1, 2, 1]
})
rows = polars_to_sql_bulk(conn, df, "##PolarsBulk")
print(f"Bulk inserted {rows} rows")

Patrones de análisis de datos

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

Consultas agregadas

Este ejemplo agrupa productos por subcategoría y calcula estadísticas de conteo y precios en SQL, luego carga el resumen en un Polars DataFrame:

def get_sales_summary(cursor) -> pl.DataFrame:
    """Get sales summary by subcategory."""
    cursor.execute("""
        SELECT
            sc.Name AS SubcategoryName,
            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 sc ON p.ProductSubcategoryID = sc.ProductSubcategoryID
        GROUP BY sc.Name
        ORDER BY ProductCount DESC
    """)
    return pl.from_arrow(cursor.arrow())

df = get_sales_summary(cursor)
print(df)

Análisis de series temporales

Carga datos de series temporales desde Microsoft SQL y añade columnas calculadas, como medias móviles, mediante expresiones de Polars.

def get_daily_sales(cursor, start_date: str, end_date: str) -> pl.DataFrame:
    """Get daily sales and compute rolling statistics."""
    cursor.execute("""
        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})

    df = pl.from_arrow(cursor.arrow())

    # Add rolling 7-day average
    df = df.with_columns(
        pl.col("Revenue").rolling_mean(window_size=7).alias("RollingAvg")
    )
    return df

sales_df = get_daily_sales(cursor, "2013-01-01", "2013-12-31")
print(sales_df)

Unir datos SQL con archivos locales

Puedes enriquecer los datos SQL de Microsoft uniéndolos con archivos CSV locales en Polars. Carga cada fuente en un DataFrame y únete a la memoria.

# Load SQL data via Arrow
cursor.execute("SELECT c.CustomerID, p.FirstName, p.LastName FROM Sales.Customer c JOIN Person.Person p ON c.PersonID = p.BusinessEntityID")
customers = pl.from_arrow(cursor.arrow())

# Load local CSV
orders = pl.read_csv("orders_export.csv")

# Join in Polars
result = customers.join(orders, on="CustomerID", how="inner")
print(result)

Patrones ETL

Crea canalizaciones de extracción, transformación y carga (ETL) combinando consultas SQL de Microsoft con transformaciones con Polars. Las expresiones de Polars se encargan del paso de transformación y bulkcopy() se encarga de la carga.

Extracción, transformación y carga

Este ejemplo extrae datos de clientes activos mediante Arrow, aplica lógica de segmentación empresarial con expresiones de Polars y carga los resultados mediante copia en bloque.

def etl_pipeline(source_cursor, dest_conn):
    """ETL pipeline using Polars transformations."""

    # Extract: derive a per-customer summary from order history via Arrow
    source_cursor.execute("""
        SELECT
            CustomerID,
            COUNT(*) AS OrderCount,
            SUM(TotalDue) AS TotalSpent
        FROM Sales.SalesOrderHeader
        WHERE OrderDate > DATEADD(YEAR, -1, (SELECT MAX(OrderDate) FROM Sales.SalesOrderHeader))
        GROUP BY CustomerID
    """)
    df = pl.from_arrow(source_cursor.arrow())
    df = df.with_columns(pl.col("TotalSpent").cast(pl.Float64))

    # Transform with Polars expressions
    df = df.with_columns([
        pl.when(pl.col("TotalSpent") > 1000).then(pl.lit("Platinum"))
          .when(pl.col("TotalSpent") > 500).then(pl.lit("Gold"))
          .when(pl.col("TotalSpent") > 100).then(pl.lit("Silver"))
          .otherwise(pl.lit("Bronze"))
          .alias("CustomerSegment"),
        (pl.col("TotalSpent") / pl.col("OrderCount").clip(lower_bound=1))
          .alias("AvgOrderValue"),
        (pl.col("TotalSpent") > 500).alias("IsHighValue")
    ])

    # Load via bulk copy into the destination table
    dest_cursor = dest_conn.cursor()
    dest_cursor.execute("""
        CREATE TABLE ##CustomerAnalytics (
            CustomerID INT,
            CustomerSegment NVARCHAR(20),
            AvgOrderValue FLOAT,
            IsHighValue BIT
        )
    """)
    dest_conn.commit()

    load_df = df.select(["CustomerID", "CustomerSegment", "AvgOrderValue", "IsHighValue"])
    polars_to_sql_bulk(dest_conn, load_df, "##CustomerAnalytics")

    return len(df)

Consejos de rendimiento

Los siguientes consejos te ayudan a sacar el máximo partido a la combinación mssql-python y Polars.

Deja que Microsoft SQL se encargue del 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 Polars para análisis y transformaciones que sean más cómodos en Python.

# Avoid: pulling all rows over the wire to aggregate locally in Polars
df = query_to_polars_arrow(cursor, "SELECT * FROM Production.Product WHERE Color IS NOT NULL")  # transfers entire table
summary = df.group_by("Color").agg(pl.col("ListPrice").sum())  # aggregation that SQL can do faster

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

Usa Arrow para todas las operaciones de lectura

La transferencia basada en flechas evita crear objetos Python intermedios, lo que reduce el consumo de memoria y mejora el rendimiento. Prefiere cursor.arrow() a la conversión manual fila por fila para cualquier conjunto de resultados de más de unas pocas filas.

# Suboptimal: Row-by-row conversion
cursor.execute("SELECT * FROM Production.TransactionHistory")
columns = [col[0] for col in cursor.description]
rows = cursor.fetchall()
df = pl.DataFrame({col: [row[i] for row in rows] for i, col in enumerate(columns)})

# Better: Arrow-based transfer
cursor.execute("SELECT * FROM Production.TransactionHistory")
df = pl.from_arrow(cursor.arrow())