Nota:
El acceso a esta página requiere autorización. Puede intentar iniciar sesión o cambiar directorios.
El acceso a esta página requiere autorización. Puede intentar cambiar los directorios.
La biblioteca pandas es la principal herramienta de análisis de datos de Python. Combinando pandas con el controlador mssql-python, puedes:
- Carga los resultados de las consultas SQL directamente en DataFrames.
- Escribe DataFrames de nuevo en Microsoft SQL de forma eficiente.
- Realizar operaciones ETL.
- Crear canalizaciones de datos.
Los ejemplos de este artículo consultan la Production.Product tabla y otras tablas en la base de datos de ejemplo de AdventureWorks. Los ejemplos que escriben datos utilizan tablas temporales para evitar modificar los datos de muestra.
Otras tablas referenciadas en ejemplos de análisis (Sales.SalesOrderHeader, Sales.SalesOrderDetail, Production.ProductSubcategory) forman parte de AdventureWorks. Sustituye tus propias tablas al adaptar estos patrones.
Lee datos en DataFrames
El controlador mssql-python devuelve las filas como objetos de Python, que conviertes en DataFrames de pandas leyendo los nombres de las columnas de cursor.description y los valores de las filas de fetchall(). Las funciones auxiliares de esta sección envuelven esa conversión en patrones reutilizables.
Consulta básica a DataFrame
Esta función ejecuta una consulta parametrizada y construye un DataFrame a partir del conjunto completo de resultados. Funciona bien para conjuntos de resultados que encajan cómodamente en la memoria.
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 su cadena de conexión usa Authentication=ActiveDirectoryDefault, el controlador usa DefaultAzureCredential, que intenta usar varios proveedores de credenciales en secuencia. La primera conexión puede ser lenta porque el SDK recorre la cadena hasta encontrar un proveedor que funcione. En producción, si sabes qué tipo de credencial utiliza tu entorno, especifícala directamente (por ejemplo, ActiveDirectoryMSI para identidad gestionada) para evitar el recorrido en cadena. Para más información, consulte Autenticación de Microsoft Entra.
Transmitir grandes conjuntos de datos
Para tablas con millones de filas, cargar todo a la vez puede agotar la memoria. El enfoque por fragmentos obtiene filas en lotes con fetchmany() y concatena los resultados, manteniendo el uso máximo de memoria proporcional a chunksize en lugar de al conjunto completo de resultados.
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)
Generador para grandes conjuntos de datos
Cuando necesites procesar datos de forma incremental sin tener el resultado completo en memoria, usa un generador. Cada uno yield produce un bloque de DataFrame que puedes procesar y descartar antes de recuperar el siguiente.
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}")
Escribe DataFrames en Microsoft SQL
Entrecomille los identificadores para evitar inyección SQL
Los nombres de tablas y columnas no pueden pasarse como parámetros de consulta en SQL. Cuando construyas sentencias SQL con identificadores dinámicos, envuelve cada nombre entre corchetes y evita cualquier carácter incrustado ] para evitar la inyección 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}]"
Las funciones auxiliares de esta sección utilizan quote_id() para todos los nombres de tablas y columnas en el SQL generado.
Insertar filas de DataFrame
El enfoque más sencillo itera sobre las filas del DataFrame y genera un INSERT por fila. El enfoque sencillo funciona para DataFrames pequeños, pero es lento para volúmenes grandes porque cada fila requiere un viaje de ida y vuelta separado al servidor.
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")
Inserción masiva con BCP (recomendada para DataFrames grandes)
Para DataFrames grandes, utiliza el método bulkcopy() del controlador, que envía filas de forma masiva a través del protocolo TDS (Tabular Data Stream), el protocolo nativo de comunicación que usa Microsoft SQL. Este enfoque es más rápido que los insertos fila por fila porque minimiza los viajes de ida y vuelta.
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")
Actualizar filas existentes desde DataFrame
Para actualizar las filas que ya existen en la tabla, itere sobre el DataFrame y ejecute instrucciones parametrizadas UPDATE. El key_column indica qué fila se debe actualizar.
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")
Patrón de upsert (fusión)
Cuando algunas filas puedan ser nuevas y otras ya existan, usa una instrucción SQL MERGE para insertar o actualizar en una sola operación.
MERGE compara cada fila entrante con la tabla objetivo usando las columnas clave. Si se encuentra una coincidencia, se actualiza; de lo contrario, se inserta.
MERGE evita verificar su existencia por separado.
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
Patrones de análisis de datos
Los siguientes ejemplos muestran tareas de análisis comunes que combinan consultas de Microsoft SQL con transformaciones pandas.
Agregar consultas a 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())
Datos de serie temporal
Utiliza la indexación de fechas y el remuestreo de Pandas para trabajar con datos de series temporales de Microsoft SQL. Para permitir operaciones como promedios móviles y remuestreo, establece la columna de fecha como índice 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()
Tablas dinámicas a partir de datos SQL
Las tablas dinámicas reorganizan los datos de filas en una matriz. Para reorganizarlo por dimensiones como año, mes y categoría, extrae los datos en bruto de Microsoft SQL y luego usa 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)
Patrones ETL
Para construir pipelines de extracción, transformación y carga, combina consultas SQL de Microsoft con transformaciones pandas para construir pipelines de extracción, transformación y carga. El conductor se encarga de la extracción y la carga mientras que los pandas se encargan del paso de transformación.
Extracción, transformación y carga
Este ejemplo extrae datos activos de clientes, aplica reglas de negocio a segmentar clientes y carga los resultados en una tabla de destinos.
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)
Patrón incremental de carga
Para canalizaciones de datos en curso, solo carga los registros que cambiaron desde la última ejecución. Este enfoque consulta la tabla de destino para obtener la marca de tiempo máxima y luego recupera solo los registros más recientes de la fuente.
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)
Consejos de rendimiento
Uso del tipo de datos adecuado
Pandas por defecto utiliza tipos de 64 bits para números, desperdiciando memoria cuando los tipos más pequeños son suficientes. Reducir enteros y flotadores, y convertir columnas de cadenas de baja cardinalidad en categorías puede reducir significativamente el uso de memoria.
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
Usa SQL para el trabajo pesado
Microsoft SQL es más rápido para agregaciones, filtrado y uniones que extraer todos tus datos en bruto por cable y procesarlos localmente en Python. Deja que Microsoft SQL haga el trabajo duro siempre que sea posible, mueve solo los datos que necesitas por la red y usa pandas para análisis y transformaciones que sean más cómodos 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
""")
Escribe por lotes
Para los DataFrames que sean demasiado grandes para una sola inserción masiva, divide el trabajo en lotes y realiza un seguimiento del progreso.
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}")