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.
Pandas kütüphanesi, Python'un birincil veri analiz aracıdır. Pandas'ı mssql-python sürücüsüyle birleştirerek şunları yapabilirsiniz:
- SQL sorgu sonuçlarını doğrudan DataFrames'e yükleyin.
- DataFrames'i Microsoft SQL'e verimli şekilde geri yaz.
- ETL işlemlerini gerçekleştirin.
- Veri boru hatları oluşturun.
Bu makaledeki örnekler, Production.Product tablo ve diğer tabloları sorgular. Veri yazan örnekler, örnek veriyi değiştirmemek için geçici tablolar kullanır.
Analiz örneklerinde (Sales.SalesOrderHeader, Sales.SalesOrderDetail, Production.ProductSubcategory) referans verilen diğer tablolar AdventureWorks'un bir parçasıdır. Bu kalıpları uyarlarken kendi tablolarınızı kullanın.
Veri Çerçevelerine Veri Okuma
mssql-python sürücüsü, satırları Python nesneleri olarak döndürür; bunları, sütun adlarını cursor.description etiketinden ve satır değerlerini fetchall() etiketinden okuyarak pandas DataFrame'lerine dönüştürürsünüz. Bu bölümdeki yardımcı fonksiyonlar bu dönüşümü yeniden kullanılabilir desenlere sarar.
DataFrame'e temel sorgu
Bu fonksiyon, parametreli bir sorgu çalıştırır ve tam sonuç kümesinden bir DataFrame oluşturur. Hafızaya rahatça oturan sonuç setleri için iyi çalışıyor.
import pandas as pd
import mssql_python
conn = mssql_python.connect(connection_string)
cursor = conn.cursor()
def query_to_dataframe(cursor, query: str, params: dict = None) -> pd.DataFrame:
"""Execute query and return results as DataFrame."""
cursor.execute(query, params or {})
# cursor.description is a list of tuples, one per column.
# Each tuple's first element is the column name.
columns = [col[0] for col in cursor.description]
# Fetch all rows
rows = cursor.fetchall()
# Convert to DataFrame
data = [tuple(row) for row in rows]
return pd.DataFrame(data, columns=columns)
# Usage: %(cat)s is a parameterized placeholder. The driver safely substitutes
# the value from the dict, which prevents SQL injection.
df = query_to_dataframe(cursor, "SELECT * FROM Production.Product WHERE ProductSubcategoryID = %(cat)s", {"cat": 5})
print(df.head())
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ı.
Büyük veri setlerini akış
Milyonlarca satırlı tablolar için, her şeyi aynı anda yüklemek hafızayı tüketebilir. Parçalı yaklaşım, satırları fetchmany() ile toplu halde getirir ve sonuçları birleştirir; böylece en yüksek bellek kullanımını tam sonuç kümesi yerine chunksize ile orantılı tutar.
def query_to_dataframe_chunked(cursor, query: str, params: dict = None,
chunksize: int = 10000) -> pd.DataFrame:
"""Load large query results in chunks for memory efficiency."""
cursor.execute(query, params or {})
columns = [col[0] for col in cursor.description]
chunks = []
while True:
rows = cursor.fetchmany(chunksize)
if not rows:
break
data = [tuple(row) for row in rows]
chunks.append(pd.DataFrame(data, columns=columns))
return pd.concat(chunks, ignore_index=True) if chunks else pd.DataFrame(columns=columns)
# Usage for large tables
df = query_to_dataframe_chunked(cursor, "SELECT * FROM Production.TransactionHistory", chunksize=50000)
Büyük veri kümeleri için üretici
Veriyi tüm sonucu bellekte tutmadan kademeli olarak işlemeniz gerekiyorsa, bir jeneratör kullanın. Her yield, bir sonrakini almadan önce işleyip atabileceğiniz bir DataFrame parçası üretir.
def query_to_dataframe_generator(cursor, query: str, params: dict = None,
chunksize: int = 10000):
"""Yield DataFrame chunks for processing without loading all data."""
cursor.execute(query, params or {})
columns = [col[0] for col in cursor.description]
while True:
rows = cursor.fetchmany(chunksize)
if not rows:
break
data = [tuple(row) for row in rows]
yield pd.DataFrame(data, columns=columns)
# Process chunks without loading entire dataset
huge_query = """
SELECT * FROM Production.TransactionHistory
UNION ALL SELECT * FROM Production.TransactionHistory
UNION ALL SELECT * FROM Production.TransactionHistory
"""
for chunk_df in query_to_dataframe_generator(cursor, huge_query):
# Process each chunk, then discard it before the next fetch
print(f"Processing chunk of {len(chunk_df)} rows")
total_cost = chunk_df["ActualCost"].sum()
print(f"Chunk total cost: {total_cost}")
DataFrame'leri 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
En basit yaklaşım, DataFrame satırları üzerinden geçer ve her satırda bir tane INSERT verir. Basit yaklaşım küçük DataFrame'ler için işe yarar, ancak büyük hacimler için yavaştır çünkü her satır sunucuya ayrı bir gidiş-dönüş gerektirir.
def dataframe_to_sql(cursor, conn, df: pd.DataFrame, table: str,
if_exists: str = "append") -> int:
"""Write DataFrame to Microsoft SQL table."""
if if_exists == "replace":
cursor.execute(f"TRUNCATE TABLE {quote_id(table)}")
columns = df.columns.tolist()
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.iterrows():
params = {col: (None if pd.isna(val) else val) for col, val in row.items()}
cursor.execute(query, params)
rows_inserted += 1
conn.commit()
return rows_inserted
# Usage
cursor.execute("""
CREATE TABLE #Products (
Name NVARCHAR(100),
ListPrice DECIMAL(10,2),
ProductSubcategoryID INT
)
""")
df = pd.DataFrame({
"Name": ["Product A", "Product B"],
"ListPrice": [29.99, 49.99],
"ProductSubcategoryID": [1, 2]
})
rows = dataframe_to_sql(cursor, conn, df, "#Products")
print(f"Inserted {rows} rows")
BCP ile toplu ekleme (büyük DataFrame'ler için önerilir)
Büyük DataFrame'ler için, Microsoft SQL'in kullandığı yerel kablo protokolü olan TDS (Tabular Data Stream) protokolü üzerinden satır toplu olarak gönderen sürücü bulkcopy() yöntemini kullanın. Bu yaklaşım, gidiş-dönüş sayısını en aza indirdiği için satır satır eklemelerden daha hızlıdır.
def dataframe_to_sql_bulk(conn, df: pd.DataFrame, table: str) -> int:
"""Bulk insert DataFrame using BCP for better performance."""
# Convert DataFrame to list of tuples, handling NaN
rows = []
for _, row in df.iterrows():
row_data = tuple(None if pd.isna(v) else v for v in row)
rows.append(row_data)
cursor = conn.cursor()
result = cursor.bulkcopy(table, rows)
conn.commit()
return result["rows_copied"]
# Usage
cursor.execute("CREATE TABLE ##PandasProducts (Name NVARCHAR(50), ListPrice DECIMAL(10,2), ProductSubcategoryID INT)")
conn.commit()
df = pd.DataFrame({
"Name": ["Product A", "Product B", "Product C"],
"ListPrice": [29.99, 49.99, 19.99],
"ProductSubcategoryID": [1, 2, 1]
})
rows = dataframe_to_sql_bulk(conn, df, "##PandasProducts")
DataFrame'den mevcut satırları güncelle
Tabloda zaten var olan satırları güncellemek için DataFrame üzerinde yineleme yapın ve parametreli UPDATE ifadeler yayınlayın. Hangi key_column satırın güncelleneceğini belirler.
def update_from_dataframe(cursor, conn, df: pd.DataFrame, table: str,
key_column: str) -> int:
"""Update existing rows based on key column."""
columns = [col for col in df.columns if col != key_column]
set_clause = ", ".join([f"{quote_id(col)} = %({col})s" for col in columns])
query = f"UPDATE {quote_id(table)} SET {set_clause} WHERE {quote_id(key_column)} = %({key_column})s"
rows_updated = 0
for _, row in df.iterrows():
params = {col: (None if pd.isna(val) else val) for col, val in row.items()}
cursor.execute(query, params)
rows_updated += cursor.rowcount
conn.commit()
return rows_updated
# Usage
cursor.execute("""
CREATE TABLE #ProductPrices (
ProductID INT PRIMARY KEY,
ListPrice DECIMAL(10,2)
);
INSERT INTO #ProductPrices VALUES (1, 29.99), (2, 49.99), (3, 19.99);
""")
conn.commit()
df_updates = pd.DataFrame({
"ProductID": [1, 2, 3],
"ListPrice": [31.99, 52.99, 21.99]
})
updated = update_from_dataframe(cursor, conn, df_updates, "#ProductPrices", "ProductID")
Upsert (birleştirme) deseni
Bazı satırlar yeni olup diğerleri zaten var olabilirse, MERGE tek bir işlemde SQL ifadesi eklemek veya güncellemek için kullanın.
MERGE her gelen satırı hedef tabloyla anahtar sütunları kullanarak karşılaştırır. Eşleşme bulunursa güncellenir; aksi takdirde eklenir.
MERGE var olup olmadığını ayrıca denetlemekten kaçınır.
def upsert_from_dataframe(cursor, conn, df: pd.DataFrame, table: str,
key_columns: list[str]) -> int:
"""Insert or update rows based on key columns. Returns total rows affected."""
all_columns = df.columns.tolist()
value_columns = [c for c in all_columns if c not in key_columns]
total_affected = 0
for _, row in df.iterrows():
params = {col: (None if pd.isna(val) else val) for col, val in row.items()}
# Build MERGE statement with quoted identifiers
key_match = " AND ".join([f"t.{quote_id(k)} = s.{quote_id(k)}" for k in key_columns])
update_set = ", ".join([f"{quote_id(c)} = s.{quote_id(c)}" for c in value_columns])
all_cols = ", ".join([quote_id(c) for c in all_columns])
all_vals = ", ".join([f"%({c})s" for c in all_columns])
cursor.execute(f"""
MERGE {quote_id(table)} AS t
USING (SELECT {', '.join([f'%({c})s AS {quote_id(c)}' for c in all_columns])}) AS s
ON {key_match}
WHEN MATCHED THEN UPDATE SET {update_set}
WHEN NOT MATCHED THEN INSERT ({all_cols}) VALUES ({all_vals});
""", params)
total_affected += cursor.rowcount
conn.commit()
return total_affected
Veri analiz desenleri
Aşağıdaki örnekler, Microsoft SQL sorgularını pandas dönüşümleriyle birleştiren yaygın analiz görevlerini göstermektedir.
Soruları DataFrame'e topla
def get_sales_summary(cursor) -> pd.DataFrame:
"""Get sales summary by category."""
return query_to_dataframe(cursor, """
SELECT
pc.Name AS CategoryName,
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 pc ON p.ProductSubcategoryID = pc.ProductSubcategoryID
GROUP BY pc.Name
ORDER BY ProductCount DESC
""")
df = get_sales_summary(cursor)
print(df.to_string())
Zaman serisi verileri
Microsoft SQL'den alınan zaman serisi verileriyle çalışmak için pandas tarih indeksleme ve yeniden örnekleme kullanın. Rolling average ve resampling gibi işlemleri etkinleştirmek için tarih sütununu DataFrame index olarak ayarlayın.
def get_daily_sales(cursor, start_date: str, end_date: str) -> pd.DataFrame:
"""Get daily sales time series."""
df = query_to_dataframe(cursor, """
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})
# Set date as index for time series operations
df["Date"] = pd.to_datetime(df["Date"])
df.set_index("Date", inplace=True)
return df
# Usage
sales_df = get_daily_sales(cursor, "2024-01-01", "2024-12-31")
# Resample to weekly
weekly = sales_df.resample("W").sum()
# Calculate rolling average
sales_df["RollingAvg"] = sales_df["Revenue"].rolling(window=7).mean()
SQL verilerinden pivot tablolar
Pivot tablolar, verileri satırlardan matris formatına dönüştürür. Yıl, ay ve kategori gibi boyutlara göre yeniden düzenlemek için Microsoft SQL'den ham verileri alın, ardından pivot_table() kullanın.
def get_sales_pivot(cursor) -> pd.DataFrame:
"""Get sales data and create pivot table."""
df = query_to_dataframe(cursor, """
SELECT
YEAR(soh.OrderDate) AS Year,
MONTH(soh.OrderDate) AS Month,
pc.Name AS CategoryName,
SUM(sod.OrderQty * sod.UnitPrice) AS Revenue
FROM Sales.SalesOrderHeader soh
JOIN Sales.SalesOrderDetail sod ON soh.SalesOrderID = sod.SalesOrderID
JOIN Production.Product p ON sod.ProductID = p.ProductID
JOIN Production.ProductSubcategory pc ON p.ProductSubcategoryID = pc.ProductSubcategoryID
GROUP BY YEAR(soh.OrderDate), MONTH(soh.OrderDate), pc.Name
""")
# Create pivot table
pivot = df.pivot_table(
values="Revenue",
index=["Year", "Month"],
columns="CategoryName",
aggfunc="sum",
fill_value=0
)
return pivot
pivot_df = get_sales_pivot(cursor)
print(pivot_df)
ETL desenleri
Ayıklama, dönüştürme ve yükleme iş hatları oluşturmak için Microsoft SQL sorgularını pandas dönüşümleriyle birleştirerek ayıklama, dönüştürme ve yükleme iş hatları oluşturun. Sürücü çıkarma ve yüklemeyi yaparken, pandas dönüşüm adımını üstlenir.
Ayıklama, dönüştürme, yükleme
Bu örnek, aktif müşteri verilerini çıkarır, müşterileri segmentlere ayırmak için iş kuralları uygular ve sonuçları hedef tabloya yükler.
def etl_pipeline(source_cursor, dest_cursor, dest_conn):
"""Simple ETL pipeline with pandas."""
# Extract
df = query_to_dataframe(source_cursor, """
SELECT
c.CustomerID,
COUNT(soh.SalesOrderID) AS OrderCount,
SUM(soh.TotalDue) AS TotalSpent
FROM Sales.Customer c
JOIN Sales.SalesOrderHeader soh ON c.CustomerID = soh.CustomerID
WHERE soh.OrderDate > DATEADD(YEAR, -1, GETDATE())
GROUP BY c.CustomerID
""")
# Transform
df["CustomerSegment"] = pd.cut(
df["TotalSpent"],
bins=[0, 100, 500, 1000, float("inf")],
labels=["Bronze", "Silver", "Gold", "Platinum"]
)
df["AvgOrderValue"] = df["TotalSpent"] / df["OrderCount"].replace(0, 1)
df["IsHighValue"] = df["TotalSpent"] > 500
# Load
dataframe_to_sql_bulk(dest_conn, df[["CustomerID", "CustomerSegment", "AvgOrderValue", "IsHighValue"]],
"#CustomerAnalytics")
return len(df)
Artan yük deseni
Devam eden veri boru hatları için, sadece son çalıştırmadan beri değişen kayıtları yükleyin. Bu yaklaşım, maksimum zaman damgası için hedef tabloda sorgular ve ardından kaynaktan yalnızca yeni kayıtları getirir.
def incremental_load(cursor, conn, source_table: str, dest_table: str,
timestamp_col: str) -> int:
"""Load only new/changed records based on timestamp."""
# Get last loaded timestamp
cursor.execute(f"SELECT MAX({quote_id(timestamp_col)}) FROM {quote_id(dest_table)}")
last_loaded = cursor.fetchval()
# Build query for new records
if last_loaded:
df = query_to_dataframe(cursor, f"""
SELECT * FROM {quote_id(source_table)}
WHERE {quote_id(timestamp_col)} > %(last)s
""", {"last": last_loaded})
else:
df = query_to_dataframe(cursor, f"SELECT * FROM {quote_id(source_table)}")
if df.empty:
return 0
# Load new records
return dataframe_to_sql_bulk(conn, df, dest_table)
Performans ipuçları
Uygun veri türlerini kullanma
Pandas sayılar için varsayılan olarak 64 bit tiplere geçer, küçük tipler yeterli olduğunda bellek kaybı yapar. Tam sayıları ve kayan sayıları aşağıya kaçırmak ve düşük kardinalite dizi sütunlarını kategoriklere dönüştürmek, bellek kullanımını önemli ölçüde azaltabilir.
def optimize_dataframe_types(df: pd.DataFrame) -> pd.DataFrame:
"""Optimize DataFrame memory usage."""
for col in df.columns:
col_type = df[col].dtype
if col_type == "int64":
# Downcast integers
df[col] = pd.to_numeric(df[col], downcast="integer")
elif col_type == "float64":
# Downcast floats
df[col] = pd.to_numeric(df[col], downcast="float")
elif col_type == "object":
# Convert to category if low cardinality
num_unique = df[col].nunique()
if num_unique / len(df) < 0.5:
df[col] = df[col].astype("category")
return df
Ağır kaldırmak için SQL kullanın
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 pandas kullanın.
# Avoid: pulling all rows over the wire to aggregate locally in pandas
df_all = query_to_dataframe(cursor, "SELECT * FROM Production.Product") # transfers entire table
summary = df_all.groupby("Color").agg({"ListPrice": "sum"}) # aggregation that SQL can do faster
# Better: push the aggregation into SQL and transfer only the summary
df = query_to_dataframe(cursor, """
SELECT Color, SUM(ListPrice) AS TotalPrice
FROM Production.Product
WHERE Color IS NOT NULL
GROUP BY Color
""")
Toplu yazma işlemleri
Tek bir toplu ekleme için fazla büyük olan DataFrame'lerde, işlemi partiler hâlinde bölün ve ilerlemeyi izleyin.
def batch_insert(cursor, conn, df: pd.DataFrame, table: str, batch_size: int = 1000):
"""Insert in batches with progress tracking."""
total = len(df)
for i in range(0, total, batch_size):
batch = df.iloc[i:i + batch_size]
dataframe_to_sql(cursor, conn, batch, table)
print(f"Inserted {min(i + batch_size, total)}/{total}")