Użyj mssql-python z Polars

Polars to wysokowydajna biblioteka ramek danych (DataFrame) napisana w języku Rust, która oferuje szybką alternatywę dla pandas, efektywnie wykorzystującą pamięć. Polars w połączeniu ze sterownikiem mssql-python pozwala Ci:

  • Ładuj wyniki zapytań SQL bezpośrednio do Polars DataFrames.
  • Używaj Apache Arrow do transferu danych bez kopiowania z Microsoft SQL.
  • Efektywnie zapisuj Polars DataFrames do Microsoft SQL.
  • Buduj wysokowydajne potoki danych z leniwą oceną.

Przykłady w tym artykule wykonują zapytania do przykładowej bazy danych AdventureWorks. Jeśli jeszcze go nie masz, zobacz przykładowe bazy danych AdventureWorks.

Odczyt danych do Polars DataFrames

Dane Microsoft SQL możesz ładować do Polars na dwa sposoby: konwersją wiersz po wierszu za pomocą standardowych metod kursora lub transferem bez kopii przez Apache Arrow. "Zero-copy" oznacza, że dane pozostają w jednym buforze pamięci, który sterownik, Arrow i Polary odczytują bezpośrednio, więc żadne wiersze nie są duplikowane do pośrednich obiektów Python. Stosuj podejście Arrow do większości obciążeń ze względu na tę efektywność.

Podstawowe zapytanie do DataFrame

To podejście pobiera wszystkie wiersze za pomocą standardowego kursora i ręcznie tworzy ramkę danych Polars. Akceptuje parametryzowane zapytania dla bezpiecznej podstawiania wartości. Działa bez PyArrow, ale jest wolniejszy dla dużych zbiorów wyników, ponieważ każda wartość przechodzi przez Python.

import polars as pl
import mssql_python

conn = mssql_python.connect(connection_string)
cursor = conn.cursor()


def query_to_polars(cursor, query: str, params: dict = None) -> pl.DataFrame:
    """Execute query and return results as Polars DataFrame."""
    cursor.execute(query, params or {})

    columns = [col[0] for col in cursor.description]
    rows = cursor.fetchall()

    data = {col: [row[i] for row in rows] for i, col in enumerate(columns)}
    return pl.DataFrame(data)

# Usage: %(color)s is a parameterized placeholder. The driver safely substitutes
# the value from the dict, which prevents SQL injection.
df = query_to_polars(cursor, "SELECT TOP 5 Name, ListPrice FROM Production.Product WHERE Color = %(color)s", {"color": "Black"})
print(df)

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

Najefektywniejszym sposobem na załadowanie danych SQL Microsoft do Polars jest Apache Arrow. Metoda arrow() sterownika mssql-python zwraca obiekt pyarrow.Table, który Polars może wykorzystać bez narzutu związanego z kopiowaniem danych.

def query_to_polars_arrow(cursor, query: str, params: dict = None) -> pl.DataFrame:
    """Execute query and load results through Arrow for best performance."""
    cursor.execute(query, params or {})
    arrow_table = cursor.arrow()
    return pl.from_arrow(arrow_table)

# Usage
df = query_to_polars_arrow(cursor, "SELECT ProductID, Name, ListPrice FROM Production.Product")
print(df)

Przesyłaj strumieniowo duże zbiory danych w partiach Arrow

W przypadku zbiorów danych, które nie mieszczą się w pamięci, używaj arrow_reader() do przetwarzania danych w partiach strumieniowych. Każda partia to rodzaj pyarrow.RecordBatch, który Polars może przetwarzać niezależnie, więc zużycie pamięci pozostaje proporcjonalne do batch_size, a nie do całego zestawu wyników.

def process_large_query(cursor, query: str, params: dict = None, batch_size: int = 50000) -> pl.DataFrame:
    """Process large query results as streaming Arrow batches."""
    cursor.execute(query, params or {})
    reader = cursor.arrow_reader(batch_size=batch_size)

    results = []
    for batch in reader:
        chunk_df = pl.from_arrow(batch)
        # Process each chunk
        results.append(chunk_df)

    return pl.concat(results) if results else pl.DataFrame()

# Usage
df = process_large_query(cursor, "SELECT * FROM Production.TransactionHistory")

Używaj LazyFrames do odroczonego wykonania

Polars LazyFrames pozwalają tworzyć sekwencję operacji (filtrowanie, grupowanie, sortowanie) bez ich natychmiastowego wykonywania. Polars optymalizuje cały łańcuch przed wykonaniem, co może być szybsze niż stosowanie każdego kroku osobno.

def query_to_lazy(cursor, query: str, params: dict = None) -> pl.LazyFrame:
    """Execute query and return a Polars LazyFrame."""
    cursor.execute(query, params or {})
    arrow_table = cursor.arrow()
    return pl.from_arrow(arrow_table).lazy()

