Gebruik mssql-python met pandas

De pandas-bibliotheek is de belangrijkste tool van Python voor data-analyse. Door pandas te combineren met de mssql-python-driver, kun je:

  • Laad de resultaten van SQL-query's rechtstreeks in DataFrames.
  • Schrijf DataFrames efficiënt terug naar Microsoft SQL.
  • Voer ETL-operaties uit.
  • Maak datapipelines.

De voorbeelden in dit artikel vragen naar de Production.Product tabel en andere tabellen in de AdventureWorks-voorbeelddatabase. Voorbeelden die data schrijven gebruiken tijdelijke tabellen om het wijzigen van voorbeeldgegevens te voorkomen.

Andere tabellen die in analysevoorbeelden worden genoemd (Sales.SalesOrderHeader, Sales.SalesOrderDetail, Production.ProductSubcategory) maken deel uit van AdventureWorks. Gebruik je eigen tabellen ter vervanging van de tabellen in deze patronen wanneer je deze aanpast.

Lees gegevens in DataFrames

De mssql-python-driver retourneert rijen als Python-objecten, die je omzet naar pandas DataFrames door kolomnamen op te halen uit cursor.description en rijwaarden uit fetchall(). De helperfuncties in deze sectie omvatten die omzetting in herbruikbare patronen.

Eenvoudige query naar DataFrame

Deze functie voert een geparametriseerde query uit en bouwt een DataFrame op van de volledige resultaatset. Het werkt goed voor resultaatsets die comfortabel in het geheugen passen.

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

Opmerking

Als je verbindingsreeks Authentication=ActiveDirectoryDefault gebruikt, gebruikt de driver DefaultAzureCredential, die meerdere referentieproviders opeenvolgend probeert. De eerste verbinding kan traag zijn omdat de SDK de keten doorloopt totdat hij een werkende provider vindt. In productie, als je weet welk type inloggegevens je omgeving gebruikt, specificeer het dan direct (bijvoorbeeld ActiveDirectoryMSI voor managed identity) om de chain walk te voorkomen. Zie Microsoft Entra-verificatie voor meer informatie.

Grote datasets streamen

Voor tabellen met miljoenen rijen kan het laden van alles in één keer het geheugen uitputten. De benadering met chunks haalt rijen in batches op met fetchmany() en voegt de resultaten samen, waarbij het piekgeheugengebruik evenredig blijft aan chunksize in plaats van aan de volledige resultaatset.

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)

Generator voor grote datasets

Wanneer je data incrementeel moet verwerken zonder het hele resultaat in het geheugen te houden, gebruik dan een generator. Elk yield produceert één DataFrame-chunk die je kunt verwerken en weggooien voordat je de volgende ophaalt.

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

Schrijf DataFrames naar Microsoft SQL

Quote-identificaties om SQL-injectie te voorkomen

Tabel- en kolomnamen kunnen niet als queryparameters in SQL worden doorgegeven. Wanneer je SQL-statements met dynamische identificaties bouwt, wikkel je elke naam tussen vierkante haken en verwijder je ingebedde ] tekens om SQL-injectie te voorkomen.

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

De hulpfuncties in deze sectie gebruiken quote_id() voor alle tabel- en kolomnamen in de gegenereerde SQL.

DataFrame-rijen invoegen

De eenvoudigste aanpak itereert over de rijen van de DataFrame en genereert één INSERT per rij. De eenvoudige aanpak werkt voor kleine DataFrames, maar is traag voor grote volumes omdat elke rij een aparte retour naar de server vereist.

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

Voor grote DataFrames gebruik je de drivermethodebulkcopy(), die rijen in bulk verzendt via het TDS (Tabular Data Stream) protocol, het native wire-protocol dat Microsoft SQL gebruikt. Deze aanpak is sneller dan rij-voor-rij invoegingen omdat het heen-en-weer reizen minimaliseert.

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

Werk bestaande rijen bij vanuit DataFrame

Om rijen die al in de tabel bestaan bij te werken, doorloop je de DataFrame en voer je geparameteriseerde UPDATE-instructies uit. Hiermee key_column wordt aangegeven welke rij bijgewerkt moet worden.

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-patroon (samenvoegen)

Wanneer sommige rijen nieuw zijn en andere al bestaan, gebruik dan een SQL-instructie MERGE om in één enkele bewerking in te voegen of bij te werken. MERGE vergelijkt elke binnenkomende rij met de doeltabel met behulp van de sleutelkolommen. Als er een match wordt gevonden, wordt het bijgewerkt; anders voegt hij zich in. MERGE voorkomt dat apart wordt gecontroleerd of iets bestaat.

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

Patronen voor data-analyse

De volgende voorbeelden tonen veelvoorkomende analysetaken die Microsoft SQL-queries combineren met pandas-transformaties.

Query’s samenvoegen in DataFrame

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

Tijdreeksgegevens

Gebruik pandas datumindexering en resampling om te werken met tijdreeksgegevens uit Microsoft SQL. Om operaties zoals rollende gemiddelden en resampling mogelijk te maken, stel je de datumkolom in als de DataFrame-index.

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

Draaitabellen uit SQL-data

Draaitabellen herstructureren data van rijen naar een matrixformaat. Om het te herorganiseren volgens dimensies zoals jaar, maand en categorie, haal je de ruwe gegevens uit Microsoft SQL en gebruik je vervolgens pivot_table().

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-patronen

Om pipelines te bouwen, te transformeren en te laden, combineer je Microsoft SQL-queries met pandas-transformaties om pipelines te extraheren, transformeren en laden. De bestuurder verzorgt de extractie en het laden, terwijl pandas de transformatiestap uitvoert.

Extraheren, transformeren, laden

Dit voorbeeld haalt actieve klantgegevens uit, past bedrijfsregels toe op segmentklanten en laadt de resultaten in een bestemmingstabel.

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)

Incrementeel belastingspatroon

Voor lopende datapijplijnen, laad alleen records die sinds de laatste run zijn veranderd. Deze benadering vraagt de bestemmingstabel op voor de maximale tijdstempel en haalt vervolgens alleen nieuwere records op van de bron.

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)

Tips voor prestaties

De juiste gegevenstypen gebruiken

Pandas gebruikt standaard 64-bits types voor getallen, waardoor geheugen wordt verspild wanneer kleinere types voldoende zijn. Het downcasten van gehele getallen en floats, en het omzetten van stringkolommen met lage kardinaliteit naar categoricals, kan het geheugengebruik aanzienlijk verminderen.

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

Gebruik SQL voor zwaar werk

Microsoft SQL is sneller voor aggregaties, filtering en joins dan al je ruwe data via de kabel ophalen en lokaal verwerken in Python. Laat Microsoft SQL waar mogelijk het zware werk doen, verplaats alleen de data die je nodig hebt over het netwerk en gebruik pandas voor analyses en transformaties die in Python handiger zijn.

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

Batch-schrijfopdrachten

Voor grote DataFrames die te groot zijn voor één bulkinsert, splits het werk op in batches en houd de voortgang bij.

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