Notitie
Voor toegang tot deze pagina is autorisatie vereist. U kunt proberen u aan te melden of de directory te wijzigen.
Voor toegang tot deze pagina is autorisatie vereist. U kunt proberen de mappen te wijzigen.
Polars is een high-performance DataFrame-bibliotheek geschreven in Rust die een snel, geheugenefficiënt alternatief voor pandas biedt. Polars gecombineerd met de mssql-python-driver laat je het volgende doen:
- Laad resultaten van SQL-query’s direct in Polars DataFrames.
- Gebruik Apache Arrow voor zero-copy datatransfer vanuit Microsoft SQL.
- Schrijf Polars DataFrames efficiënt terug naar Microsoft SQL.
- Bouw krachtige data-pijplijnen met luie evaluatie.
De voorbeelden in dit artikel zoeken de AdventureWorks voorbeelddatabase op. Als je het nog niet hebt, bekijk dan de voorbeelddatabases van AdventureWorks.
Laad gegevens in Polars-DataFrames
Je kunt Microsoft SQL-gegevens op twee manieren in Polars laden: rij-voor-rij conversie via standaard cursormethoden, of zero-copy transfer via Apache Arrow. "Zero-copy" betekent dat de data in één enkele geheugenbuffer blijft die de driver, Arrow en Polars direct lezen, zodat geen rijen worden gedupliceerd in tussenliggende Python-objecten. Gebruik de Arrow-aanpak voor de meeste workloads vanwege deze efficiëntie.
Eenvoudige query voor een DataFrame
Deze aanpak haalt alle rijen op met de standaardcursor en bouwt handmatig een Polars DataFrame. Het accepteert geparametriseerde queries voor veilige waardesubstitutie. Het werkt zonder PyArrow, maar is trager voor grote resultatensets omdat elke waarde door Python gaat.
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)
Opmerking
Als je verbindingsreeks Authentication=ActiveDirectoryDefault gebruikt, gebruikt de driver DefaultAzureCredential, die meerdere referentieproviders opeenvolgend probeert. De eerste verbinding kan traag zijn omdat de SDK de keten doorloopt totdat hij een werkende provider vindt. In productie, als je weet welk type inloggegevens je omgeving gebruikt, specificeer het dan direct (bijvoorbeeld ActiveDirectoryMSI voor managed identity) om de chain walk te voorkomen. Zie Microsoft Entra-verificatie voor meer informatie.
Gebruik Arrow voor zero-copy transfer (aanbevolen)
De meest efficiënte manier om Microsoft SQL-gegevens in Polars te laden is via Apache Arrow. De mssql-python-driver methode arrow() levert een pyarrow.Table die Polars kan gebruiken zonder kopieeroverhead.
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)
Stream grote datasets met Arrow-batches
Voor datasets die niet in het geheugen passen, gebruik je arrow_reader() om data in streaming-batches te verwerken. Elke batch is een pyarrow.RecordBatch die Polars onafhankelijk kan verwerken, zodat het geheugengebruik evenredig blijft aan batch_size in plaats van aan de volledige resultatenset.
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")
Gebruik LazyFrames voor uitgestelde uitvoering
Polars LazyFrames laat je een keten van bewerkingen bouwen (filteren, groeperen, sorteren) zonder ze direct uit te voeren. Polars optimaliseert de volledige keten voordat het wordt uitgevoerd, wat sneller kan zijn dan elke stap afzonderlijk toe te passen.
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)
Schrijf Polars DataFrames naar Microsoft SQL
Quote-identificaties om SQL-injectie te voorkomen
Tabel- en kolomnamen kunnen niet als queryparameters in SQL worden doorgegeven. Wanneer je SQL-statements met dynamische identificaties bouwt, wikkel je elke naam tussen vierkante haken en verwijder je ingebedde ] tekens om SQL-injectie te voorkomen.
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}]"
De hulpfuncties in deze sectie gebruiken quote_id() voor alle tabel- en kolomnamen in gegenereerde SQL.
DataFrame-rijen invoegen
De rij-voor-rij benadering itereert over het DataFrame met iter_rows(named=True) en voert er één INSERT uit per rij. Deze aanpak is eenvoudig maar traag voor grote volumes omdat elke rij een retour naar de server vereist.
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")
Bulksgewijs invoegen (aanbevolen voor grote DataFrames)
Voor grote DataFrames gebruik je de drivermethode bulkcopy() om rijen in bulk te verzenden via het TDS (Tabular Data Stream) protocol, het native wire-protocol dat Microsoft SQL gebruikt. Deze aanpak beperkt het heen-en-weerverkeer en is sneller dan rij-voor-rij invoegingen.
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")
Patronen voor data-analyse
De volgende voorbeelden tonen veelvoorkomende analysetaken die Microsoft SQL-queries combineren met Polars-transformaties.
Aggregate queries
Dit voorbeeld groepeert producten per subcategorie en berekent aantal- en prijsstatistieken in SQL, waarna de samenvatting wordt geladen in een 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)
Tijdreeksanalyse
Laadtijdreeksgegevens van Microsoft SQL en voeg berekende kolommen toe zoals rollend gemiddelden met behulp van Polars-expressies.
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)
Verbind SQL-gegevens met lokale bestanden
Je kunt Microsoft SQL-gegevens verrijken door deze te joinen met lokale CSV-bestanden in Polars. Laad elke bron in een DataFrame en voeg deze samen in het werkgeheugen.
# 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)
ETL-patronen
Bouw extract, transform en laad (ETL) pijplijnen door Microsoft SQL-queries te combineren met Polars-transformaties. Polars-expressies verzorgen de transformatiestap, en bulkcopy() verzorgt het laden.
Extraheren, transformeren, laden
Dit voorbeeld haalt actieve klantgegevens uit via Arrow, past bedrijfssegmentatielogica toe met Polars-expressies en laadt de resultaten met bulk copy.
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)
Tips voor prestaties
De volgende tips helpen je om het meeste uit de combinatie van mssql-python en Polars te halen.
Laat Microsoft SQL het zware werk doen
Microsoft SQL is sneller voor aggregaties, filtering en joins dan al je ruwe data via de kabel ophalen en lokaal verwerken in Python. Laat Microsoft SQL waar mogelijk het zware werk doen, verplaats alleen de data die je nodig hebt over het netwerk en gebruik Polars voor analyses en transformaties die handiger zijn in 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
""")
Gebruik Arrow voor alle leesbewerkingen
Pijl-gebaseerde overdracht voorkomt het creëren van tussenliggende Python-objecten, wat het geheugengebruik vermindert en de doorvoer verbetert. Gebruik bij voorkeur cursor.arrow() in plaats van handmatige rij-voor-rijconversie voor resultaatsets van meer dan een paar rijen.
# 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())