# Build a query plan without executing immediately
lf = query_to_lazy(cursor, "SELECT SalesOrderID, CustomerID, TotalDue, OrderDate FROM Sales.SalesOrderHeader")
result = (
    lf.filter(pl.col("TotalDue") > 100)
    .group_by("CustomerID")
    .agg([
        pl.col("TotalDue").sum().alias("TotalSpent"),
        pl.col("SalesOrderID").count().alias("OrderCount")
    ])
    .sort("TotalSpent", descending=True)
    .collect()  # Execute the optimized plan
)
print(result)

Zapisuj Polars 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

Podejście polegające na przetwarzaniu wiersz po wierszu iteruje po ramce danych DataFrame za pomocą iter_rows(named=True) i wykonuje jedno INSERT dla każdego wiersza. To podejście jest proste, ale wolne przy dużych wolumenach danych, ponieważ każdy wiersz wymaga komunikacji w obie strony z serwerem.

def polars_to_sql(cursor, conn, df: pl.DataFrame, table: str) -> int:
    """Write Polars DataFrame to Microsoft SQL table."""
    columns = df.columns
    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.iter_rows(named=True):
        params = {k: (None if v is None else v) for k, v in row.items()}
        cursor.execute(query, params)
        rows_inserted += 1

    conn.commit()
    return rows_inserted

# Usage
cursor.execute("CREATE TABLE #PolarsInsert (Name NVARCHAR(100), Price DECIMAL(10,2), CategoryID INT)")
df = pl.DataFrame({
    "Name": ["Product A", "Product B"],
    "Price": [29.99, 49.99],
    "CategoryID": [1, 2]
})
rows = polars_to_sql(cursor, conn, df, "#PolarsInsert")
print(f"Inserted {rows} rows")

Dla dużych DataFrame'ów użyj metody sterownika bulkcopy() do przesyłania wierszy masowo przez protokół TDS (Tabular Data Stream), natywny protokół przewodowy używany przez Microsoft SQL. Ta metoda ogranicza liczbę połączeń zwrotnych i jest szybsza niż wstawianie wiersz po wierszu.

def polars_to_sql_bulk(conn, df: pl.DataFrame, table: str) -> int:
    """Bulk insert Polars DataFrame using BCP for best performance."""
    rows = [tuple(None if v is None else v for v in row) for row in df.iter_rows()]

    cursor = conn.cursor()
    result = cursor.bulkcopy(table, rows)
    conn.commit()
    return result["rows_copied"]

# Usage
cursor.execute("CREATE TABLE ##PolarsBulk (Name NVARCHAR(50), Price FLOAT, CategoryID INT)")
conn.commit()
df = pl.DataFrame({
    "Name": ["Product A", "Product B", "Product C"],
    "Price": [29.99, 49.99, 19.99],
    "CategoryID": [1, 2, 1]
})
rows = polars_to_sql_bulk(conn, df, "##PolarsBulk")
print(f"Bulk inserted {rows} rows")

Wzorce analizy danych

Poniższe przykłady pokazują typowe zadania analityczne łączące zapytania Microsoft SQL z transformacjami Polars.

Zapytania agregowane

Ten przykład grupuje produkty według podkategorii i oblicza statystyki liczby oraz cen w SQL, a następnie ładuje podsumowanie do Polars DataFrame:

def get_sales_summary(cursor) -> pl.DataFrame:
    """Get sales summary by subcategory."""
    cursor.execute("""
        SELECT
            sc.Name AS SubcategoryName,
            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 sc ON p.ProductSubcategoryID = sc.ProductSubcategoryID
        GROUP BY sc.Name
        ORDER BY ProductCount DESC
    """)
    return pl.from_arrow(cursor.arrow())

df = get_sales_summary(cursor)
print(df)

Analiza szeregów czasowych

Załaduj dane szeregów czasowych z Microsoft SQL i dodaj kolumny obliczeniowe, takie jak średnie kroczące, za pomocą wyrażeń biblioteki Polars.

