Použijte mssql-python s Polars

Polars je vysoce výkonná knihovna DataFrame napsaná v Rustu, která poskytuje rychlou a paměťově efektivní alternativu k pandas. Polars v kombinaci s ovladačem mssql-python vám umožňuje:

  • Načítání výsledků SQL dotazů přímo do Polars DataFrames.
  • Použijte Apache Arrow pro přenos dat bez kopírování z Microsoft SQL.
  • Efektivně zapisujte Polars DataFrames zpět do Microsoft SQL.
  • Vytvářejte vysoce výkonná datová potrubí s odloženým vyhodnocováním.

Příklady v tomto článku se dotazují na vzorovou databázi AdventureWorks . Pokud ji ještě nemáte, podívejte se na ukázkové databáze AdventureWorks.

Čtěte data do Polars DataFrames

Microsoft SQL data můžete do Polars nahrát dvěma způsoby: převodem řádek po řádku pomocí standardních metod kurzoru nebo přenosem bez kopií přes Apache Arrow. "Zero-copy" znamená, že data zůstávají v jednom paměťovém bufferu, který ovladač, šipka i polární moduly čtou přímo, takže žádné řádky nejsou duplikovány do mezilehlých Python objektů. Použijte přístup Arrow pro většinu pracovních zátěží právě kvůli této efektivitě.

Základní dotaz na DataFrame

Tento přístup načítá všechny řádky pomocí standardního kurzoru a ručně vytváří Polars DataFrame. Přijímá parametrizované dotazy pro bezpečnou substituci hodnot. Funguje to bez PyArrow, ale je to pomalejší u velkých sad výsledků, protože každá hodnota prochází 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

Pokud váš připojovací řetězec používá Authentication=ActiveDirectoryDefault, ovladač používá DefaultAzureCredential, který zkouší více poskytovatelů přihlašovacích údajů za sebou. První spojení může být pomalé, protože SDK prochází řetězec, dokud nenajde funkčního poskytovatele. V produkci, pokud víte, jaký typ přihlašovacích údajů vaše prostředí používá, zadejte ho přímo (například ActiveDirectoryMSI pro spravovanou identitu), abyste se vyhnuli tzv. chain walk. Další informace naleznete v tématu ověřování Microsoft Entra.

Nejefektivnějším způsobem, jak načíst data Microsoft SQL do Polars, je Apache Arrow. Metoda arrow() ovladače mssql-python vrací pyarrow.Table, který může Polars využít bez režie kopírování.

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)

Streamujte velké datové sady pomocí Arrow batches

Pro datové sady, které se nevejdou do paměti, použijte arrow_reader() ke zpracování dat po dávkách v režimu streamování. Každá dávka je pyarrow.RecordBatch, kterou může Polars zpracovat nezávisle, takže využití paměti zůstává úměrné batch_size, nikoli celé sadě výsledků.

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

Použijte LazyFrames pro odložené provádění

Polars LazyFrames vám umožní vytvořit řetězec operací (filtrovat, seskupovat, třídit) bez okamžitého spouštění. Polars optimalizuje celý řetězec před vykonáním, což může být rychlejší než aplikace každého kroku zvlášť.

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)

Zápis datových rámců Polars do Microsoft SQL

Identifikátory uvozovek pro zabránění SQL injekci

Názvy tabulek a sloupců nelze v SQL předávat jako parametry dotazu. Když vytváříte SQL příkazy s dynamickými identifikátory, zabalte každé jméno do hranatých závorek a unikněte všem vloženým ] znacím znakům, abyste zabránili SQL injekci.

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

Pomocné funkce v této sekci používají quote_id() pro všechny názvy tabulek a sloupců ve generovaném SQL.

Vložit řádky do DataFrame

Přístup po řádcích prochází DataFrame pomocí iter_rows(named=True) a pro každý řádek provede jedno INSERT. Tento přístup je jednoduchý, ale při zpracování velkých objemů dat pomalý, protože každý řádek vyžaduje samostatnou komunikaci se serverem.

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

Pro velké DataFrames použijte metodu ovladače bulkcopy() k hromadnému odesílání řádků přes protokol TDS (Tabular Data Stream), nativní wire protokol, který používá Microsoft SQL. Tento přístup minimalizuje zpáteční cesty a je rychlejší než vkládání řádek po řadě.

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

Vzorce analýzy dat

Následující příklady ukazují běžné analytické úkoly, které kombinují Microsoft SQL dotazy s transformacemi Polars.

Agregované dotazy

Tento příklad seskupuje produkty podle podkategorie a počítá statistiky počtu a cen v SQL, poté načte souhrn 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)

Analýza časových řad

Načtěte data časových řad z databáze Microsoft SQL a přidejte vypočítané sloupce, například klouzavé průměry, pomocí výrazů 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)

Spojit SQL data s lokálními soubory

Microsoft SQL data můžete obohatit tím, že je spojíte s lokálními CSV soubory v Polars. Načtěte každý zdroj do DataFrame a spojte je v paměti.

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

ETL vzory

Vytvářejte pipeline pro extrakci, transformaci a načítání (ETL) kombinací dotazů Microsoft SQL s transformacemi Polars. Výrazy v Polars mají na starosti krok transformace a bulkcopy() se stará o načítání.

Extrahování, transformace, načítání

Tento příklad extrahuje aktivní zákaznická data pomocí Arrow, aplikuje logiku obchodní segmentace pomocí výrazů Polars a načte výsledky pomocí hromadného kopírování.

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)

Tipy týkající se výkonu

Následující tipy vám pomohou co nejlépe využít kombinaci mssql-python a Polars.

Nechte Microsoft SQL zvládnout těžkou práci

Microsoft SQL je rychlejší pro agregace, filtrování a spojování než stahování všech surových dat přes kabel a zpracování lokálně v Python. Nechte Microsoft SQL dělat těžkou práci, kdykoli je to možné, přesouvejte jen ta data, která potřebujete, a používejte Polars pro analýzu a transformace, které jsou v Pythonu pohodlnější.

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

Používejte šipku pro všechny čtecí operace

Přenos na základě šipek se vyhýbá vytváření mezipředmětů v Python, což snižuje využití paměti a zlepšuje propustnost. Preferujte cursor.arrow() před manuálním převodem řádek po řádku pro jakoukoli množinu výsledků větší než několik řádků.

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