DuckDBでmssql-pythonを使いましょう

DuckDBは、データをコピーせずにApache Arrowテーブルに直接クエリできる進行中のSQL分析エンジンです。 DuckDBとmssql-pythonドライバーを組み合わせることで、以下が可能になります:

  • PandasやPolarにデータをロードせずに、MicrosoftのSQL結果セットに対して分析的なSQLクエリを実行しましょう。
  • コピーオーバーヘッドゼロのメモリ内のArrowテーブルをクエリします。
  • MicrosoftのSQLデータとローカルファイル(CSV、Parquet、JSON)を単一のDuckDBクエリで結合できます。
  • DuckDBを通じてMicrosoftのSQLデータをParquet、CSV、その他の形式にエクスポートできます。

前提条件

  • Python 3.10 以降。
  • mssql-pythonduckdbpyarrowパッケージです。 すべて pip install mssql-python duckdb pyarrowでインストールしてください。
  • オペレーティング システム固有の 1 回限りの前提条件をインストールします。 Windowsユーザーはこのステップをスキップできます。 プラットフォームの詳細は 「Install mssql-python」をご覧ください。
    apk add libtool krb5-libs krb5-dev
    

SQL データベースを作成する

以下のいずれかのプラットフォームでSQLデータベースを作成または接続してください:

この記事の例は AdventureWorks サンプルデータベースをクエリします。 もしまだ持っていなければ、 AdventureWorksのサンプルデータベースをご覧ください。

依存関係のインストール

pip install mssql-python duckdb pyarrow

DuckDBでMicrosoft SQLデータをクエリ

基本的なワークフローは、mssql-pythonでクエリを実行し、結果をArrowテーブルとして取得し、そのArrowテーブルにDuckDB SQLでクエリを送ります。

基本パターン

まず接続を確立し、Arrowテーブルとしてデータを取得します。

import duckdb
import mssql_python

conn = mssql_python.connect(
    "Server=<server>.database.windows.net;"
    "Database=<database>;"
    "Authentication=ActiveDirectoryDefault;"
    "Encrypt=yes"
)
cursor = conn.cursor()
# Fetch Microsoft SQL data as Arrow
cursor.execute("SELECT * FROM Production.Product WHERE ListPrice > 0")
products = cursor.arrow()

# Query the Arrow table with DuckDB
result = duckdb.sql("""
    SELECT Color, COUNT(*) AS ProductCount, AVG(ListPrice) AS AvgPrice
    FROM products
    GROUP BY Color
    ORDER BY ProductCount DESC
""")
print(result.fetchdf())

DuckDBはproducts ArrowテーブルをPython変数名で参照します。 DuckDBのストレージにはデータがコピーされません。

集約とフィルター

DuckDBのSQLを使ってArrowのデータをグループ化・集約しましょう。

cursor.execute("SELECT * FROM Sales.SalesOrderHeader")
orders = cursor.arrow()

# Top customers by total spend
top_customers = duckdb.sql("""
    SELECT
        CustomerID,
        COUNT(*) AS OrderCount,
        SUM(TotalDue) AS TotalSpent,
        AVG(TotalDue) AS AvgOrderValue
    FROM orders
    GROUP BY CustomerID
    HAVING SUM(TotalDue) > 10000
    ORDER BY TotalSpent DESC
    LIMIT 20
""")
print(top_customers.fetchdf())

複数の Microsoft SQL の結果を結合する

Microsoft SQLから複数のテーブルを取得し、クロスサーバークエリを書かずにDuckDBで結合できます。

# Fetch two tables
cursor.execute("SELECT * FROM Production.Product")
products = cursor.arrow()

cursor.execute("SELECT * FROM Production.ProductSubcategory")
subcategories = cursor.arrow()

# Join in DuckDB
result = duckdb.sql("""
    SELECT
        s.Name AS Subcategory,
        COUNT(*) AS ProductCount,
        ROUND(AVG(p.ListPrice), 2) AS AvgPrice
    FROM products p
    JOIN subcategories s ON p.ProductSubcategoryID = s.ProductSubcategoryID
    GROUP BY s.Name
    ORDER BY AvgPrice DESC
""")
print(result.fetchdf())

Microsoft SQLデータをローカルファイルと結合する

DuckDBはCSV、Parquet、JSONファイルをネイティブに読み取ることができます。 SQL Serverのデータとローカルファイルを1つのクエリで組み合わせます。

CSVファイルで結合してください

CSVファイルを読み込み、Microsoft SQLのデータと結合します。

import csv
from pathlib import Path

cursor.execute("SELECT CustomerID, PersonID FROM Sales.Customer")
customers = cursor.arrow()

csv_path = Path("customer_regions.csv")
with csv_path.open("w", newline="", encoding="utf-8") as file:
    writer = csv.writer(file)
    writer.writerow(["CustomerID", "Region", "Segment"])
    writer.writerows([
        (1, "West", "Premium"),
        (2, "East", "Standard"),
        (3, "Central", "Basic"),
    ])

try:
    result = duckdb.sql("""
        SELECT c.CustomerID, c.PersonID, f.Region, f.Segment
        FROM customers c
        JOIN read_csv_auto('customer_regions.csv') f ON c.CustomerID = f.CustomerID
    """)
    print(result.fetchdf())
