Utilisez mssql-python avec Polars

Polars est une bibliothèque DataFrame haute performance écrite en Rust qui offre une alternative rapide et économe en mémoire aux pandas. Les polars combinés au pilote mssql-python vous permettent :

  • Chargez directement les résultats des requêtes SQL dans les DataFrames Polars.
  • Utilisez Apache Arrow pour le transfert de données sans copie depuis Microsoft SQL.
  • Réécrivez efficacement les DataFrames de Polars dans Microsoft SQL Server.
  • Construisez des pipelines de données haute performance avec une évaluation paresseuse.

Les exemples de cet article interrogent la base de données d’exemple AdventureWorks . Si vous ne l’avez pas déjà, consultez les bases de données d’exemple AdventureWorks.

Lire les données dans les DataFrames Polars

Vous pouvez charger des données SQL Microsoft dans les Polar de deux manières : conversion ligne par ligne via des méthodes standard de curseur, ou transfert sans copie via Apache Arrow. « Zero-copy » signifie que les données restent dans un seul tampon mémoire que le pilote, Arrow et Polars lisent directement, de sorte qu’aucune ligne ne soit dupliquée en objets Python intermédiaires. Utilisez l’approche Arrow pour la plupart des charges de travail grâce à cette efficacité.

Requête de base vers DataFrame

Cette approche récupère toutes les lignes avec le curseur standard et construit manuellement un DataFrame Polars. Il accepte les requêtes paramétrées pour une substitution de valeur sûre. Il fonctionne sans PyArrow mais est plus lent pour les grands ensembles de résultats car chaque valeur passe par 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

Si votre chaîne de connexion utilise Authentication=ActiveDirectoryDefault, le pilote utilise DefaultAzureCredential, ce qui essaie plusieurs fournisseurs d’identifiants en séquence. La première connexion peut être lente car le SDK parcourt la chaîne jusqu’à ce qu’il trouve un fournisseur fonctionnel. En production, si vous savez quel type d’identifiant votre environnement utilise, spécifiez-le directement (par exemple, ActiveDirectoryMSI pour l’identité gérée) afin d’éviter la marche en chaîne. Pour plus d’informations, consultez Authentification Microsoft Entra.

La façon la plus efficace de charger des données SQL Microsoft dans Polars est via Apache Arrow. La méthode arrow() du pilote mssql-python renvoie un pyarrow.Table que Polars peut exploiter sans surcoût de copie.

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)

Diffusez de grands ensembles de données avec des lots Arrow

Pour les jeux de données qui ne rentrent pas en mémoire, utilisez arrow_reader() pour traiter les données en lots en streaming. Chaque lot est un pyarrow.RecordBatch que Polars peut traiter indépendamment, de sorte que l’utilisation de la mémoire reste proportionnelle à batch_size plutôt qu’au jeu de résultats complet.

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

Utilisez LazyFrames pour l’exécution différée

Les polars LazyFrames permettent de construire une chaîne d’opérations (filtre, groupe, tri) sans les exécuter immédiatement. Polars optimise toute la chaîne avant d’exécuter, ce qui peut être plus rapide que d’appliquer chaque étape individuellement.

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)

Écrire des DataFrames Polars vers Microsoft SQL

Mettre les identifiants entre guillemets pour éviter l’injection SQL

Les noms de tables et de colonnes ne peuvent pas être passés comme paramètres de requête en SQL. Lorsque vous construisez des instructions SQL avec des identifiants dynamiques, enroulez chaque nom entre crochets et éliminez les caractères intégrés ] pour éviter l’injection 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}]"

Les fonctions d’assistance de cette section utilisent quote_id() pour tous les noms de table et de colonne dans le SQL généré.

Insérer des lignes de DataFrame

L’approche ligne par ligne parcourt le DataFrame avec iter_rows(named=True) et exécute un INSERT par ligne. Cette approche est simple mais lente pour les gros volumes car chaque ligne nécessite un aller-retour jusqu’au serveur.

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

Pour les DataFrames volumineux, utilisez la méthode bulkcopy() du pilote pour envoyer des lignes en masse via le protocole TDS (Tabular Data Stream), le protocole natif de communication utilisé par Microsoft SQL. Cette approche minimise les allers-retours et est plus rapide que les insertions ligne par ligne.

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

Schémas d’analyse des données

Les exemples suivants montrent des tâches d’analyse courantes qui combinent des requêtes Microsoft SQL avec des transformations Polaires.

Requêtes agrégées

Cet exemple regroupe les produits par sous-catégorie et calcule les statistiques de comptage et de prix en SQL, puis charge le résumé dans 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)

Analyse de série chronologique

Chargez les données de séries temporelles depuis Microsoft SQL et ajoutez des colonnes calculées comme des moyennes roulantes en utilisant des expressions polaires.

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)

Relier les données SQL avec des fichiers locaux

Vous pouvez enrichir les données SQL de Microsoft en les associant à des fichiers CSV locaux dans Polars. Chargez chaque source dans un DataFrame et rejoignez la mémoire.

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

Motifs ETL

Construisez des pipelines d’extraction, transformation et chargement (ETL) en combinant les requêtes Microsoft SQL avec les transformations Polars. Les expressions Polars gèrent l’étape de transformation, et bulkcopy() le chargement.

Extraire, transformer, charger

Cet exemple extrait les données des clients actifs à l’aide d’Arrow, applique une logique de segmentation métier avec des expressions Polars, et charge les résultats au moyen d’une opération de copie en masse.

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)

Astuces pour les performances

Les conseils suivants vous aident à tirer le meilleur parti de la combinaison mssql-python et Polars.

Laissez Microsoft SQL gérer le gros du travail

Microsoft SQL est plus rapide pour les agrégations, le filtrage et les jointures que de récupérer toutes vos données brutes via le réseau et de les traiter localement en Python. Laissez Microsoft SQL faire le gros du travail dès que possible, ne transférez que les données dont vous avez besoin sur le réseau et utilisez Polars pour l’analyse et les transformations plus pratiques en 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
""")

Utilisez Arrow pour toutes les opérations de lecture

Le transfert basé sur des flèches évite de créer des objets Python intermédiaires, ce qui réduit la consommation de mémoire et améliore le débit. Privilégiez cursor.arrow() à la conversion manuelle ligne par ligne pour tout ensemble de résultats de plus de quelques lignes.

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