Usa mssql-python con Polar

Polars è una libreria DataFrame ad alte prestazioni scritta in Rust che offre un'alternativa veloce ed efficiente in termini di memoria ai panda. Polars, in combinazione con il driver mssql-python, consente di:

  • Carica direttamente i risultati delle query SQL nei DataFrame Polars.
  • Usa Apache Arrow per il trasferimento dati a zero copia da Microsoft SQL.
  • Scrivi i DataFrame di Polars in Microsoft SQL in modo efficiente.
  • Crea pipeline di dati ad alte prestazioni con lazy evaluation.

Gli esempi in questo articolo interrogano il AdventureWorks database di esempio. Se non lo possiedi già, consulta i database di esempio di AdventureWorks.

Leggi i dati nei Polars DataFrame

Puoi caricare dati SQL Microsoft in Polar in due modi: conversione riga per riga tramite metodi standard del cursore, oppure trasferimento zero-copy tramite Apache Arrow. "Zero-copy" significa che i dati rimangono in un unico buffer di memoria che driver, Arrow e Polar leggono direttamente, così nessuna riga viene duplicata in oggetti Python intermedi. Usa l'approccio Arrow per la maggior parte dei carichi di lavoro grazie a questa efficienza.

Query di base su DataFrame

Questo approccio recupera tutte le righe con il cursore standard e costruisce manualmente un DataFrame Polars. Accetta query parametrizzate per una sostituzione sicura dei valori. Funziona senza PyArrow ma è più lento per grandi set di risultati perché ogni valore passa attraverso 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

Se la tua stringa di connessione usa Authentication=ActiveDirectoryDefault, il driver usa DefaultAzureCredential, che prova più fornitori di credenziali in sequenza. La prima connessione può essere lenta perché l'SDK percorre la catena finché non trova un fornitore funzionante. In produzione, se sai quale tipo di credenziale utilizza il tuo ambiente, specificalo direttamente (ad esempio, ActiveDirectoryMSI per l'identità gestita) per evitare il chain walk. Per altre informazioni, vedere Autenticazione di Microsoft Entra.

Il modo più efficiente per caricare dati SQL Microsoft in Polar è tramite Apache Arrow. Il metodo arrow() del driver mssql-python restituisce un pyarrow.Table che Polars può utilizzare senza overhead di copia.

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)

Trasferisci grandi set di dati in batch Arrow

Per i set di dati che non possono essere caricati in memoria, usa arrow_reader() per elaborare i dati in batch di streaming. Ogni lotto è un pyarrow.RecordBatch che Polars può elaborare indipendentemente, quindi l'uso della memoria rimane proporzionale a batch_size anziché all'intero set di risultati.

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

Usa LazyFrames per l'esecuzione differita

Polars LazyFrames ti permette di costruire una catena di operazioni (filtro, raggruppamento, ordinamento) senza eseguirle immediatamente. Polars ottimizza l'intera catena prima di eseguire, il che può essere più veloce rispetto ad applicare ogni passaggio singolarmente.

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)

Scrittura dei DataFrame di Polars in Microsoft SQL Server

Identificatori di quote per prevenire l'iniezione SQL

I nomi di tabelle e colonne non possono essere passati come parametri di query in SQL. Quando costruisci istruzioni SQL con identificatori dinamici, avvolgi ogni nome tra parentesi quadrate e evita eventuali caratteri incorporati ] per evitare l'iniezione 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}]"

Le funzioni helper in questa sezione utilizzano quote_id() per i nomi di tutte le tabelle e colonne nell'SQL generato.

Inserisci righe DataFrame

L'approccio riga per riga itera sul DataFrame con iter_rows(named=True) ed esegue uno INSERT per riga. Questo approccio è semplice ma lento per volumi grandi perché ogni riga richiede un viaggio di andata e ritorno al server.

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

Per i DataFrame di grandi dimensioni, usa il metodo del driver bulkcopy() per inviare righe in blocco tramite il protocollo TDS (Tabular Data Stream), il protocollo nativo a livello di rete usato da Microsoft SQL. Questo approccio minimizza i viaggi di andata e ritorno ed è più veloce rispetto agli inserti fila per fila.

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

Modelli di analisi dei dati

I seguenti esempi mostrano compiti di analisi comuni che combinano query Microsoft SQL con trasformazioni Polars.

Query di aggregazione

Questo esempio raggruppa i prodotti per sottocategoria e calcola le statistiche di conteggio e prezzo in SQL, quindi carica il riassunto in un 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)

Analisi delle serie temporali

Carica i dati delle serie temporali da Microsoft SQL e aggiungi colonne calcolate come medie rotanti usando espressioni polari.

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)

Unisci dati SQL con file locali

Puoi arricchire i dati SQL di Microsoft unendoli ai file CSV locali in Polars. Carica ogni sorgente in un DataFrame e unisci la memoria.

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

Modelli ETL

Costruisci pipeline di estrazione, trasformazione e caricamento (ETL) combinando query Microsoft SQL con trasformazioni Polars. Le espressioni di Polars gestiscono la fase di trasformazione e bulkcopy() gestisce il caricamento.

Estrarre, trasformare, caricare

Questo esempio estrae dati clienti attivi tramite Arrow, applica la logica di segmentazione aziendale con espressioni polari e carica i risultati usando copia di massa.

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)

Suggerimenti per le prestazioni

I seguenti consigli ti aiutano a sfruttare al massimo la combinazione mssql-python e Polars.

Lascia che Microsoft SQL si occupi del lavoro pesante

Microsoft SQL è più veloce per aggregazioni, filtri e join rispetto a recuperare tutti i dati grezzi tramite wire e processarlo localmente in Python. Lascia che Microsoft SQL faccia il lavoro pesante ogni volta che puoi, sposta solo i dati necessari sulla rete e usa i Polari per analisi e trasformazioni più comode in 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
""")

Usa Arrow per tutte le operazioni di lettura

Il trasferimento basato su frecce evita la creazione di oggetti Python intermedi, riducendo l'uso di memoria e migliorando la produttività. Preferisci cursor.arrow() alla conversione manuale riga per riga per qualsiasi set di risultati con più di poche righe.

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