Nota
O acesso a esta página requer autorização. Pode tentar iniciar sessão ou alterar os diretórios.
O acesso a esta página requer autorização. Pode tentar alterar os diretórios.
A biblioteca pandas é a principal ferramenta de análise de dados do Python. Ao combinar pandas com o driver mssql-python, pode:
- Carregue os resultados das consultas SQL diretamente nos DataFrames.
- Escreva DataFrames de volta para Microsoft SQL de forma eficiente.
- Realizar operações ETL.
- Criar pipelines de dados.
Os exemplos deste artigo consultam a Production.Product tabela e outras tabelas na base de dados de exemplo do AdventureWorks. Os exemplos que gravam dados usam tabelas temporárias para evitar a modificação de dados de exemplo.
Outras tabelas referenciadas em exemplos de análise (Sales.SalesOrderHeader, Sales.SalesOrderDetail, Production.ProductSubcategory) fazem parte do AdventureWorks. Substitua as suas próprias tabelas ao adaptar estes padrões.
Leia dados em DataFrames
O driver mssql-python devolve linhas sob a forma de objetos Python, que pode converter em DataFrames do pandas lendo os nomes das colunas de cursor.description e os valores das linhas de fetchall(). As funções auxiliares nesta secção envolvem essa conversão em padrões reutilizáveis.
Consulta básica para DataFrame
Esta função executa uma consulta parametrizada e constrói um DataFrame a partir do conjunto completo de resultados. Funciona bem para conjuntos de resultados que cabem confortavelmente na memória.
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
Se a sua cadeia de ligação usar Authentication=ActiveDirectoryDefault, o controlador usa DefaultAzureCredential, que tenta vários fornecedores de credenciais em sequência. A primeira conexão pode ser lenta porque o SDK percorre a cadeia até encontrar um provedor funcional. Em produção, se souber que tipo de credencial o seu ambiente utiliza, especifique-o diretamente (por exemplo, ActiveDirectoryMSI para identidade gerida) para evitar o chain walk. Para obter mais informações, consulte Autenticação do Microsoft Entra.
Transmitir grandes conjuntos de dados
Para tabelas com milhões de linhas, carregar tudo de uma vez pode esgotar a memória. A abordagem em blocos obtém linhas em lotes com fetchmany() e concatena os resultados, mantendo a utilização máxima de memória proporcional a chunksize, em vez de ao 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)
Gerador para grandes conjuntos de dados
Quando precisar de processar dados de forma incremental sem manter o resultado completo na memória, use um gerador. Cada um yield produz um bloco de DataFrame que podes processar e descartar antes de buscar o próximo.
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}")
Escrever DataFrames para Microsoft SQL
Coloque os identificadores entre aspas para evitar a injeção de SQL
Nomes de tabelas e colunas não podem ser passados como parâmetros de consulta em SQL. Quando constróis instruções SQL com identificadores dinâmicos, envolve cada nome em colchetes quadrados e escapa de quaisquer caracteres embutidos ] para evitar a injeção 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}]"
As funções auxiliares nesta secção utilizam quote_id() para todos os nomes de tabelas e colunas no SQL gerado.
Inserir linhas do DataFrame
A abordagem mais simples itera sobre as linhas do DataFrame e executa uma INSERT por linha. A abordagem simples funciona para DataFrames pequenos, mas é lenta para volumes grandes porque cada linha requer uma ida e volta separada até ao 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")
Inserção em massa com BCP (recomendada para DataFrames grandes)
Para DataFrames grandes, utilize o método bulkcopy() do controlador, que envia linhas em massa através do protocolo TDS (Tabular Data Stream), o protocolo nativo de comunicação que o Microsoft SQL utiliza. Esta abordagem é mais rápida do que as inserções linha a linha porque minimiza as idas e voltas.
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")
Atualizar linhas existentes a partir do DataFrame
Para atualizar as linhas que já existem na tabela, itere sobre o DataFrame e emita instruções parametrizadas UPDATE .
key_column identifica que linha atualizar.
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")
Padrão upsert (fusão)
Quando algumas linhas podem ser novas e outras já existir, use uma instrução SQL MERGE para inserir ou atualizar numa única operação.
MERGE compara cada linha de entrada com a tabela alvo usando as colunas-chave. Se for encontrada uma correspondência, é atualizada; caso contrário, é inserida.
MERGE evita verificar a existência separadamente.
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
Padrões de análise de dados
Os exemplos seguintes mostram tarefas de análise comuns que combinam consultas SQL do Microsoft com transformações pandas.
Agregar consultas para 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())
Dados de séries temporais
Usa a indexação de datas e reamostragem do pandas para trabalhar com dados de séries temporais do Microsoft SQL. Para permitir operações como médias móveis e reamostragem, defina a coluna de data como o í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()
Tabelas dinâmicas a partir de dados SQL
As tabelas dinâmicas reorganizam dados de linhas para um formato de matriz. Para reorganizar por dimensões como ano, mês e categoria, retira os dados brutos do Microsoft SQL e depois 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)
Padrões ETL
Para construir pipelines de extração, transformação e carregamento, combine consultas do Microsoft SQL com transformações em pandas para construir pipelines de extração, transformação e carregamento. O condutor trata da extração e do carregamento enquanto os pandas tratam da etapa de transformação.
Extrair, transformar, carregar
Este exemplo extrai dados ativos de clientes, aplica regras de negócio para segmentar clientes e carrega os resultados numa tabela 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)
Padrão incremental de carga
Para pipelines de dados em curso, carregue apenas registos que mudaram desde a última execução. Esta abordagem consulta a tabela de destino para obter o carimbo temporal máximo e depois recupera apenas registos mais recentes da fonte.
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)
Sugestões de desempenho
Usar tipos de dados apropriados
O Pandas utiliza por defeito tipos de 64 bits para números, desperdiçando memória quando tipos mais pequenos são suficientes. Reduzir inteiros e floats, e converter colunas de cadeia de baixa cardinalidade em categóricas, pode reduzir significativamente o uso de memória.
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 trabalho pesado
O Microsoft SQL é mais rápido para agregações, filtragem e junções do que transferir todos os seus dados brutos pela rede e processá-los localmente em Python. Deixa o Microsoft SQL fazer o trabalho pesado sempre que possível, move apenas os dados de que precisas pela rede e usa pandas para análises e transformações que sejam mais convenientes em 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
""")
Batch escreve
Para DataFrames demasiado grandes para uma única inserção em massa, divida o trabalho em lotes e monitorize o progresso.
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}")