Utilize o mssql-python com o Polars

Polars é uma biblioteca DataFrame de alto desempenho escrita em Rust que oferece uma alternativa rápida e eficiente em termos de memória aos pandas. O Polars, combinado com o driver mssql-python, permite-lhe:

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

Os exemplos deste artigo consultam a AdventureWorks base de dados de exemplo. Se ainda não o tens, vê as bases de dados de exemplo do AdventureWorks.

Leia dados em DataFrames Polars

Podes carregar dados SQL da Microsoft no Polars de duas formas: conversão linha a linha através de métodos padrão de cursor, ou transferência zero-copy através do Apache Arrow. "Zero-copy" significa que os dados permanecem num único buffer de memória que o driver, Arrow e Polars leem diretamente, pelo que nenhuma linha é duplicada em objetos Python intermédios. Use a abordagem Arrow para a maioria das cargas de trabalho devido a esta eficiência.

Consulta básica para DataFrame

Esta abordagem recolhe 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 todos os valores passam 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 ligação usar Authentication=ActiveDirectoryDefault, o controlador usa DefaultAzureCredential, que tenta vários fornecedores de credenciais em sequência. A primeira conexão pode ser lenta porque o SDK percorre a cadeia até encontrar um provedor funcional. Em produção, se souber que tipo de credencial o seu ambiente utiliza, especifique-o diretamente (por exemplo, ActiveDirectoryMSI para identidade gerida) para evitar o chain walk. Para obter mais informações, consulte Autenticação do Microsoft Entra.

A forma mais eficiente de carregar dados Microsoft SQL no Polars é através do Apache Arrow. O método arrow() do driver mssql-python retorna pyarrow.Table, que o Polars pode consumir sem sobrecusto 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)

Transmitir grandes conjuntos de dados em lotes 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, pelo que a utilização de memória se mantém proporcional a batch_size em vez de 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

Os Polars LazyFrames permitem-te construir uma cadeia de operações (filtro, agrupar, ordenar) sem as executar 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)

Gravar DataFrames do Polars no Microsoft SQL

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

Nomes de tabelas e colunas não podem ser passados como parâmetros de consulta em SQL. Quando constróis instruções SQL com identificadores dinâmicos, envolve cada nome em colchetes quadrados e escapa de quaisquer caracteres embutidos ] para evitar a injeção 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 secção utilizam quote_id() para todos os nomes de tabelas e colunas no SQL gerado.

Inserir linhas do DataFrame

A abordagem linha a linha itera sobre o DataFrame com iter_rows(named=True) e executa um INSERT por linha. Esta abordagem é simples mas lenta para volumes grandes porque cada linha requer uma ida e volta até ao 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 controlador para enviar linhas em bloco através do protocolo TDS (Tabular Data Stream), o protocolo nativo de comunicação que o Microsoft SQL utiliza. Esta abordagem minimiza as viagens de ida e volta e é mais rápida do que as inserções fila a 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")

Padrões de análise de dados

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

Consultas agregadas

Este exemplo agrupa os produtos por subcategoria e calcula estatísticas de contagem e preços em SQL, carregando depois o resumo num 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éries temporais

Carregar dados de séries temporais do Microsoft SQL e adicionar 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)

Juntar dados SQL com ficheiros locais

Pode enriquecer dados SQL da Microsoft juntando-os com ficheiros CSV locais em Polars. Carregue cada fonte para um DataFrame e faça a junção em 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

Crie pipelines de extração, transformação e carregamento (ETL) ao combinar consultas de Microsoft SQL com transformações de Polars. As expressões Polars tratam da etapa de transformação, e o bulkcopy() trata do carregamento.

Extrair, transformar, carregar

Este exemplo extrai dados de clientes ativos com o Arrow, aplica lógica de segmentação de negócio 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)

Sugestões de desempenho

As dicas seguintes ajudam-no a tirar o máximo proveito da combinação mssql-python e Polars.

Deixe o Microsoft SQL tratar do trabalho pesado

O Microsoft SQL é mais rápido para agregações, filtragem e junções do que transferir todos os seus dados brutos pela rede e processá-los localmente em Python. Deixa o Microsoft SQL fazer o trabalho pesado sempre que possível, move apenas os dados de que precisas pela rede e usa os Polars para análises e transformações 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 Arrow evita a criação de objetos Python intermédios, o que reduz a utilização de memória e melhora a taxa de transferência. Prefira cursor.arrow() em vez da conversão manual linha a linha para qualquer conjunto de resultados com mais do que 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())