Observação
O acesso a essa página exige autorização. Você pode tentar entrar ou alterar diretórios.
O acesso a essa página exige autorização. Você pode tentar alterar os diretórios.
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.
Use o Arrow para transferência sem cópias (recomendado)
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")
Inserção em massa (recomendada para DataFrames grandes)
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())