用mssql-python配合DuckDB

DuckDB 是一个正在进行的 SQL 分析引擎,可以直接查询 Apache Arrow 表而无需复制数据。 将 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 安装全部。
  • 安装一次性操作系统特定的先决条件。 Windows 用户可以跳过这一步。 完整平台详情请参见 “安装 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 表,然后用 DuckDB SQL 查询该 Arrow 表。

基本模式

首先建立连接,并以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 通过 Arrow 表的 Python 变量名 products 引用该表。 没有数据被复制到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 数据与本地文件合并为单一查询。

用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 文件联接

加载 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

基于箭头的传输避免了创建中间的 Python 对象,从而减少内存使用并提升吞吐量。 在将数据传递给DuckDB时,优先选择 cursor.arrow() 逐行转换而非手动转换。

对大型数据集使用流式处理

对于超出可用内存的结果集,使用 arrow_reader() 参数 batch_size 进行数据增量处理。