Notatka
Dostęp do tej strony wymaga autoryzacji. Może spróbować zalogować się lub zmienić katalogi.
Dostęp do tej strony wymaga autoryzacji. Możesz spróbować zmienić katalogi.
Biblioteka pandas to główne narzędzie analizy danych w Python. Łącząc pandas ze sterownikiem mssql-python, możesz:
- Ładuj wyniki zapytań SQL bezpośrednio do DataFrames.
- Efektywnie zapisuj DataFrames z powrotem do Microsoft SQL.
- Wykonuj operacje ETL.
- Tworz potoki danych.
Przykłady w tym artykule wykonują zapytania względem tabeli Production.Product i innych tabel w przykładowej bazie danych AdventureWorks. Przykłady zapisujące dane używają tabel tymczasowych, aby uniknąć modyfikacji danych próbnych.
Inne tabele cytowane w przykładach analiz (Sales.SalesOrderHeader, Sales.SalesOrderDetail, Production.ProductSubcategory) są częścią AdventureWorks. Zastąp te tabele własnymi podczas dostosowywania tych wzorców.
Odczyt danych do DataFrames
Sterownik mssql-python zwraca wiersze jako obiekty języka Python, które można przekształcić w ramki danych pandas, odczytując nazwy kolumn z cursor.description, a wartości wierszy z fetchall(). Funkcje pomocnicze w tej sekcji opakują tę konwersję w wzorce wielokrotnego użytku.
Podstawowe zapytanie do DataFrame
Funkcja ta wykonuje parametryzowane zapytanie i buduje DataFrame z pełnego zbioru wyników. Sprawdza się dobrze w zestawach wyników, które wygodnie mieszczą się w pamięci.
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
Jeśli w parametrach połączenia użyto Authentication=ActiveDirectoryDefault, sterownik używa DefaultAzureCredential, który po kolei próbuje użyć wielu dostawców poświadczeń. Pierwsze połączenie może być wolne, ponieważ SDK przechodzi przez łańcuch, aż znajdzie dostawcę, który działa. W środowisku produkcyjnym, jeśli wiesz, jakiego typu poświadczeń używa środowisko, wskaż go bezpośrednio (na przykład ActiveDirectoryMSI w przypadku tożsamości zarządzanej), aby uniknąć przechodzenia przez łańcuch. Aby uzyskać więcej informacji, zobacz Microsoft Entra authentication (Uwierzytelnianie w usłudze Microsoft Entra).
Strumieniowanie dużych zbiorów danych
W tabelach z milionami wierszy ładowanie wszystkiego naraz może wyczerpać pamięć. Podejście oparte na porcjowaniu pobiera wiersze partiami przy użyciu fetchmany() i łączy wyniki, utrzymując szczytowe zużycie pamięci na poziomie proporcjonalnym do chunksize, a nie do pełnego zbioru wyników.
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 dla dużych zbiorów danych
Gdy musisz przetwarzać dane stopniowo, bez przechowywania całego wyniku w pamięci, użyj generatora. Każdy yield z nich generuje jeden fragment DataFrame, który można przetworzyć i odrzucić przed pobraniem kolejnego.
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}")
Zapisuj DataFrames do Microsoft SQL
Identyfikatory cytatów, aby zapobiec wstrzykiwaniu SQL
Nazwy tabel i kolumn nie mogą być przekazywane jako parametry zapytań w SQL. Gdy budujesz instrukcje SQL z dynamicznymi identyfikatorami, otul każdą nazwę nawiasem kwadratowym i uciekaj przed osadzonymi ] znakami, aby zapobiec wstrzykiwaniu SQL.
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}]"
Funkcje pomocnicze w tej sekcji używają elementu quote_id() dla wszystkich nazw tabel i kolumn w wygenerowanym kodzie SQL.
Wstawianie wierszy DataFrame
Najprostsze podejście iteruje po wierszach obiektu DataFrame i wykonuje jedno INSERT na wiersz. To proste podejście sprawdza się w przypadku małych obiektów DataFrame, ale jest wolne przy dużych ilościach danych, ponieważ każdy wiersz wymaga oddzielnego cyklu komunikacji z serwerem.
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")
Wstawianie zbiorcze za pomocą BCP (zalecane w przypadku dużych ramek danych)
W przypadku dużych ramek danych DataFrame użyj metody bulkcopy() sterownika, która wysyła wiersze zbiorczo za pośrednictwem protokołu TDS (Tabular Data Stream), natywnego protokołu komunikacyjnego używanego przez Microsoft SQL Server. To podejście jest szybsze niż wstawianie wiersz po wierszu, ponieważ minimalizuje liczbę połączeń tam i z powrotem.
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")
Zaktualizuj istniejące wiersze z DataFrame
Aby zaktualizować wiersze, które już istnieją w tabeli, iteruj po obiekcie DataFrame i wykonuj parametryzowane instrukcje UPDATE.
key_column wskazuje, który wiersz należy zaktualizować.
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")
Wzorzec upsert (merge)
Gdy niektóre wiersze mogą być nowe, a inne już istnieją, użyj instrukcji SQL MERGE do wstawienia lub aktualizacji w jednej operacji.
MERGE porównuje każdy przychodzący wiersz z tabelą docelową za pomocą kolumn kluczowych. Jeśli zostanie znalezione dopasowanie, dane zostaną zaktualizowane; w przeciwnym razie zostanie wstawiony nowy rekord.
MERGE pozwala uniknąć osobnego sprawdzania, czy coś istnieje.
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
Wzorce analizy danych
Poniższe przykłady pokazują typowe zadania analityczne łączące zapytania Microsoft SQL z transformacjami pandas.
Zapytania agregujące 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())
Dane szeregów czasowych
Użyj indeksowania dat i ponownego próbkowania w bibliotece pandas do pracy z danymi szeregów czasowych z Microsoft SQL. Aby umożliwić operacje takie jak średnie kroczące i ponowne próbkowanie, ustaw kolumnę daty jako indeks 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()
Tabele przestawne z danych SQL
Tabele przestawne przekształcają dane z wierszy w formę macierzy. Aby zorganizować je według wymiarów takich jak rok, miesiąc i kategoria, pobierz surowe dane z Microsoft SQL, a następnie użyj 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)
Wzorce ETL
Aby budować potoki do ekstrakcji, transformacji i ładowania, połącz zapytania Microsoft SQL z transformacjami pandas, aby budować potoki ekstrakcji, transformacji i ładowania. Kierowca zajmuje się wyciąganiem i ładowaniem, podczas gdy Pandas zajmuje się etapem transformacji.
Wyodrębnianie, przekształcanie, ładowanie
Ten przykład wyodrębnia dane aktywnych klientów, stosuje reguły biznesowe do segmentacji klientów i ładuje wyniki do tabeli docelowej.
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)
Wzorzec obciążenia przyrostowego
W przypadku bieżących potoków danych ładuj tylko rekordy, które zmieniły się od ostatniego uruchomienia. To podejście zapytuje tabelę docelową o maksymalną liczbę znaczników czasu, a następnie pobiera tylko nowsze rekordy ze źródła.
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)
Porady dotyczące wydajności
Używanie odpowiednich typów danych
Pandas domyślnie stosuje 64-bitowe typy liczb, marnując pamięć, gdy wystarczają mniejsze typy. Konwersja liczb całkowitych i liczb zmiennoprzecinkowych na typy o mniejszym rozmiarze oraz konwersja kolumn tekstowych o niskiej kardynalności na typ kategoryczny mogą znacząco zmniejszyć zużycie pamięci.
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
Używaj SQL do ciężkiego podnoszenia
Microsoft SQL jest szybszy w agregacjach, filtrowaniu i połączeniach niż pobieranie wszystkich surowych danych przez kabel i przetwarzanie ich lokalnie w Python. Pozwól Microsoft SQL wykonać ciężką pracę, kiedy tylko to możliwe, przesuwaj tylko potrzebne dane przez sieć i używaj Pandas do analiz i transformacji, które są wygodniejsze w 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
""")
Zapisy wsadowe
Dla dużych DataFrame'ów, które są zbyt duże na pojedynczy wkład zbiorczy, podziel pracę na partie i śledź postępy.
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}")