Verwenden Sie mssql-python mit Polars

Polars ist eine leistungsstarke DataFrame-Bibliothek, geschrieben in Rust, die eine schnelle, speichereffiziente Alternative zu Pandas bietet. Polars in Kombination mit dem mssql-python-Treiber ermöglicht es:

  • Laden Sie SQL-Abfrageergebnisse direkt in Polars DataFrames.
  • Verwenden Sie Apache Arrow für den Zero-Copy-Datentransfer von Microsoft SQL.
  • Schreiben Sie Polars DataFrames effizient zurück an Microsoft SQL.
  • Bauen Sie Hochleistungs-Datenpipelines mit lazy evaluation.

Die Beispiele in diesem Artikel fragen die AdventureWorks Beispieldatenbank ab. Falls du sie noch nicht hast, schau dir die AdventureWorks-Beispieldatenbanken an.

Lesen Sie Daten in Polars DataFrames ein

Sie können Microsoft SQL-Daten auf zwei Arten in Polars laden: Zeilen-für-Zeilen-Konvertierung mit Standard-Cursor-Methoden oder Zero-Copy-Transfer über Apache Arrow. "Zero-copy" bedeutet, dass die Daten in einem einzigen Speicherpuffer bleiben, den der Treiber, Arrow und Polars direkt lesen, sodass keine Zeilen in zwischenliegende Python-Objekte dupliziert werden. Verwenden Sie den Arrow-Ansatz für die meisten Workloads wegen dieser Effizienz.

Grundanfrage an DataFrame

Dieser Ansatz ruft alle Zeilen mit dem Standardcursor ab und erstellt manuell einen Polars DataFrame. Es akzeptiert parametrisierte Abfragen zur sicheren Wertsubstitution. Es funktioniert ohne PyArrow, ist aber bei großen Ergebnismengen langsamer, weil jeder Wert durch Python läuft.

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

Wenn Ihre Verbindungszeichenfolge Authentication=ActiveDirectoryDefault verwendet, nutzt der Treiber DefaultAzureCredential, das mehrere Anmeldeinformationsanbieter nacheinander ausprobiert. Die erste Verbindung kann langsam sein, weil das SDK die Kette durchläuft, bis es einen funktionierenden Anbieter findet. In der Produktion gilt: Wenn du weißt, welchen Zugangsdatentyp deine Umgebung verwendet, gib ihn direkt an (zum Beispiel ActiveDirectoryMSI für eine verwaltete Identität), um den Chain Walk zu vermeiden. Weitere Informationen finden Sie unter Microsoft Entra-Authentifizierung.

Die effizienteste Methode, Microsoft SQL-Daten in Polars zu laden, ist über Apache Arrow. Die Methode arrow() des mssql-python-Treibers gibt ein pyarrow.Table zurück, das Polars ohne Zero-Copy-Overhead nutzen kann.

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)

Streame große Datensätze mit Arrow-Batches

Für Datensätze, die nicht in den Arbeitsspeicher passen, verwenden Sie arrow_reader(), um Daten in Streaming-Batches zu verarbeiten. Jeder Batch ist ein pyarrow.RecordBatch, den Polars unabhängig verarbeiten kann, sodass der Speicherverbrauch proportional zu batch_size bleibt und nicht zur vollständigen Ergebnismenge.

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

Verwenden Sie LazyFrames für eine aufgeschobene Ausführung

Mit Polars LazyFrames konntest du eine Kette von Operationen (Filtern, Gruppen, Sortieren) aufbauen, ohne sie sofort auszuführen. Polars optimiert die gesamte Kette vor der Ausführung, was schneller sein kann, als jeder Schritt einzeln anzuwenden.

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)

Schreibe Polars DataFrames in Microsoft SQL

Quote-Identifikatoren zur Verhinderung von SQL-Injektion

Tabellen- und Spaltennamen können in SQL nicht als Abfrageparameter übermittelt werden. Wenn Sie SQL-Anweisungen mit dynamischen Identifikatoren erstellen, packen Sie jeden Namen in eckigen Klammern und vermeiden Sie eingebettete ] Zeichen, um SQL-Injektionen zu verhindern.

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

Die Hilfsfunktionen in diesem Abschnitt verwenden quote_id() für alle Tabellen- und Spaltennamen im generierten SQL.

DataFrame-Zeilen einfügen

Der Zeilen-für-Zeilen-Ansatz iteriert über den DataFrame mit iter_rows(named=True) und führt pro Zeile einen INSERT aus. Dieser Ansatz ist unkompliziert, aber für große Volumina langsam, da jede Reihe eine Hin- und Rückfahrt zum Server erfordert.

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

Für große DataFrames verwenden Sie die Treibermethodebulkcopy(), um Zeilen in großen Mengen über das TDS-Protokoll (Tabular Data Stream) zu senden, das native Wire-Protokoll, das Microsoft SQL verwendet. Dieser Ansatz minimiert Roundtrips und ist schneller als zeilenweise Einfügungen.

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

Datenanalysemuster

Die folgenden Beispiele zeigen gängige Analyseaufgaben, die Microsoft SQL-Abfragen mit Polars-Transformationen kombinieren.

Aggregierte Abfragen

Dieses Beispiel gruppiert Produkte nach Unterkategorie und berechnet Anzahl- und Preisstatistiken in SQL, bevor es die Zusammenfassung in einen Polars DataFrame lädt:

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)

Zeitreihenanalyse

Laden Sie Zeitreihendaten aus Microsoft SQL und fügen Sie berechnete Spalten wie rollende Durchschnitte mit Polars-Ausdrücken hinzu.

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)

Verknüpfe SQL-Daten mit lokalen Dateien

Du kannst Microsoft SQL-Daten bereichern, indem du sie mit lokalen CSV-Dateien in Polars verknüpfst. Lade jede Quelle in ein DataFrame und verbinde sie in den Speicher.

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

Erstellen Sie Extract, Transform and Load (ETL)-Pipelines, indem Sie Microsoft SQL-Abfragen mit Polars-Transformationen kombinieren. Polars-Ausdrücke übernehmen den Transformationsschritt und bulkcopy() übernehmen das Laden.

Extrahieren, Transformieren, Laden

Dieses Beispiel extrahiert aktive Kundendaten über Arrow, wendet Geschäftssegmentierungslogik mit Polars-Ausdrücken an und lädt die Ergebnisse mittels Massenkopie.

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)

Leistungstipps

Die folgenden Tipps helfen dir, das Beste aus der Kombination von mssql-python und Polars herauszuholen.

Lass Microsoft SQL die schwere Arbeit übernehmen

Microsoft SQL ist schneller für Aggregationen, Filtern und Joins, als alle Rohdaten über die Leitung zu ziehen und lokal in Python zu verarbeiten. Lass Microsoft SQL wann immer möglich die schwere Arbeit übernehmen, verschiebe nur die benötigten Daten über das Netzwerk und nutze Polars für Analysen und Transformationen, die in Python bequemer sind.

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

Verwenden Sie Arrow für alle Leseoperationen

Pfeilbasierte Übertragungen vermeiden das Erstellen von Zwischen-Python-Objekten, was den Speicherbedarf reduziert und den Durchsatz verbessert. Bevorzugen Sie bei jeder Ergebnismenge mit mehr als einigen Zeilen cursor.arrow() gegenüber einer manuellen zeilenweisen Konvertierung.

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