Gunakan mssql-python dengan Polars

Polars adalah pustaka DataFrame berperforma tinggi yang ditulis dalam Rust yang menyediakan alternatif panda yang cepat dan hemat memori. Polar yang dikombinasikan dengan driver mssql-python memungkinkan Anda:

  • Muat hasil kueri SQL langsung ke Polars DataFrames.
  • Gunakan Apache Arrow untuk transfer data tanpa salinan dari Microsoft SQL.
  • Tulis Polars DataFrames kembali ke Microsoft SQL secara efisien.
  • Bangun alur data berperforma tinggi dengan evaluasi lambat.

Contoh dalam artikel ini menggunakan kueri pada database sampel AdventureWorks ini. Jika Anda belum memilikinya, lihat contoh database AdventureWorks.

Membaca data ke dalam Polars DataFrames

Anda dapat memuat data Microsoft SQL ke Polar dengan dua cara: konversi baris demi baris melalui metode kursor standar, atau transfer tanpa salinan melalui Apache Arrow. "Zero-copy" berarti data tetap berada dalam satu buffer memori yang dibaca langsung oleh driver, Arrow, dan Polars, sehingga tidak ada baris yang diduplikasi menjadi objek Python perantara. Gunakan pendekatan Arrow untuk sebagian besar beban kerja karena efisiensinya.

Kueri dasar ke DataFrame

Pendekatan ini mengambil semua baris dengan kursor standar dan membuat Polars DataFrame secara manual. Ini menerima kueri berparameter untuk substitusi nilai yang aman. Ini bekerja tanpa PyArrow tetapi lebih lambat untuk kumpulan hasil yang besar karena setiap nilai melewati 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

Jika string koneksi Anda menggunakan Authentication=ActiveDirectoryDefault, driver menggunakan DefaultAzureCredential, yang mencoba beberapa penyedia kredensial secara berurutan. Koneksi pertama bisa lambat karena SDK menelusuri rantai hingga menemukan penyedia yang berfungsi. Dalam lingkungan produksi, jika Anda mengetahui jenis kredensial yang digunakan oleh lingkungan Anda, tentukan secara langsung (misalnya, ActiveDirectoryMSI untuk identitas terkelola) agar terhindar dari penelusuran berantai. Untuk informasi selengkapnya, lihat Autentikasi Microsoft Entra.

Cara paling efisien untuk memuat data Microsoft SQL ke dalam Polars adalah melalui Apache Arrow. Metode arrow() milik driver mssql-python mengembalikan pyarrow.Table yang dapat digunakan oleh Polars dengan overhead tanpa penyalinan.

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)

Streaming dataset besar dengan batch Arrow

Untuk himpunan data yang tidak muat dalam memori, gunakan arrow_reader() untuk memproses data dalam batch streaming. Setiap batch merupakan pyarrow.RecordBatch yang dapat diproses Polars secara independen, sehingga penggunaan memori tetap proporsional terhadap batch_size, bukan terhadap seluruh kumpulan hasil.

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

Menggunakan LazyFrames untuk eksekusi yang ditangguhkan

Polars LazyFrames memungkinkan Anda membangun rantai operasi (filter, grup, sortir) tanpa segera menjalankannya. Polars mengoptimalkan rantai penuh sebelum dieksekusi, yang bisa lebih cepat daripada menerapkan setiap langkah satu per satu.

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)

Menulis Polars DataFrames ke Microsoft SQL

Pengidentifikasi kutipan untuk mencegah injeksi SQL

Nama tabel dan kolom tidak dapat diteruskan sebagai parameter kueri di SQL. Saat Anda membuat pernyataan SQL dengan pengidentifikasi dinamis, bungkus setiap nama dalam tanda kurung siku dan keluarkan karakter yang disematkan ] untuk mencegah injeksi 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}]"

Fungsi pembantu di bagian ini digunakan quote_id() untuk semua nama tabel dan kolom di SQL yang dihasilkan.

Sisipkan baris DataFrame

Pendekatan baris demi baris mengulangi DataFrame dengan iter_rows(named=True) dan mengeksekusi satu INSERT per baris. Pendekatan ini mudah tetapi lambat untuk volume besar karena setiap baris memerlukan perjalanan pulang pergi ke 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")

Untuk DataFrame besar, gunakan metode driver bulkcopy() untuk mengirim baris secara massal melalui protokol TDS (Tabular Data Stream), protokol kawat asli yang digunakan Microsoft SQL. Pendekatan ini meminimalkan perjalanan pulang pergi dan lebih cepat daripada sisipan baris demi baris.

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

Pola analisis data

Contoh berikut menunjukkan tugas analisis umum yang menggabungkan kueri Microsoft SQL dengan transformasi Polars.

Kueri agregat

Contoh ini mengelompokkan produk berdasarkan subkategori dan menghitung jumlah dan statistik harga dalam SQL, lalu memuat ringkasan ke dalam 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)

Analisis data deret waktu

Muat data deret waktu dari Microsoft SQL dan tambahkan kolom komputasi seperti rata-rata bergulir menggunakan ekspresi 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)

Gabungkan data SQL dengan file lokal

Anda dapat memperkaya data Microsoft SQL dengan menggabungkannya dengan file CSV lokal di Polars. Muat setiap sumber ke dalam DataFrame dan gabungkan dalam memori.

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

Pola ETL

Bangun alur ekstrak, transformasi, dan muat (ETL) dengan menggabungkan kueri Microsoft SQL dengan transformasi Polars. Ekspresi Polars menangani langkah transformasi, dan bulkcopy() menangani pemuatan.

Mengekstrak, mengubah, memuat

Contoh ini mengekstrak data pelanggan aktif melalui Arrow, menerapkan logika segmentasi bisnis dengan ekspresi Polars, dan memuat hasilnya menggunakan salinan massal.

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)

Tips kinerja

Kiat-kiat berikut membantu Anda mendapatkan hasil maksimal dari kombinasi mssql-python dan Polars.

Biarkan Microsoft SQL menangani pekerjaan berat

Microsoft SQL lebih cepat untuk agregasi, pemfilteran, dan gabungan daripada menarik semua data mentah Anda melalui kabel dan pemrosesan secara lokal di Python. Biarkan Microsoft SQL melakukan pekerjaan berat bila memungkinkan, pindahkan hanya data yang Anda butuhkan di seluruh jaringan dan gunakan Polars untuk analisis dan transformasi yang lebih nyaman di 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
""")

Gunakan Arrow untuk semua operasi pembacaan

Transfer berbasis panah menghindari pembuatan objek Python perantara, yang mengurangi penggunaan memori dan meningkatkan throughput. Gunakan cursor.arrow() daripada konversi manual per baris untuk set hasil apa pun yang berisi lebih dari beberapa baris.

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