Not
Bu sayfaya erişim yetkilendirme gerektiriyor. Oturum açmayı veya dizinleri değiştirmeyi deneyebilirsiniz.
Bu sayfaya erişim yetkilendirme gerektiriyor. Dizinleri değiştirmeyi deneyebilirsiniz.
Polars, Rust ile yazılmış yüksek performanslı bir DataFrame kütüphanesidir ve pandalara hızlı ve bellek açısından verimli bir alternatif sunar. Polars ile mssql-python sürücüsü birleştiğinde:
- SQL sorgu sonuçlarını doğrudan Polars DataFrames'e yükleyin.
- Microsoft SQL'den sıfır kopya veri aktarımı için Apache Arrow kullanın.
- Polars DataFrames'i Microsoft SQL'e verimli şekilde geri yaz.
- Tembel değerlendirmeyle yüksek performanslı veri boru hatları oluşturun.
Bu makaledeki örnekler örnek veritabanını AdventureWorks sorgulamaktadır. Henüz sahip değilseniz, AdventureWorks örnek veritabanlarına bakabilirsiniz.
Polars DataFrames'e veri okuma
Microsoft SQL verilerini Polars'a iki şekilde yükleyebilirsiniz: standart imleç yöntemleriyle satır satır dönüştürme veya Apache Arrow ile sıfır kopya aktarımı. "Sıfır kopyalama", verinin sürücü, Arrow ve Polars'ın doğrudan okunduğu tek bir bellek tamponunda kalması anlamına gelir; böylece satırlar ara Python nesnelerine çoğaltılmaz. Bu verimlilik nedeniyle çoğu iş yükü için Ok yaklaşımını kullanın.
DataFrame'e temel sorgu
Bu yaklaşım, standart imleçle tüm satırları getirir ve manuel olarak bir Polars DataFrame oluşturur. Güvenli değer yerine koyma için parametrizlenmiş sorguları kabul eder. PyArrow olmadan çalışıyor ama büyük sonuç kümelerinde daha yavaş çünkü her değer Python'dan geçiyor.
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
Bağlantı dizeniz Authentication=ActiveDirectoryDefault kullanıyorsa, sürücü DefaultAzureCredential kullanır; bu da birden çok kimlik bilgisi sağlayıcısını sırayla dener. İlk bağlantı yavaş olabilir çünkü SDK çalışan bir sağlayıcı bulana kadar zincirde yürür. Üretimde, ortamınızın hangi kimlik bilgisi türünü kullandığını biliyorsanız, zincirde dolaşmayı önlemek için bunu doğrudan belirtin (örneğin, yönetilen kimlik için ActiveDirectoryMSI). Daha fazla bilgi için bkz . Microsoft Entra kimlik doğrulaması.
Sıfır kopya transferi için Arrow kullanın (önerilir)
Microsoft SQL verilerini Polars'a yüklemenin en verimli yolu Apache Arrow aracılığıdır. mssql-python sürücüsünün arrow() yöntemi, Polars'ın sıfır kopyalama ek yükü olmadan işleyebileceği bir pyarrow.Table döndürür.
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)
Büyük veri kümelerini Arrow batch’leriyle akış olarak iletin
Belleğe sığmayan veri setleri için, verileri akış şeklinde partiler hâlinde işlemek üzere arrow_reader() kullanın. Her parti, Polars'ın bağımsız olarak işleyebileceği bir pyarrow.RecordBatch olduğundan, bellek kullanımı tam sonuç kümesi yerine batch_size ile orantılı kalır.
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")
Gecikmiş yürütme için LazyFrames kullanın
Polars LazyFrames, hemen çalıştırmadan bir işlem zinciri (filtre, gruplama, sıralama) oluşturmanıza olanak tanır. Polars, tüm zinciri çalıştırmadan önce optimize eder, bu da her adımı ayrı ayrı uygulamaktan daha hızlı olabilir.
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 Veri Çerçevelerini Microsoft SQL'e Yazma
SQL enjeksiyonunu önlemek için tanımlayıcıları tırnak içine alın
Tablo ve sütun adları SQL'de sorgu parametresi olarak iletilemez. Dinamik tanımlayıcılarla SQL ifadeleri oluştururken, her ismi karegüzel parantezlere sarın ve gömülü ] karakterlerden kaçınarak SQL enjeksiyonunu önleyin.
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}]"
Bu bölümdeki yardımcı fonksiyonlar, oluşturulan SQL'deki tüm tablo ve sütun adlarını kullanır quote_id() .
DataFrame satırlarını ekle
Satır satır yaklaşımı, DataFrame üzerinde iter_rows(named=True) ile yinelenir ve her satır için bir INSERT çalıştırır. Bu yaklaşım basittir ancak büyük hacimler için yavaştır çünkü her satır sunucuya gidiş-dönüş gerektirir.
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")
Toplu ekleme (büyük DataFrame'ler için önerilir)
Büyük DataFrame'ler için, sürücünün bulkcopy() yöntemini kullanarak Microsoft SQL'in kullandığı yerel kablo protokolü olan TDS (Tabular Data Stream) protokolü üzerinden satır toplu gönderebilirsiniz. Bu yaklaşım, gidiş-dönüş işlemlerini en aza indirir ve satır sıra eklemelerden daha hızlıdır.
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")
Veri analiz desenleri
Aşağıdaki örnekler, Microsoft SQL sorgularını Polars dönüşümleriyle birleştiren yaygın analiz görevlerini göstermektedir.
Toplulaştırma sorguları
Bu örnek, ürünleri alt kategorilere göre gruplar ve SQL'de sayım ile fiyat istatistiklerini hesaplar, ardından özeti bir Polars DataFrame'e yükler:
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)
Zaman serisi analizi
Microsoft SQL'den zaman serisi verilerini yükleyin ve Polars ifadeleriyle yuvarlanan ortalamalar gibi hesaplanan sütunları ekleyin.
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 verilerini yerel dosyalarla birleştir
Microsoft SQL verilerini Polars'taki yerel CSV dosyalarıyla birleştirerek zenginleştirebilirsiniz. Her kaynağı bir DataFrame'e yükleyin ve belleğe ekleyin.
# 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 desenleri
Microsoft SQL sorgularını Polars dönüşümleriyle birleştirerek çıkar, dönüştürme ve yükleme (ETL) boru hatları oluşturun. Polar ifadeleri dönüşüm adımını ve bulkcopy() yüklemeyi yönetir.
Ayıklama, dönüştürme, yükleme
Bu örnek, Arrow üzerinden aktif müşteri verilerini çıkarır, Polars ifadeleriyle iş segmentasyon mantığını uygular ve sonuçları toplu kopya ile yükler.
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)
Performans ipuçları
Aşağıdaki ipuçları, mssql-python ve Polars kombinasyonundan en iyi şekilde yararlanmanıza yardımcı olur.
Bırakın Microsoft SQL ağır işleri üstlensin
Microsoft SQL, toplama, filtreleme ve birleştirme işlemlerinde, tüm ham verilerinizi kablo üzerinden çekip yerel olarak Python'da işlemektense daha hızlı. Mümkün olduğunda Microsoft SQL'in ağır işi yapmasına izin verin, sadece ihtiyacınız olan verileri ağ üzerinden taşıyın ve Python'da daha kullanışlı analiz ve dönüşümler için Polars kullanın.
# 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
""")
Tüm okuma işlemleri için Arrow kullanın
Ok tabanlı transfer, ara Python nesneleri oluşturmayı engeller; bu da bellek kullanımını azaltır ve veri verimliliğini artırır. Birkaç satırdan büyük tüm sonuç kümeleri için manuel, satır satır dönüştürme yerine cursor.arrow() tercih edin.
# 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())