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-python、duckdb和pyarrow软件包。 使用pip install mssql-python duckdb pyarrow安装全部。 - 安装一次性操作系统特定的先决条件。 Windows 用户可以跳过这一步。 完整平台详情请参见 “安装 mssql-python”。
创建 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 进行数据增量处理。