Kommentar
Åtkomst till den här sidan kräver auktorisering. Du kan prova att logga in eller ändra kataloger.
Åtkomst till den här sidan kräver auktorisering. Du kan prova att ändra kataloger.
Polars är ett högpresterande DataFrame-bibliotek skrivet i Rust som erbjuder ett snabbt och minneseffektivt alternativ till pandas. Polars i kombination med mssql-python-drivrutinen låter dig:
- Ladda SQL-frågeresultat direkt i Polars DataFrames.
- Använd Apache Arrow för nollkopieringsöverföring från Microsoft SQL.
- Skriv Polars DataFrames tillbaka till Microsoft SQL effektivt.
- Bygg högpresterande datapipeliner med lat utvärdering.
Exemplen i denna artikel gör frågor mot AdventureWorks exempeldatabasen. Om du inte redan har det, se AdventureWorks exempeldatabaser.
Läs in data i Polars DataFrames
Du kan ladda in Microsoft SQL-data i Polars på två sätt: rad-för-rad-konvertering via standardmetoder för markörer, eller zero-copy transfer via Apache Arrow. "Zero-copy" innebär att datan stannar i en enda minnesbuffert som drivrutinen, Arrow och Polars alla läser direkt, så inga rader dupliceras till mellanliggande Python-objekt. Använd Arrow-metoden för de flesta arbetsbelastningar på grund av denna effektivitet.
Grundläggande fråga till DataFrame
Denna metod hämtar alla rader med standardmarkören och bygger manuellt en Polars DataFrame. Den accepterar parameteriserade frågor för säker värdesubstitution. Det fungerar utan PyArrow men är långsammare för stora resultatmängder eftersom varje värde passerar genom 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
Om din anslutningssträng använder Authentication=ActiveDirectoryDefault, använder drivrutinen DefaultAzureCredential, som försöker med flera autentiseringsleverantörer i följd. Den första anslutningen kan vara långsam eftersom SDK:n går igenom kedjan tills den hittar en fungerande leverantör. I produktion, om du vet vilken typ av behörighet din miljö använder, ange det direkt (till exempel ActiveDirectoryMSI för managed identity) för att undvika kedjevandring. Mer information finns i Microsoft Entra-autentisering.
Använd Arrow för överföring utan kopiering (rekommenderas)
Det mest effektiva sättet att ladda in Microsoft SQL-data i Polars är via Apache Arrow. mssql-python-drivrutinens arrow() metod returnerar en pyarrow.Table som Polars kan använda utan någon kopieringsöverhead.
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)
Strömma stora datamängder med Arrow-batcher
För datauppsättningar som inte ryms i minnet, använd arrow_reader() för att bearbeta data i batchar i strömmande form. Varje batch är en pyarrow.RecordBatch som Polars kan konsumera oberoende, så minnesanvändningen förblir proportionell mot batch_size snarare än hela resultatuppsättningen.
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")
Använd LazyFrames för uppskjuten exekvering
Polars LazyFrames låter dig bygga en kedja av operationer (filtrera, gruppera, sortera) utan att köra dem direkt. Polars optimerar hela kedjan innan den utförs, vilket kan gå snabbare än att tillämpa varje steg individuellt.
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)
Skriv Polars DataFrames till Microsoft SQL
Citatidentifierare för att förhindra SQL-injektion
Tabell- och kolumnnamn kan inte skickas som frågeparametrar i SQL. När du bygger SQL-satser med dynamiska identifierare, slå in varje namn inom hakparenteser och undvik inbäddade ] tecken för att förhindra SQL-injektion.
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}]"
Hjälpfunktionerna i detta avsnitt används quote_id() för alla tabell- och kolumnnamn i genererad SQL.
Infoga rader i DataFrame
Rad-för-rad-metoden itererar över dataramen med iter_rows(named=True) och kör en INSERT per rad. Denna metod är enkel men långsam för stora volymer eftersom varje rad kräver en rundresa till servern.
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")
Massinfogning (rekommenderas för stora DataFrames)
För stora DataFrames använder du drivrutinens bulkcopy() metod för att skicka rader i bulk över TDS (Tabular Data Stream)-protokollet, det inbyggda trådprotokollet som Microsoft SQL använder. Denna metod minimerar rundresor och är snabbare än rad-för-rad-insatser.
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")
Dataanalysmönster
Följande exempel visar vanliga analysuppgifter som kombinerar Microsoft SQL-frågor med Polars-transformationer.
Aggregerade frågor
Detta exempel grupperar produkter efter underkategori och beräknar antals- och prisstatistik i SQL, och laddar sedan sammanfattningen i en 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)
Tidsserieanalys
Ladda tidsseriedata från Microsoft SQL och lägg till beräknade kolumner som rullande medelvärden med hjälp av Polars-uttryck.
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)
Koppla ihop SQL-data med lokala filer
Du kan berika Microsoft SQL-data genom att koppla ihop dem med lokala CSV-filer i Polars. Ladda varje källa i en DataFrame och koppla in i minnet.
# 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-mönster
Bygg extrakt, transformera och ladda (ETL) pipelines genom att kombinera Microsoft SQL-frågor med Polars-transformationer. Polars-uttryck hanterar transformationssteget och bulkcopy() hanterar belastningen.
Extrahera, transformera, läsa in
Detta exempel extraherar aktiv kunddata via Arrow, tillämpar affärssegmenteringslogik med Polars-uttryck och laddar resultaten med bulk-kopiering.
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)
Prestandatips
Följande tips hjälper dig att få ut det mesta av kombinationen mssql-python och Polars.
Låt Microsoft SQL ta hand om det tunga arbetet
Microsoft SQL är snabbare vid aggregering, filtrering och sammanslagningar än att hämta all din rådata över nätverket och bearbeta den lokalt i Python. Låt Microsoft SQL göra det tunga arbetet när det är möjligt, flytta bara den data du behöver över nätverket och använd Polars för analys och transformationer som är mer bekväma i 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
""")
Använd Arrow för alla läsoperationer
Pilbaserad överföring undviker att skapa mellanliggande Python-objekt, vilket minskar minnesanvändningen och förbättrar genomströmningen. Föredra cursor.arrow() framför manuell rad-för-rad-konvertering för alla resultatuppsättningar som är större än några rader.
# 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())