Používejte mssql-python s pandas

Knihovna pandas je hlavním nástrojem pro analýzu dat v Python. Kombinací pandas s ovladačem mssql-python můžete:

  • Načte výsledky SQL dotazů přímo do DataFrames.
  • Efektivně zapisujte DataFrames zpět do Microsoft SQL.
  • Provádějte ETL operace.
  • Vytvářejte datové kanály.

Příklady v tomto článku dotazují tabulku Production.Product a další tabulky v databázi AdventureWorks. Příklady, které zapisují data, používají dočasné tabulky, aby se zabránilo úpravám vzorkových dat.

Další tabulky uvedené v příkladech analýz (Sales.SalesOrderHeader, Sales.SalesOrderDetail, Production.ProductSubcategory) jsou součástí AdventureWorks. Při přizpůsobování těchto vzorů nahraďte své vlastní tabulky.

Čtěte data do DataFrames

Ovladač mssql-python vrací řádky jako objekty jazyka Python, které převedete na datové rámce pandas DataFrame načtením názvů sloupců z cursor.description a hodnot řádků z fetchall(). Pomocné funkce v této sekci přenášejí tuto konverzi do opakovaně použitelných vzorů.

Základní dotaz na DataFrame

Tato funkce vykoná parametrizovaný dotaz a vytvoří DataFrame z celé výsledné množiny. Funguje dobře pro sady výsledků, které se pohodlně vejdou do paměti.

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

Pokud váš připojovací řetězec používá Authentication=ActiveDirectoryDefault, ovladač používá DefaultAzureCredential, který zkouší více poskytovatelů přihlašovacích údajů za sebou. První spojení může být pomalé, protože SDK prochází řetězec, dokud nenajde funkčního poskytovatele. V produkci, pokud víte, jaký typ přihlašovacích údajů vaše prostředí používá, zadejte ho přímo (například ActiveDirectoryMSI pro spravovanou identitu), abyste se vyhnuli tzv. chain walk. Další informace naleznete v tématu ověřování Microsoft Entra.

Dotaz na DataFrame pomocí šipky

cursor.arrow()vrací množinu výsledků jako , pyarrow.Tablekterou pandas převede bez vytvoření Python objektu pro každou hodnotu. Názvy a typy sloupců pocházejí z výsledné množiny, takže je nepřestavujete z cursor.description:

cursor.execute("SELECT Name, ListPrice, ProductSubcategoryID FROM Production.Product")
df = cursor.arrow().to_pandas()

Aby využití paměti zůstalo úměrné velikosti dávky, streamujte pomocí arrow_reader() a zřetězujte:

cursor.execute("SELECT * FROM Production.TransactionHistory")

with cursor.arrow_reader(batch_size=50000) as reader:
    df = pd.concat([batch.to_pandas() for batch in reader], ignore_index=True)

Použijte čtecí objekt jako správce kontextu, aby se kurzor na straně serveru uvolnil i v případě, že smyčku přeruší výjimka.

Streamujte velké datové sady

U tabulek s miliony řádků může načítání všeho najednou vyčerpat paměť. Přístup po částech načítá řádky dávkově pomocí fetchmany() a zřetězí výsledky, přičemž špičkové využití paměti zůstává úměrné chunksize, nikoli celé sadě výsledků.

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)

Generátor pro velké datové sady

Když potřebujete zpracovávat data postupně, aniž byste celý výsledek uchovávali v paměti, použijte generátor. Každý z nich yield vytvoří jeden blok DataFrame, který můžete zpracovat a zahodit před načtením dalšího.

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

Zápis DataFrames do Microsoft SQL

Identifikátory uvozovek pro zabránění SQL injekci

Názvy tabulek a sloupců nelze v SQL předávat jako parametry dotazu. Když vytváříte SQL příkazy s dynamickými identifikátory, zabalte každé jméno do hranatých závorek a unikněte všem vloženým ] znacím znakům, abyste zabránili SQL injekci.

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

Pomocné funkce v této sekci používají quote_id() pro všechny názvy tabulek a sloupců ve generovaném SQL.

Vložit řádky do DataFrame

Nejjednodušší přístup prochází řádky DataFrame a pro každý řádek provede jedno INSERT. Jednoduchý přístup funguje pro malé DataFrame, ale je pomalý pro velké objemy, protože každý řádek vyžaduje samostatnou cestu tam a zpět na server.

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

