Poznámka:
Přístup k této stránce vyžaduje autorizaci. Můžete se zkusit přihlásit nebo změnit adresáře.
Přístup k této stránce vyžaduje autorizaci. Můžete zkusit změnit adresáře.
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.
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")
Hromadné vkládání pomocí BCP (doporučeno pro velké datové rámce)
Pro velké DataFrames použijte metodu ovladačebulkcopy(), která posílá řádky hromadně přes protokol TDS (Tabular Data Stream), nativní wire protokol, který používá Microsoft SQL. Tento přístup je rychlejší než vkládání řádek po řádku, protože minimalizuje cesty tam a zpět.
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, což snižuje paměť, když stačí menší typy. Snižování celých čísel a plovoucích čísel a převod sloupců řetězců s nízkou kardinálností na kategorické může výrazně snížit využití 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}")