Verwenden Sie mssql-python mit pandas

Die pandas-Bibliothek ist Pythons primäres Datenanalysetool. Indem Sie pandas mit dem mssql-python-Treiber kombinieren, können Sie:

  • Laden Sie SQL-Abfrageergebnisse direkt in DataFrames.
  • Schreiben Sie DataFrames effizient zurück zu Microsoft SQL.
  • Führen Sie ETL-Operationen durch.
  • Erstellen Sie Datenpipelines.

Die Beispiele in diesem Artikel führen Abfragen für die Production.Product Tabelle und andere Tabellen in der AdventureWorks-Beispieldatenbank aus. Beispiele, die Daten schreiben, verwenden temporäre Tabellen, um Beispieldaten nicht zu verändern.

Weitere in Analysebeispielen referenzierte Tabellen (Sales.SalesOrderHeader, Sales.SalesOrderDetail, Production.ProductSubcategory) sind Teil von AdventureWorks. Ersetzen Sie Ihre eigenen Tabellen, wenn Sie diese Muster anpassen.

Daten in DataFrames einlesen

Der mssql-python-Treiber gibt Zeilen als Python-Objekte zurück, die Sie in pandas-DataFrames umwandeln, indem Sie Spaltennamen aus cursor.description und Zeilenwerte aus fetchall() lesen. Die Hilfsfunktionen in diesem Abschnitt packen diese Umwandlung in wiederverwendbare Muster.

Grundanfrage an DataFrame

Diese Funktion führt eine parametrisierte Abfrage aus und erstellt einen DataFrame aus dem vollständigen Ergebnisset. Es ist gut geeignet für Ergebnismengen, die problemlos in den Arbeitsspeicher 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())

Note

Wenn Ihre Verbindungszeichenfolge Authentication=ActiveDirectoryDefault verwendet, nutzt der Treiber DefaultAzureCredential, das mehrere Anmeldeinformationsanbieter nacheinander ausprobiert. Die erste Verbindung kann langsam sein, weil das SDK die Kette durchläuft, bis es einen funktionierenden Anbieter findet. In der Produktion gilt: Wenn du weißt, welchen Zugangsdatentyp deine Umgebung verwendet, gib ihn direkt an (zum Beispiel ActiveDirectoryMSI für eine verwaltete Identität), um den Chain Walk zu vermeiden. Weitere Informationen finden Sie unter Microsoft Entra-Authentifizierung.

Große Datensätze streamen

Bei Tabellen mit Millionen von Zeilen kann das Laden von allem auf einmal den Speicher verbrauchen. Der Ansatz mit Chunking ruft Zeilen mit fetchmany() stapelweise ab und fügt die Ergebnisse zusammen, wobei die maximale Speichernutzung proportional zu chunksize statt zum vollständigen Ergebnissatz bleibt.

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 für große Datensätze

Wenn Sie Daten schrittweise verarbeiten müssen, ohne das gesamte Ergebnis im Speicher zu speichern, verwenden Sie einen Generator. Jeder yield erzeugt einen DataFrame-Chunk, den du verarbeiten und verwerfen kannst, bevor du den nächsten abrufst.

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

DataFrames in Microsoft SQL schreiben

Quote-Identifikatoren zur Verhinderung von SQL-Injektion

Tabellen- und Spaltennamen können in SQL nicht als Abfrageparameter übermittelt werden. Wenn Sie SQL-Anweisungen mit dynamischen Identifikatoren erstellen, packen Sie jeden Namen in eckigen Klammern und vermeiden Sie eingebettete ] Zeichen, um SQL-Injektionen zu verhindern.

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

Die Hilfsfunktionen in diesem Abschnitt verwenden quote_id() für alle Tabellen- und Spaltennamen im generierten SQL.

DataFrame-Zeilen einfügen

Der einfachste Ansatz iteriert über DataFrame-Zeilen und gibt pro Zeile eine INSERT aus. Der einfache Ansatz funktioniert für kleine DataFrames, ist aber für große Volumina langsam, da jede Zeile eine separate Rundreise zum Server benötigt.

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

