Remarque
L’accès à cette page nécessite une autorisation. Vous pouvez essayer de vous connecter ou de modifier des répertoires.
L’accès à cette page nécessite une autorisation. Vous pouvez essayer de modifier des répertoires.
La bibliothèque pandas est l'outil principal d'analyse de données de Python. En combinant pandas avec le pilote mssql-python, vous pouvez :
- Chargez directement les résultats des requêtes SQL dans DataFrames.
- Réécrire efficacement des DataFrames dans Microsoft SQL Server.
- Effectuer des opérations ETL.
- Créer des pipelines de données.
Les exemples de cet article interrogent la Production.Product table et d’autres tables de la base de données d’exemple AdventureWorks. Les exemples qui écrivent des données utilisent des tables temporaires pour éviter de modifier les données d’échantillon.
D’autres tableaux référencés dans les exemples d’analyse (Sales.SalesOrderHeader, Sales.SalesOrderDetail, Production.ProductSubcategory) font partie d’AdventureWorks. Remplacez ces modèles par vos propres tables lorsque vous les adaptez.
Lire les données dans DataFrames
Le pilote mssql-python retourne les lignes sous forme d’objets Python, que vous convertissez en DataFrames pandas en lisant les noms des colonnes à partir de cursor.description et les valeurs des lignes à partir de fetchall(). Les fonctions d’assistance dans cette section enveloppent cette conversion en motifs réutilisables.
Requête de base vers DataFrame
Cette fonction exécute une requête paramétrée et construit un DataFrame à partir de l’ensemble complet des résultats. Cela fonctionne bien pour des ensembles de résultats qui tiennent aisément en mémoire.
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
Si votre chaîne de connexion utilise Authentication=ActiveDirectoryDefault, le pilote utilise DefaultAzureCredential, ce qui essaie plusieurs fournisseurs d’identifiants en séquence. La première connexion peut être lente car le SDK parcourt la chaîne jusqu’à ce qu’il trouve un fournisseur fonctionnel. En production, si vous savez quel type d’identifiant votre environnement utilise, spécifiez-le directement (par exemple, ActiveDirectoryMSI pour l’identité gérée) afin d’éviter la marche en chaîne. Pour plus d’informations, consultez Authentification Microsoft Entra.
Flux de grands ensembles de données
Pour des tables de millions de lignes, tout charger en même temps peut épuiser la mémoire. L’approche par blocs récupère les lignes par lots avec fetchmany() et concatène les résultats, en maintenant une utilisation maximale de la mémoire proportionnelle à chunksize plutôt qu’à l’ensemble complet des résultats.
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)
Générateur pour de grands ensembles de données
Lorsque vous devez traiter des données de manière incrémentale sans garder l’ensemble du résultat en mémoire, utilisez un générateur. Chacun yield produit un bloc de DataFrame que vous pouvez traiter et supprimer avant de récupérer le suivant.
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}")
Écrire DataFrames vers Microsoft SQL
Mettre les identifiants entre guillemets pour éviter l’injection SQL
Les noms de tables et de colonnes ne peuvent pas être passés comme paramètres de requête en SQL. Lorsque vous construisez des instructions SQL avec des identifiants dynamiques, enroulez chaque nom entre crochets et éliminez les caractères intégrés ] pour éviter l’injection 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}]"
Les fonctions d’assistance de cette section utilisent quote_id() pour tous les noms de table et de colonne dans le SQL généré.
Insérer des lignes de DataFrame
L’approche la plus simple parcourt les lignes du DataFrame et génère un INSERT par ligne. L’approche simple fonctionne pour les petits DataFrames mais est lente pour les gros volumes car chaque ligne nécessite un aller-retour séparé vers le serveur.
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")
Insertion en bloc avec BCP (recommandée pour les grands cadres de données)
Pour les DataFrames volumineux, utilisez la méthode bulkcopy() du pilote, qui envoie les lignes en bloc via le protocole TDS (Tabular Data Stream), le protocole réseau natif utilisé par Microsoft SQL. Cette approche est plus rapide que les insertions ligne par ligne, car elle minimise les allers-retours.
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")
Mettre à jour les lignes existantes depuis DataFrame
Pour mettre à jour les lignes déjà présentes dans le tableau, itérez sur le DataFrame et émettez des instructions paramétrées UPDATE . Le key_column identifie la ligne à mettre à jour.
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")
Modèle upsert (mise à jour ou insertion)
Lorsque certaines lignes peuvent être nouvelles et que d’autres existent déjà, utilisez une instruction SQL MERGE pour insérer ou mettre à jour dans une seule opération.
MERGE compare chaque ligne entrante à la table cible en utilisant les colonnes clés. Si une correspondance est trouvée, elle est mise à jour ; sinon, elle est insérée.
MERGE évite de vérifier l’existence au préalable.
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
Schémas d’analyse des données
Les exemples suivants montrent des tâches d’analyse courantes qui combinent des requêtes Microsoft SQL avec des transformations pandas.
Agréger les requêtes vers 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())
Données de série chronologique
Utilisez l’indexation et le rééchantillonnage de dates de Pandas pour travailler avec les données de séries temporelles de Microsoft SQL. Pour permettre des opérations comme les moyennes mobiles et le rééchantillonnage, définissez la colonne date comme l’indice 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()
Tables croisés dynamiques à partir de données SQL
Les tableaux croisés dynamiques remodelent les données des lignes en format matriciel. Pour réorganiser par dimensions comme année, mois et catégorie, extrais les données brutes de Microsoft SQL, puis utilise 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)
Motifs ETL
Pour construire des pipelines d’extraction, transformation et chargement, combinez les requêtes Microsoft SQL avec les transformations pandas pour construire des pipelines d’extraction, de transformation et de chargement. Le conducteur s’occupe de l’extraction et du chargement tandis que Pandas s’occupe de l’étape de transformation.
Extraire, transformer, charger
Cet exemple extrait les données clients actifs, applique des règles métier pour segmenter les clients, et charge les résultats dans une table de destination.
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)
Schéma de charge incrémental
Pour les pipelines de données en cours, chargez uniquement les enregistrements qui ont changé depuis la dernière exécution. Cette approche interroge la table de destination pour obtenir le timestamp maximal, puis récupère uniquement les enregistrements plus récents depuis la source.
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)
Astuces pour les performances
Utiliser les types de données appropriés
Pandas utilise par défaut des types 64 bits pour les nombres, gaspillant de la mémoire lorsque les types plus petits suffisent. Réduire les entiers et floats, ainsi que convertir les colonnes de chaînes de faible cardinalité en catégoriques, peut réduire considérablement la consommation de mémoire.
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
Utilisez SQL pour le travail lourd
Microsoft SQL est plus rapide pour les agrégations, le filtrage et les jointures que de récupérer toutes vos données brutes via le réseau et de les traiter localement en Python. Laissez Microsoft SQL faire le gros du travail dès que possible, ne déplacez que les données dont vous avez besoin sur le réseau et utilisez pandas pour des analyses et transformations plus pratiques en 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
""")
Écriture par lots
Pour les DataFrames trop volumineux pour une seule insertion en masse, divisez le travail en lots et suivez la progression.
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}")