Используйте mssql-python с pandas

Библиотека pandas — основной инструмент анализа данных Python. Комбинируя pandas с драйвером mssql-python, вы можете:

  • Загрузите результаты SQL-запроса напрямую в DataFrames.
  • Эффективно записывайте объекты DataFrame обратно в Microsoft SQL.
  • Выполняйте операции ETL.
  • Создавайте конвейеры данных.

В примерах в этой статье выполняются запросы к таблице Production.Product и другим таблицам в примерной базе данных AdventureWorks. Примеры, записывающие данные, используют временные таблицы, чтобы избежать изменения выборочных данных.

Другие таблицы, на которые ссылаются в примерах анализа (Sales.SalesOrderHeader, Sales.SalesOrderDetail, Production.ProductSubcategory), являются частью AdventureWorks. Подставляйте собственные таблицы при адаптации этих шаблонов.

Считывайте данные в датафреймы

Драйвер mssql-python возвращает строки как объекты Python, которые вы конвертируете в pandas DataFrames, считывая имена столбцов из cursor.description и значения строк из fetchall(). Вспомогательные функции в этом разделе оборачивают это преобразование в повторно используемые паттерны.

Базовый запрос к DataFrame

Эта функция выполняет параметризованный запрос и строит DataFrame из полного набора результатов. Это хорошо работает для наборов результатов, которые удобно помещаются в память.

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())

Замечание

Если в строке подключения используется Authentication=ActiveDirectoryDefault, драйвер использует DefaultAzureCredential, который поочередно проверяет несколько поставщиков учетных данных. Первое соединение может быть медленным, потому что SDK идёт по цепочке, пока не найдёт работающего провайдера. В продакшене, если вы знаете, какой тип учетных данных использует ваша среда, укажите его напрямую (например, ActiveDirectoryMSI для управляемой идентичности), чтобы избежать цепной ходьбы. Дополнительные сведения см. в разделе проверки подлинности Microsoft Entra.

Поток больших наборов данных

Для таблиц с миллионами строк загрузка всего сразу может исчерпать память. Подход с обработкой по частям извлекает строки пакетами с помощью fetchmany() и объединяет результаты, при этом пиковое потребление памяти остаётся пропорциональным chunksize, а не всему набору результатов.

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)

Генератор больших наборов данных

Когда нужно обрабатывать данные постепенно, не храня весь результат в памяти, используйте генератор. Каждый yield создаёт один фрагмент DataFrame, который можно обработать и отбросить перед получением следующего.

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}")

Запись фреймов данных в Microsoft SQL

Идентификаторы кавычок для предотвращения SQL-инъекции

Имена таблиц и столбцов не могут передаваться как параметры запроса в SQL. При создании SQL-выражений с динамическими идентификаторами заключайте каждое имя в квадратные скобки и экранируйте все встроенные символы ], чтобы предотвратить 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}]"

Вспомогательные функции в этом разделе используют quote_id() для всех имен таблиц и столбцов в сгенерированных SQL-запросах.

Вставить строки DataFrame

Самый простой подход перебирает строки DataFrame и создаёт по одному INSERT для каждой строки. Простой подход работает для небольших DataFrames, но медленный для больших объёмов, так как каждая строка требует отдельного кругового перехода к серверу.

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")

Для больших DataFrames используйте метод драйвера bulkcopy(), который передаёт строки пакетно по протоколу TDS (Tabular Data Stream) — собственному сетевому протоколу, используемому Microsoft SQL. Этот подход быстрее, чем вставка строка за строкой, поскольку он сокращает количество циклов обмена с базой данных.

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")

Обновление существующих строк из DataFrame

Чтобы обновить строки, которые уже существуют в таблице, переберите строки DataFrame и выполняйте параметризованные операторы UPDATE. key_column указывает, какую строку нужно обновить.

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")

Упсерт (слияние) паттерн

Когда некоторые строки могут быть новыми, а другие уже существуют, используйте SQL-оператор MERGE для вставки или обновления в одной операции. MERGE сравнивает каждую входящую строку с целевой таблицей, используя ключевые столбцы. Если совпадение найдено, выполняется обновление; в противном случае выполняется вставка. MERGE позволяет избежать отдельной проверки наличия.

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

Шаблоны анализа данных

Следующие примеры демонстрируют распространённые задачи анализа, которые сочетают запросы Microsoft SQL с преобразованиями панда.

Агрегированные запросы к 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())

Данные временных рядов

Используйте индексацию и пересэмплирование дат Pandas для работы с временными рядами данных из Microsoft SQL. Чтобы включить операции, такие как скользящие средние и редискретизация, установите столбец даты как индекс 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()

Сводные таблицы из данных SQL

Сводные таблицы преобразуют данные из строк в матричный формат. Чтобы реорганизовать его по измерениям, таким как год, месяц и категория, возьмите исходные данные из Microsoft SQL, затем используйте 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)

Паттерны ETL

Для построения извлечения, трансформации и загрузки конвейеров комбинируйте запросы Microsoft SQL с преобразованиями Pandas для построения конвейеров извлечения, трансформации и загрузки. Драйвер занимается извлечением и загрузкой, а панда — этапом трансформации.

Извлечение, преобразование, загрузка

Этот пример извлекает активные данные клиентов, применяет бизнес-правила для сегментирования клиентов и загружает результаты в таблицу получателей.

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)

Инкрементальный паттерн нагрузки

Для постоянно выполняемых конвейеров данных загружайте только те записи, которые изменились с момента последнего запуска. Этот подход запрашивает таблицу назначения для максимальной временной метки, затем получает только новые записи из источника.

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)

Советы по производительности

Использование соответствующих типов данных

Pandas по умолчанию использует для чисел 64-битные типы, расходуя память впустую, когда достаточно типов меньшего размера. Приведение целочисленных типов и чисел с плавающей точкой к типам меньшей разрядности, а также преобразование строковых столбцов с низкой кардинальностью в категориальный тип может значительно уменьшить потребление памяти.

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

Используйте SQL для тяжёлой работы

Microsoft SQL работает быстрее для агрегирования, фильтрации и объединений, чем перемещение всех необработанных данных по проводу и локальную обработку в Python. Пусть Microsoft SQL выполняет основную работу, когда это возможно, перемещайте только нужные данные по сети и используйте pandas для анализа и трансформаций, которые удобнее в 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 пишет

Для больших таблиц DataFrame, которые слишком велики для одной массовой вставки, разбейте работу на пакеты и отслеживайте прогресс.

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}")