Für große DataFrames verwenden Sie die Treibermethodebulkcopy(), die Zeilen in großen Mengen über das TDS-Protokoll (Tabular Data Stream) sendet, das native Wire-Protokoll, das Microsoft SQL verwendet. Dieser Ansatz ist schneller als zeilenweise Einfügungen, da er die Anzahl der Roundtrips minimiert.

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

Aktualisieren Sie bestehende Zeilen aus DataFrame

Um bereits vorhandene Zeilen in der Tabelle zu aktualisieren, iterieren Sie über das DataFrame und geben Sie parametrisierte UPDATE Anweisungen aus. key_column bestimmt, welche Zeile aktualisiert wird.

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-(Zusammenführungs-)Muster

Wenn einige Zeilen neu sind und andere bereits existieren, verwenden Sie eine SQL-Anweisung MERGE , um in einer einzigen Operation einzufügen oder zu aktualisieren. MERGE vergleicht jede eingehende Zeile mit der Zieltabelle anhand der Schlüsselspalten. Wenn eine Übereinstimmung gefunden wird, wird aktualisiert; andernfalls wird eingefügt. MERGE vermeidet, separat zu prüfen, ob etwas vorhanden ist.

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

Datenanalysemuster

Die folgenden Beispiele zeigen gängige Analyseaufgaben, die Microsoft SQL-Abfragen mit Pandas-Transformationen kombinieren.

Aggregierte Abfragen zu 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())

Zeitreihendaten

Verwenden Sie Pandas Datumsindexierung und -Resampling, um mit Zeitreihendaten aus Microsoft SQL zu arbeiten. Um Operationen wie gleitende Durchschnitte und Resampling zu ermöglichen, legen Sie die Datumsspalte als Index des DataFrame fest.

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

Pivot-Tabellen aus SQL-Daten

Pivot-Tabellen formen Daten von Zeilen in ein Matrixformat um. Um es nach Dimensionen wie Jahr, Monat und Kategorie neu zu organisieren, ziehe die Rohdaten aus Microsoft SQL und benutze pivot_table()dann .

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

Um Pipelines zu extrahieren, zu transformieren und zu laden, kombinieren Sie Microsoft-SQL-Abfragen mit Pandas-Transformationen, um Pipelines zu extrahieren, zu transformieren und zu laden. Der Fahrer übernimmt die Extraktion und das Beladen, während Pandas den Transformationsschritt übernimmt.

Extrahieren, Transformieren, Laden

Dieses Beispiel extrahiert aktive Kundendaten, wendet Geschäftsregeln auf segmentierte Kunden an und lädt die Ergebnisse in eine Zieltabelle.

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)

Inkrementelles Lastmuster

Für laufende Datenpipelines laden Sie nur Datensätze, die sich seit dem letzten Durchlauf geändert haben. Dieser Ansatz fragt die Zieltabelle nach dem maximalen Zeitstempel ab und holt dann nur neuere Datensätze von der Quelle.

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)

Leistungstipps

Verwenden geeigneter Datentypen

Pandas verwendet standardmäßig 64-Bit-Typen für Zahlen und verschwendet Speicher, wenn kleinere Typen ausreichen. Das Herunterwerfen von Ganzzahlen und Floatzahlen sowie das Umwandeln von String-Spalten mit niedriger Kardinalitität in Kategorien kann den Speicherbedarf erheblich reduzieren.

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

Nutze SQL für schwere Arbeit

Microsoft SQL ist schneller für Aggregationen, Filtern und Joins, als alle Rohdaten über die Leitung zu ziehen und lokal in Python zu verarbeiten. Lass Microsoft SQL wann immer möglich die schwere Arbeit übernehmen, trage nur die benötigten Daten über das Netzwerk und nutze Pandas für Analysen und Transformationen, die in Python bequemer sind.

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

Stapel-Schreibvorgänge

Bei großen DataFrames, die für einen einzelnen Bulk-Insert zu groß sind, teilen Sie die Verarbeitung in Batches auf und verfolgen Sie den Fortschritt.

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