Use mssql-python com o Polars

Polars é uma biblioteca de DataFrames de alto desempenho, escrita em Rust, que oferece uma alternativa rápida e com uso eficiente de memória ao pandas. O Polars, combinado com o driver mssql-python, permite que você:

  • Carregue os resultados da consulta SQL diretamente nos DataFrames do Polars.
  • Use o Apache Arrow para transferência de dados sem cópia zero do Microsoft SQL.
  • Escreva os DataFrames Polars de volta para o Microsoft SQL de forma eficiente.
  • Construa pipelines de dados de alto desempenho com avaliação preguiçosa.

Os exemplos neste artigo consultam o AdventureWorks banco de dados de exemplo. Se você ainda não tem, veja bancos de dados de exemplo do AdventureWorks.

Leia dados em DataFrames Polars

Você pode carregar dados do Microsoft SQL no Polars de duas maneiras: conversão linha a linha por meio dos métodos padrão de cursor ou transferência sem cópia por meio do Apache Arrow. "Zero-copy" significa que os dados permanecem em um único buffer de memória que o driver, Arrow e Polars leem diretamente, para que nenhuma linha seja duplicada em objetos Python intermediários. Use a abordagem Arrow para a maioria das cargas de trabalho por causa dessa eficiência.

Consulta básica para DataFrame

Essa abordagem recupera todas as linhas com o cursor padrão e constrói manualmente um DataFrame Polars. Aceita consultas parametrizadas para substituição segura de valores. Funciona sem o PyArrow, mas é mais lento para conjuntos de resultados grandes porque todo valor passa pelo 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

Se a sua cadeia de conexão usa Authentication=ActiveDirectoryDefault, o driver usa DefaultAzureCredential, que tenta vários provedores de credenciais em sequência. A primeira conexão pode ser lenta porque o SDK percorre a cadeia até encontrar um provedor funcionando. Em produção, se você sabe qual tipo de credencial seu ambiente usa, especifique-o diretamente (por exemplo, ActiveDirectoryMSI para identidade gerenciada) para evitar a caminhada em cadeia. Para obter mais informações, consulte Autenticação do Microsoft Entra.

A maneira mais eficiente de carregar dados SQL do Microsoft no Polars é através do Apache Arrow. O método arrow() do driver mssql-python retorna um(a) pyarrow.Table que o Polars pode consumir sem sobrecarga de cópia.

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 dados em lotes do Arrow

Para conjuntos de dados que não cabem na memória, use arrow_reader() para processar dados em lotes de streaming. Cada lote é um pyarrow.RecordBatch que o Polars pode consumir independentemente, de modo que o uso de memória permaneça proporcional a batch_size, e não ao 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")

Use o LazyFrames para execução diferida

Polars LazyFrames permite que você construa uma cadeia de operações (filtro, grupo, ordenação) sem executá-las imediatamente. Polars otimiza toda a cadeia antes de executar, o que pode ser mais rápido do que aplicar cada passo 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)

Escreva DataFrames Polars para Microsoft SQL

Coloque os identificadores entre aspas para evitar injeção de SQL

Nomes de tabelas e colunas não podem ser passados como parâmetros de consulta em SQL. Quando você cria instruções SQL com identificadores dinâmicos, coloque cada nome entre colchetes e escape quaisquer caracteres ] incorporados para evitar injeção de 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}]"

As funções auxiliares nesta seção usam quote_id() para todos os nomes de tabelas e colunas no SQL gerado.

Inserir linhas de DataFrame

A abordagem linha por linha itera sobre o DataFrame com iter_rows(named=True) e executa um INSERT por linha. Essa abordagem é simples, mas lenta para volumes grandes porque cada linha exige uma ida e volta até o 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 DataFrames grandes, use o método bulkcopy() do driver para enviar linhas em lote pelo protocolo TDS (Tabular Data Stream), o protocolo nativo de comunicação que o Microsoft SQL usa. Essa abordagem minimiza idas e voltas e é mais rápida do que inserções fileira por fileira.

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

Padrões de análise de dados

Os exemplos a seguir mostram tarefas comuns de análise que combinam consultas SQL do Microsoft com transformações Polars.

Consultas agregadas

Este exemplo agrupa produtos por subcategoria e calcula estatísticas de contagem e preços em SQL, depois carrega o resumo em um 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álise de série temporal

Carregue dados de séries temporais do Microsoft SQL e adicione colunas calculadas, como médias móveis, usando expressões do 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 dados SQL com arquivos locais

Você pode enriquecer dados SQL da Microsoft juntando-os com arquivos CSV locais em Polars. Carregue cada fonte de dados em um DataFrame e faça a junção na memória.

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

Padrões ETL

Construa pipelines de extração, transformação e carregamento (ETL) combinando consultas SQL do Microsoft com transformações Polars. As expressões do Polars lidam com a etapa de transformação, e bulkcopy() lida com o carregamento.

Extrair, transformar e carregar

Este exemplo extrai dados ativos de clientes através do Arrow, aplica lógica de segmentação de negócios com expressões Polars e carrega os resultados usando cópia em massa.

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)

Dicas de desempenho

As dicas a seguir ajudam você a aproveitar ao máximo a combinação mssql-python e Polars.

Deixe o Microsoft SQL cuidar do trabalho pesado

O Microsoft SQL é mais rápido para agregações, filtros e junções do que transferir todos os dados brutos pela rede e processá-los localmente em Python. Deixe o Microsoft SQL fazer o trabalho pesado sempre que possível, mova apenas os dados que você precisa pela rede e use Polars para análises e transformações que sejam mais convenientes em 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
""")

Use o Arrow para todas as operações de leitura

A transferência baseada em seta evita criar objetos Python intermediários, o que reduz o uso de memória e melhora a taxa de transferência. Prefira cursor.arrow() à conversão manual linha por linha para qualquer conjunto de resultados com mais de algumas linhas.

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