def get_daily_sales(cursor, start_date: str, end_date: str) -> pl.DataFrame:
    """Get daily sales and compute rolling statistics."""
    cursor.execute("""
        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})

    df = pl.from_arrow(cursor.arrow())

    # Add rolling 7-day average
    df = df.with_columns(
        pl.col("Revenue").rolling_mean(window_size=7).alias("RollingAvg")
    )
    return df

sales_df = get_daily_sales(cursor, "2013-01-01", "2013-12-31")
print(sales_df)

Połącz dane SQL z plikami lokalnymi

Możesz wzbogacić dane Microsoft SQL, łącząc je z lokalnymi plikami CSV w Polars. Załaduj każde źródło do DataFrame'u i łącz w pamięci.

# Load SQL data via Arrow
cursor.execute("SELECT c.CustomerID, p.FirstName, p.LastName FROM Sales.Customer c JOIN Person.Person p ON c.PersonID = p.BusinessEntityID")
customers = pl.from_arrow(cursor.arrow())

# Load local CSV
orders = pl.read_csv("orders_export.csv")

# Join in Polars
result = customers.join(orders, on="CustomerID", how="inner")
print(result)

Wzorce ETL

Twórz potoki wyodrębniania, przekształcania i ładowania (ETL), łącząc zapytania Microsoft SQL z transformacjami Polars. Wyrażenia biblioteki Polars odpowiadają za etap transformacji, a bulkcopy() za ładowanie danych.

Wyodrębnianie, przekształcanie, ładowanie

Ten przykład wyodrębnia dane aktywnych klientów za pomocą Arrow, stosuje logikę segmentacji biznesowej za pomocą wyrażeń Polars i ładuje wyniki przy użyciu operacji kopiowania zbiorczego.

def etl_pipeline(source_cursor, dest_conn):
    """ETL pipeline using Polars transformations."""

    # Extract: derive a per-customer summary from order history via Arrow
    source_cursor.execute("""
        SELECT
            CustomerID,
            COUNT(*) AS OrderCount,
            SUM(TotalDue) AS TotalSpent
        FROM Sales.SalesOrderHeader
        WHERE OrderDate > DATEADD(YEAR, -1, (SELECT MAX(OrderDate) FROM Sales.SalesOrderHeader))
        GROUP BY CustomerID
    """)
    df = pl.from_arrow(source_cursor.arrow())
    df = df.with_columns(pl.col("TotalSpent").cast(pl.Float64))

    # Transform with Polars expressions
    df = df.with_columns([
        pl.when(pl.col("TotalSpent") > 1000).then(pl.lit("Platinum"))
          .when(pl.col("TotalSpent") > 500).then(pl.lit("Gold"))
          .when(pl.col("TotalSpent") > 100).then(pl.lit("Silver"))
          .otherwise(pl.lit("Bronze"))
          .alias("CustomerSegment"),
        (pl.col("TotalSpent") / pl.col("OrderCount").clip(lower_bound=1))
          .alias("AvgOrderValue"),
        (pl.col("TotalSpent") > 500).alias("IsHighValue")
    ])

    # Load via bulk copy into the destination table
    dest_cursor = dest_conn.cursor()
    dest_cursor.execute("""
        CREATE TABLE ##CustomerAnalytics (
            CustomerID INT,
            CustomerSegment NVARCHAR(20),
            AvgOrderValue FLOAT,
            IsHighValue BIT
        )
    """)
    dest_conn.commit()

    load_df = df.select(["CustomerID", "CustomerSegment", "AvgOrderValue", "IsHighValue"])
    polars_to_sql_bulk(dest_conn, load_df, "##CustomerAnalytics")

    return len(df)

Porady dotyczące wydajności

Poniższe wskazówki pomogą Ci w pełni wykorzystać połączenie mssql-python i Polars.

Pozwól Microsoft SQL wykonać ciężką pracę

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 Polars do analiz i transformacji, które są wygodniejsze w Python.

# Avoid: pulling all rows over the wire to aggregate locally in Polars
df = query_to_polars_arrow(cursor, "SELECT * FROM Production.Product WHERE Color IS NOT NULL")  # transfers entire table
summary = df.group_by("Color").agg(pl.col("ListPrice").sum())  # aggregation that SQL can do faster

# Better: push the aggregation into SQL and transfer only the summary
df = query_to_polars_arrow(cursor, """
    SELECT Color AS Category, SUM(ListPrice) AS TotalAmount
    FROM Production.Product
    WHERE Color IS NOT NULL
    GROUP BY Color
""")

Używaj Arrow do wszystkich operacji odczytu

Transfer oparty na strzałkach unika tworzenia pośrednich obiektów Python, co zmniejsza zużycie pamięci i poprawia przepustowość. Preferować cursor.arrow() zamiast ręcznej konwersji wiersz po wierszu dla każdego zbioru wyników większych niż kilka wierszy.

# Suboptimal: Row-by-row conversion
cursor.execute("SELECT * FROM Production.TransactionHistory")
columns = [col[0] for col in cursor.description]
rows = cursor.fetchall()
df = pl.DataFrame({col: [row[i] for row in rows] for i, col in enumerate(columns)})

# Better: Arrow-based transfer
cursor.execute("SELECT * FROM Production.TransactionHistory")
df = pl.from_arrow(cursor.arrow())