Kommentar
Åtkomst till den här sidan kräver auktorisering. Du kan prova att logga in eller ändra kataloger.
Åtkomst till den här sidan kräver auktorisering. Du kan prova att ändra kataloger.
Pandas-biblioteket är Python:s primära verktyg för dataanalys. Genom att kombinera pandas med mssql-python-drivrutinen kan du:
- Ladda SQL-frågeresultat direkt i DataFrames.
- Skriv DataFrames tillbaka till Microsoft SQL effektivt.
- Utför ETL-operationer.
- Skapa datapipelines.
Exemplen i denna artikel använder frågor mot tabellen Production.Product och andra tabeller i exempeldatabasen AdventureWorks. Exempel som skriver data använder temporära tabeller för att undvika att ändra exempeldata.
Andra tabeller som refereras till i analysexempel (Sales.SalesOrderHeader, Sales.SalesOrderDetail, ) Production.ProductSubcategoryingår i AdventureWorks. Ersätt tabellerna med dina egna när du anpassar dessa mönster.
Läs in data i DataFrames
mssql-python-drivrutinen returnerar rader som Python-objekt, vilka du konverterar till pandas DataFrames genom att läsa kolumnnamn från cursor.description och radvärden från fetchall(). Hjälpfunktionerna i detta avsnitt omvandlar den omvandlingen till återanvändbara mönster.
Grundläggande fråga till DataFrame
Denna funktion kör en parameteriserad fråga och bygger en DataFrame från hela resultatuppsättningen. Det fungerar bra för resultatuppsättningar som får plats i minnet utan problem.
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
Om din anslutningssträng använder Authentication=ActiveDirectoryDefault, använder drivrutinen DefaultAzureCredential, som försöker med flera autentiseringsleverantörer i följd. Den första anslutningen kan vara långsam eftersom SDK:n går igenom kedjan tills den hittar en fungerande leverantör. I produktion, om du vet vilken typ av behörighet din miljö använder, ange det direkt (till exempel ActiveDirectoryMSI för managed identity) för att undvika kedjevandring. Mer information finns i Microsoft Entra-autentisering.
Strömma stora datamängder
För tabeller med miljontals rader kan det ta slut på minnet att ladda allt på en gång. Den chunkade metoden hämtar rader i batcher med fetchmany() och sammanfogar resultaten, så att toppminnesanvändningen är proportionell mot chunksize snarare än hela resultatuppsättningen.
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 stora datamängder
När du behöver bearbeta data inkrementellt utan att hålla hela resultatet i minnet, använd en generator. Varje yield del producerar en DataFrame-chunk som du kan bearbeta och kassera innan du hämtar nästa.
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}")
Skriv DataFrames till Microsoft SQL
Citatidentifierare för att förhindra SQL-injektion
Tabell- och kolumnnamn kan inte skickas som frågeparametrar i SQL. När du bygger SQL-satser med dynamiska identifierare, slå in varje namn inom hakparenteser och undvik inbäddade ] tecken för att förhindra SQL-injektion.
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}]"
Hjälpfunktionerna i detta avsnitt används quote_id() för alla tabell- och kolumnnamn i genererad SQL.
Infoga rader i DataFrame
Det enklaste tillvägagångssättet itererar över DataFrame-rader och ger en INSERT per rad. Den enkla metoden fungerar för små DataFrames men är långsam för stora volymer eftersom varje rad kräver en separat rundresa till servern.
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")
Massinfogning med BCP (rekommenderas för stora DataFrames)
För stora DataFrames använder du drivrutinsmetodenbulkcopy(), som skickar rader i bulk över TDS (Tabular Data Stream)-protokollet, det inbyggda trådprotokollet som Microsoft SQL använder. Denna metod är snabbare än rad-för-rad-insättningar eftersom den minimerar rundturer.
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")
Uppdatera befintliga rader från DataFrame
För att uppdatera rader som redan finns i tabellen, iterera över DataFrame och ge parametriserade UPDATE satser.
key_column identifierar vilken rad som ska uppdateras.
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-mönster (sammanfogning)
När vissa rader kan vara nya och andra redan finns, använd en SQL-sats MERGE för att infoga eller uppdatera i en enda operation.
MERGE jämför varje inkommande rad med måltabellen med hjälp av nyckelkolumnerna. Om en träff hittas uppdateras den, annars infogas den.
MERGE undviker att kontrollera om något existerar separat.
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
Dataanalysmönster
Följande exempel visar vanliga analysuppgifter som kombinerar Microsoft SQL-frågor med pandas-transformationer.
Aggregerade frågor till 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())
Tidsseriedata
Använd datumindexering i pandas och omsampling för att arbeta med tidsseriedata från Microsoft SQL. För att möjliggöra operationer som rullande medelvärden och omprovning, sätt datumkolumnen som 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()
Pivottabeller från SQL-data
Pivottabeller omformar data från rader till ett matrisformat. För att omorganisera den efter dimensioner som år, månad och kategori, hämta rådata från Microsoft SQL och använd pivot_table()sedan .
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-mönster
För att bygga pipelines för extrahering, transformering och inläsning kan du kombinera Microsoft SQL-frågor med pandas-transformationer. Föraren sköter uthämtning och lastning medan pandas sköter omvandlingssteget.
Extrahera, transformera, läsa in
Detta exempel extraherar aktiv kunddata, tillämpar affärsregler på segmentkunder och laddar in resultaten i en destinationstabell.
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)
Inkrementellt lastmönster
För pågående datapipelines, ladda endast poster som ändrats sedan senaste körning. Denna metod frågar destinationstabellen efter maximal tidsstämpel och hämtar sedan endast nyare poster från källan.
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)
Prestandatips
Använda lämpliga datatyper
Pandas använder som standard 64-bitars typer för siffror, vilket slösar minne när mindre typer räcker. Nedtoning av heltal och flyttal, samt att konvertera strängkolumner med låg kardinalitet till kategoriska, kan avsevärt minska minnesanvändningen.
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
Använd SQL för tungt arbete
Microsoft SQL är snabbare vid aggregering, filtrering och sammanslagningar än att hämta all din rådata över nätverket och bearbeta den lokalt i Python. Låt Microsoft SQL göra det tunga arbetet när det är möjligt, flytta bara den data du behöver över nätverket och använd pandas för analyser och transformationer som är mer bekväma i Python.
# 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
""")
Satsvisa skrivningar
För stora DataFrames som är för stora för en enda bulkinsats, dela upp arbetet i batcher och följ framsteg.
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}")