Megjegyzés
Az oldalhoz való hozzáféréshez engedély szükséges. Megpróbálhat bejelentkezni vagy módosítani a címtárat.
Az oldalhoz való hozzáféréshez engedély szükséges. Megpróbálhatja módosítani a címtárat.
A Polars egy nagy teljesítményű DataFrame könyvtár, amely Rust nyelven írt, és gyors, memória-hatékony alternatívát kínál a pandák számára. A Polars és az mssql-python driverrel kombinálva lehetővé teszik:
- Töltsd be az SQL lekérdezési eredményeket közvetlenül a Polars DataFrames-be.
- Használd az Apache Arrow-t a Microsoft SQL-ből történő nulla másolású adatátvitelhez.
- Írd vissza a Polars DataFrame-eket hatékonyan a Microsoft SQL-re.
- Építs nagy teljesítményű adatfeldolgozási folyamatokat lusta kiértékeléssel.
A cikkben szereplő példák a AdventureWorks mintaadatbázist kérdezik. Ha még nincs meg, nézd meg az AdventureWorks mintaadatbázisokat.
Adatok beolvasása Polars DataFrame-ekbe
A Microsoft SQL adatait kétféleképpen töltheted be a Polars-ba: soronként konvertálva szabványos kurzormóddal, vagy nulla másolat átvitelsel Apache Arrow-val. A "nulla másolat" azt jelenti, hogy az adatok egyetlen memóriapufferben maradnak, amelyet az illesztőprogram, az Arrow és a Polars mind közvetlenül olvasnak, így nem duplikálódnak a sorok köztes Python objektumokba. A legtöbb munkaterhelésnél használd az Arrow megközelítést ennek a hatékonyságnak köszönhetően.
Alapvető lekérdezés a DataFrame-hez
Ez a megközelítés az összes sort a szabványos kurzorral kéri le, és manuálisan épít fel egy Polars-adatkeretet. Paraméterezett lekérdezéseket fogad el biztonságos értékhelyettesítéshez. PyArrow nélkül működik, de nagy eredményhalmazoknál lassabb, mert minden érték átmegy a Python-on.
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)
Megjegyzés:
Ha a kapcsolati karakterlánc Authentication=ActiveDirectoryDefault használ, az illesztőprogram DefaultAzureCredential használ, amely több hitelesítőadat-szolgáltatót próbál ki egymás után. Az első kapcsolat lassú lehet, mert az SDK végigjárja a láncot, amíg meg nem talál egy működő szolgáltatót. A termelésben, ha tudod, melyik hitelesítéstípust használja a környezeted, közvetlenül megadd (például ActiveDirectoryMSI menedzselt identitásnál), hogy elkerüld a láncos sétát. További információ: Microsoft Entra-hitelesítés.
Használd az Arrow-t nulla másolat átvitelhez (ajánlott)
A leghatékonyabb módja annak, hogy Microsoft SQL adatokat töltsünk be a Polars-ba, az Apache Arrow-on keresztül. Az mssql-python illesztőprogram arrow() metódusa egy pyarrow.Table objektumot ad vissza, amelyet a Polars másolás nélküli többletterhelés nélkül fel tud használni.
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)
Nagy adathalmazok streamelése Arrow batches-szel
Olyan adatkészletek esetén, amelyek nem férnek el a memóriában, az adatok streaming kötegekben történő feldolgozásához használd a arrow_reader() elemet. Minden köteg egy pyarrow.RecordBatch, amelyet a Polars önállóan tud feldolgozni, így a memóriahasználat a batch_size méretével marad arányos, nem pedig a teljes eredményhalmaz méretével.
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")
Használd a LazyFrames-t a késleltetett végrehajtáshoz
A Polars LazyFrames lehetővé teszi, hogy egy műveletláncot építs ki (szűrő, csoportosítás, rendezés) anélkül, hogy azonnal lefuttatnád őket. A Polars optimalizálja a teljes láncot a végrehajtás előtt, ami gyorsabb lehet, mint az egyes lépések külön-külön alkalmazása.
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)
Polars-adatkeretek írása a Microsoft SQL-be
Lássa el idézőjelekkel az azonosítókat az SQL-befecskendezés megelőzése érdekében
A táblák és oszlopnevek nem adhatók át lekérdezési paraméterként SQL-ben. Ha dinamikus azonosítók használatával hoz létre SQL-utasításokat, tegyen minden nevet szögletes zárójelek közé, és escape-elje a névben szereplő ] karaktereket az SQL-injekció megelőzése érdekében.
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}]"
Az ebben a szakaszban ismertetett segédfüggvények a generált SQL-ben az összes tábla- és oszlopnévhez a quote_id() elemet használják.
DataFrame-sorok beszúrása
A soronkénti megközelítés a iter_rows(named=True) használatával végigiterál a DataFrame sorain, és soronként egy INSERT műveletet hajt végre. Ez a megközelítés egyszerű, de nagy adatmennyiségnél lassú, mert minden sor külön oda-vissza kommunikációt igényel a szerverrel.
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")
Tömeges beszúrás (nagy adatkeretekhez ajánlott)
Nagy adatkeretek esetén a meghajtó bulkcopy() módszerét használjuk, hogy sorokat tömegesen küldjenek a TDS (Tabular Data Stream) protokollon, amely a Microsoft SQL által használt natív vezetékes protokoll. Ez a megközelítés minimalizálja a oda-vissza utakat, és gyorsabb, mint sorról sorra történő beillesztés.
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")
Adatelemzési minták
Az alábbi példák olyan gyakori elemzési feladatokat mutatnak, amelyek a Microsoft SQL lekérdezéseket Polars transzformációkkal kombinálják.
Aggregált lekérdezések
Ez a példa alkategóriák szerint csoportosítja a termékeket, és SQL-ben számolja ki a szám- és árstatisztikákat, majd betölti az összefoglalót egy Polars DataFrame-be:
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)
Idősorozat-elemzések
Töltsd be az idősoros adatokat a Microsoft SQL-ből, és hozzáadj kiszámított oszlopokat, például gördülő átlagokat Polars kifejezésekkel.
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)
SQL adatok helyi fájlokkal való összekapcsolása
Gazdagíthatod a Microsoft SQL adatait azzal, hogy helyi CSV fájlokkal kombinálod őket a Polars-ban. Töltsd be az összes forrást egy DataFrame-be, és csatlakozzunk a memóriába.
# 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 minták
Kivonás, átalakítás és betöltés (ETL) pipeline-okat építsünk Microsoft SQL lekérdezések és Polars transzformációk kombinálásával. A Polars-kifejezések végzik az átalakítási lépést, a betöltést pedig a bulkcopy().
Kinyerés, átalakítás, betöltés
Ez a példa az aktív ügyféladatokat Arrow-on keresztül nyeri ki, üzleti szegmentációs logikát alkalmaz Polars kifejezésekkel, és az eredményeket tömeges másolattal tölti be.
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)
Teljesítménnyel kapcsolatos tippek
Az alábbi tippek segítenek a lehető legtöbbet hozni az mssql-python és a Polars kombinációból.
Hagyd, hogy a Microsoft SQL végezze a nehéz munkát
A Microsoft SQL gyorsabb az aggregációkhoz, szűréshez és csatlakozásokhoz, mint amikor az összes nyers adatot áthúznánk a vezetéken keresztül és helyben dolgoznánk fel Python-ban. Hagyd, hogy a Microsoft SQL végezze a nehezebb munkát, amikor csak lehet, csak a szükséges adatokat mozgasd át a hálózaton, és használd a Polars-t az elemzéshez és transzformációkhoz, amelyek kényelmesebbek Python-ban.
# 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
""")
Minden olvasási művelethez használd az Arrow-t
Az Arrow-alapú átvitel elkerüli a köztes Python objektumok létrehozását, ami csökkenti a memóriahasználatot és javítja az áteresztőképességet. Részesítse előnyben a(z) cursor.arrow() használatát a kézi, soronkénti átalakítással szemben minden olyan eredményhalmaz esetén, amely néhány sornál nagyobb.
# 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())