cursor.bulkcopy_arrow()čte data Arrow DataFrame přímo, takže neiteruje řádky v Python. Nejprve přetypujte tabulku na cílové datové typy sloupců, protože pyarrow odvodí float64 pro číselný sloupec a ovladač nemůže tento typ namapovat na decimal, numeric nebo money:

import pyarrow as pa

def dataframe_to_sql_arrow(conn, df: pd.DataFrame, table: str, schema: pa.Schema) -> int:
    """Bulk insert a DataFrame through the Arrow path."""
    arrow_table = pa.Table.from_pandas(df, preserve_index=False).cast(schema)
    cursor = conn.cursor()
    result = cursor.bulkcopy_arrow(table, arrow_table)
    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]
})

target = pa.schema([
    pa.field("Name", pa.string()),
    pa.field("ListPrice", pa.decimal128(10, 2)),
    pa.field("ProductSubcategoryID", pa.int32()),
])

rows = dataframe_to_sql_arrow(conn, df, "##PandasProducts", target)
print(f"Bulk inserted {rows} rows")

NaN hodnoty se na této cestě stávají SQL NULL , takže je nemusíte nejdřív nahrazovat. Použijte Table.cast() místo předání schématu do Table.from_pandas(), protože Table.from_pandas() nedokáže přímo převést sloupec typu float na decimal128.

Hromadné vložení s BCP

Pokud DataFrame obsahuje objekty jazyka Python, které neodpovídají typu Arrow, použijte místo toho bulkcopy(), které přijímá řádkové n-tice:

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

Aktualizujte stávající řádky z DataFrame

Pro aktualizaci řádků, které již v tabulce existují, přejděte přes DataFrame a zadejte parametrizované UPDATE příkazy. key_column určuje, který řádek se má aktualizovat.

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

Vzor upsert (merge)

Když některé řádky mohou být nové a jiné už existují, použijte MERGE SQL příkaz k vložení nebo aktualizaci v jedné operaci. MERGE porovnává každý příchozí řádek s cílovou tabulkou pomocí klíčových sloupců. Pokud je nalezena shoda, provede aktualizaci; v opačném případě vloží záznam. MERGE eliminuje nutnost samostatně kontrolovat existenci.

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

Vzorce analýzy dat

Následující příklady ukazují běžné analytické úlohy, které kombinují Microsoft SQL dotazy s transformacemi pandas.

Agregované dotazy do 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())

Časová řada

Použijte v pandas indexování podle data a převzorkování pro práci s daty časových řad z Microsoft SQL. Pro umožnění operací jako jsou klouzavé průměry a opětovné vzorkování nastavte sloupec data jako index DataFrame.

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

Kontingenční tabulky z SQL dat

Kontingenční tabulky přetvářejí data z řádků do maticového formátu. Pro reorganizaci podle dimenzí jako rok, měsíc a kategorie si vytáhněte surová data z Microsoft SQL a použijte 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 vzory

Pro vytvoření pipeline pro extrakci, transformaci a načítání kombinujte dotazy Microsoft SQL s transformacemi pandas pro tvorbu pipeline pro extrakci, transformaci a načítání. Řidič zajišťuje extrakci a nakládání, zatímco Pandas zajišťuje krok transformace.

Extrahování, transformace, načítání

Tento příklad extrahuje aktivní data o zákaznících, aplikuje obchodní pravidla pro segmentaci zákazníků a načítá výsledky do cílové tabulky.

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)

Inkrementální zátěžový vzor

Pro průběžná datová potrubí načítejte pouze záznamy, které se změnily od posledního spuštění. Tento přístup dotazuje cílovou tabulku na maximální časové razítko a poté ze zdroje načítá pouze novější záznamy.

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)

Tipy týkající se výkonu

Použití vhodných datových typů

Pandas ve výchozím nastavení používá 64bitové typy pro čísla, která zabírají více paměti, než menší typy vyžadují. Snižujte celá čísla a plovoucí a převádějte sloupce řetězců s nízkou kardinalitou na kategorie , aby se snížila spotřeba paměti.

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

Používejte SQL pro těžké zdvíhání

Microsoft SQL je rychlejší pro agregace, filtrování a spojování než stahování všech surových dat přes kabel a zpracování lokálně v Python. Nechte Microsoft SQL dělat těžkou práci, kdykoli je to možné, přesouvejte síť jen ta data, která potřebujete, a používejte pandas pro analýzu a transformace, které jsou v Python pohodlnější.

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

Dávkové zápisy

U velkých DataFrames, které jsou příliš velké pro jedno hromadné vložení, rozdělte zpracování do dávek a sledujte postup.

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