Pandas ile mssql-python kullanın

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")

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}")