finally:
    csv_path.unlink(missing_ok=True)

パーケットファイルで結合します

Parquetファイルを読み込み、Microsoft SQLのデータと結合します。

from pathlib import Path

import pyarrow as pa
import pyarrow.parquet as pq

cursor.execute("SELECT ProductID, Name, ListPrice FROM Production.Product")
products = cursor.arrow()

parquet_path = Path("order_history.parquet")
order_history = pa.table({
    "ProductID": [1, 2, 680],
    "OrderDate": ["2024-06-01", "2024-03-15", "2024-01-10"],
    "Quantity": [10, 5, 3],
})
pq.write_table(order_history, parquet_path)

try:
    result = duckdb.sql("""
        SELECT p.Name, p.ListPrice, h.OrderDate, h.Quantity
        FROM products p
        JOIN read_parquet('order_history.parquet') h ON p.ProductID = h.ProductID
        WHERE h.OrderDate >= '2024-01-01'
    """)
    print(result.fetchdf())
finally:
    parquet_path.unlink(missing_ok=True)

Microsoft SQL データのエクスポート

DuckDBのCOPY文を使って、Microsoft SQLデータをさまざまなファイル形式にエクスポートしてください。

Parquet にエクスポート

データをApache Parquet形式にエクスポートしてください。

cursor.execute("SELECT * FROM Production.Product")
products = cursor.arrow()

duckdb.sql("COPY products TO 'products.parquet' (FORMAT PARQUET)")

CSV にエクスポート

カンマ区切られた値ファイルにデータをエクスポートします:

cursor.execute("SELECT * FROM Sales.SalesOrderHeader")
orders = cursor.arrow()

duckdb.sql("COPY orders TO 'orders.csv' (FORMAT CSV, HEADER)")

パーティション化された Parquet をエクスポート

分散分析用にパーティション分割されたParquetファイルへのデータエクスポート:

import shutil
from pathlib import Path

cursor.execute("SELECT * FROM Sales.SalesOrderHeader")
orders = cursor.arrow()

output_dir = Path("sales_data")
shutil.rmtree(output_dir, ignore_errors=True)

duckdb.sql("""
    COPY (SELECT *, YEAR(OrderDate) AS OrderYear FROM orders)
    TO 'sales_data'
    (FORMAT PARQUET, PARTITION_BY (OrderYear))
""")

大規模な結果セットをストリーミングする

大規模なデータセットの場合、 arrow_reader() を使ってすべての行を一度にメモリに読み込むことなく、ストリーミングバッチでデータを処理します:

cursor.execute("SELECT * FROM Production.TransactionHistory")
reader = cursor.arrow_reader(batch_size=50000)

# Process each batch with DuckDB
total_rows = 0
for batch in reader:
    result = duckdb.sql("""
        SELECT ProductID, SUM(ActualCost) AS TotalCost
        FROM batch
        GROUP BY ProductID
    """)
    total_rows += batch.num_rows
    print(f"Processed {total_rows} rows")

ストリーミング結果を蓄積する

すべてのバッチをまとめるには、各バッチを永続的なDuckDB接続に登録し、結果を段階的に蓄積します。

cursor.execute("SELECT * FROM Production.TransactionHistory")
reader = cursor.arrow_reader(batch_size=50000)

duck = duckdb.connect()
duck.execute("CREATE TABLE transactions (ProductID INT, ActualCost DOUBLE, Quantity INT)")

for batch in reader:
    duck.execute("INSERT INTO transactions SELECT ProductID, ActualCost, Quantity FROM batch")

# Query the accumulated data
result = duck.sql("""
    SELECT ProductID, SUM(ActualCost) AS TotalCost, SUM(Quantity) AS TotalQty
    FROM transactions
    GROUP BY ProductID
    ORDER BY TotalCost DESC
    LIMIT 10
""")
print(result.fetchdf())
duck.close()

パフォーマンスに関するヒント

重労働はMicrosoft SQLに任せましょう

Microsoft SQLは、すべての生データをワイヤーで引き寄せるよりも、フィルタリング、結合、集約がより速いです。 DuckDBは、すでに取得済みの結果セットの二次解析に使うべきであり、SQL Serverクエリ最適化の代替としては使いません。

# Suboptimal: Pull all rows, filter in DuckDB
cursor.execute("SELECT * FROM Sales.SalesOrderHeader")
orders = cursor.arrow()
result = duckdb.sql("SELECT * FROM orders WHERE TotalDue > 1000")

# Better: Filter in Microsoft SQL, analyze in DuckDB
cursor.execute("SELECT * FROM Sales.SalesOrderHeader WHERE TotalDue > 1000")
orders = cursor.arrow()
result = duckdb.sql("SELECT CustomerID, SUM(TotalDue) FROM orders GROUP BY CustomerID")

すべての読み取り操作にはArrowを使用してください

Arrowベースの転送は中間のPythonオブジェクトを作成することを避け、メモリ使用を削減しスループットを向上させます。 DuckDBにデータを渡す際は、手動の行ごとの変換よりも cursor.arrow() を優先します。

大規模なデータセットにはストリーミングを活用してください

利用可能なメモリより大きな結果セットの場合は、arrow_reader()パラメータを持つbatch_sizeを使ってデータを段階的に処理します。