mssql-python 驱动提供 Apache Arrow 提取方法,用于从 Microsoft SQL 和 Azure SQL 数据库中实现高性能列式数据检索。
Apache Arrow 是一个跨语言的内存列式数据开发平台。 驱动程序将 ODBC 结果集直接转换为 C++ 的 Arrow 格式,绕过 Python 对象创建以提升性能。
Arrow 集成可实现:
- 零拷贝数据传输至 Polars、pandas 和 DuckDB。 “零复制”是指数据始终保存在单个内存缓冲区中,由驱动程序写入,使用这些数据的库可直接从中读取,因此无需将各行数据复制到中间的 Python 对象中。
- 通过
RecordBatchReader流式传输结果集,而无需将所有结果加载到内存中。 - 列式数据格式非常适合分析和机器学习工作负载。
- 相比逐行 Python 对象创建,内存占用更低。
游标方法
使用 Arrow 获取方法需要 pyarrow 软件包。 使用 pip install pyarrow 安装它。 如果 pyarrow 未安装,调用任意箭头方法都会触发一个 ImportError。
mssql-python 驱动在光标对象中添加了三种用于 Arrow 数据访问的方法。 这三种方法都会将 ODBC 结果集转换为驱动的 C++ 层中的 Arrow 格式,避免创建中间的 Python 对象。
-
arrow()将整个结果集返回为一个内存表。 用法最为简单。 -
arrow_batch()一次返回一批行,使你能够手动控制循环。 -
arrow_reader()返回一个自动生成批次的迭代器。 最适合流式传输大型结果集。
使用 cursor.arrow(batch_size=8192)
将整个结果集作为单个 pyarrow.Table 获取。 这种方法最简单,当整个结果集能放入内存时效果良好。
import mssql_python
conn = mssql_python.connect(connection_string)
cursor = conn.cursor()
cursor.execute("SELECT ProductID, Name, ListPrice FROM Production.Product")
table = cursor.arrow()
print(type(table)) # <class 'pyarrow.lib.Table'>
print(table.num_rows) # Number of rows fetched
print(table.num_columns) # Number of columns
print(table.schema) # Column names and Arrow types
print(table.to_pandas()) # Convert to pandas DataFrame
注释
如果你的连接字符串使用 Authentication=ActiveDirectoryDefault,驱动程序会使用 DefaultAzureCredential,该机制会按顺序尝试多个凭据提供程序。 第一次连接可能很慢,因为SDK会在链路上走动,直到找到可用的提供者。 在生产环境中,如果你知道环境使用的是哪种凭据类型,可以直接指定它(例如,对于托管标识可指定 ActiveDirectoryMSI),以避免遍历凭据链。 有关详细信息,请参阅 Microsoft Entra 身份验证。
使用 cursor.arrow_batch(batch_size=8192)
获取一个最多包含 batch_size 行的 pyarrow.RecordBatch 在需要精细控制每次提取行数的自定义批处理循环中,请使用此方法。
cursor.execute("SELECT * FROM Production.TransactionHistory")
while True:
batch = cursor.arrow_batch(batch_size=10000)
if batch.num_rows == 0:
break
# Process each batch
print(f"Fetched {batch.num_rows} rows")
使用 cursor.arrow_reader(batch_size=8192)
返回 a pyarrow.RecordBatchReader ,生成 RecordBatch 对象直到结果集耗尽。 对于大型结果集,这种方法是内存效率最高的选择。
cursor.execute("SELECT * FROM Production.TransactionHistory")
reader = cursor.arrow_reader(batch_size=50000)
for batch in reader:
# Process streaming batches without loading all data
print(f"Batch: {batch.num_rows} rows")
常见模式
Arrow 表可直接集成流行的 Python 数据库。 以下示例展示了如何在不复制数据的情况下,将 Arrow 数据传递给 pandas、Polars、DuckDB 和文件格式。
将结果加载到 pandas 中
cursor.execute("SELECT * FROM Production.Product")
table = cursor.arrow()
# Convert to pandas with zero-copy where possible
df = table.to_pandas()
print(df.head())
将结果加载到 Polars 中
import polars as pl
cursor.execute("SELECT * FROM Production.Product")
table = cursor.arrow()
df = pl.from_arrow(table)
print(df)
DuckDB 查询结果
DuckDB 可以直接在 SQL 中查询 Arrow 表,无需复制数据。 当您需要对已经是Arrow格式的结果集进行SQL式分析时,这项功能非常有用。
import duckdb
cursor.execute("SELECT * FROM Sales.SalesOrderHeader")
arrow_table = cursor.arrow()
# Query the Arrow table with DuckDB SQL
result = duckdb.sql("SELECT CustomerID, SUM(TotalDue) FROM arrow_table GROUP BY CustomerID")
print(result.fetchall())
将大型结果集流式写入 Parquet
对于大型结果集,可以直接将 Arrow 批处理流到一个 Parquet 文件,而无需将整个数据集加载到内存中。 它 ParquetWriter 会逐步写入每一批。
import pyarrow.parquet as pq
cursor.execute("SELECT * FROM Production.TransactionHistory")
reader = cursor.arrow_reader(batch_size=100000)
# Write streaming batches to a Parquet file
writer = None
for batch in reader:
if writer is None:
writer = pq.ParquetWriter("output.parquet", batch.schema)
writer.write_batch(batch)
if writer:
writer.close()
导出到其他格式
PyArrow 内置了 CSV 和 Arrow IPC 文件格式(也称为 Feather V2)的写入工具。 Arrow IPC 文件准确保存 Arrow 类型,且读取速度快。
import pyarrow as pa
import pyarrow.csv as pcsv
cursor.execute("SELECT * FROM Production.Product")
table = cursor.arrow()
# Write to CSV
pcsv.write_csv(table, "products.csv")
# Write to an Arrow IPC file
with pa.ipc.new_file("products.arrow", table.schema) as writer:
writer.write_table(table)
数据类型映射
Arrow 抓取方法将 Microsoft SQL 类型映射到 C++ 级别的 Arrow 类型。
| Microsoft SQL 类型 | 箭头类型 |
|---|---|
| int, smallint, tinyint, bigint |
int32、int16、int8、int64 |
| float、real |
float64、float32 |
| 十进制, 数字 | decimal128 |
| 比特 | bool |
| char, varchar, nchar, nvarchar | utf8 |
| 文本, ntext | large_utf8 |
| 二进制, 变数 |
binary、large_binary |
| 日期 | date32 |
| time | time64[us] |
| datetime,datetime2,smalldatetime | timestamp[us] |
| datetimeoffset | timestamp[us, tz=UTC] |
| uniqueidentifier |
utf8 (大写字符串) |
| xml | utf8 |
注释
驱动程序将 datetimeoffset 类型转换为UTC,因为箭头列需要固定的时区。 该驱动程序在转换过程中将 Microsoft SQL 中每个单元格的时区信息规范化为 UTC。
sql_variant 类型不受 Arrow 提取方法支持,并会引发“不支持的数据类型”异常。 对于返回fetchone()列的查询,请使用标准 fetchmany()、 fetchall()或 sql_variant 。
性能注意事项
Arrow 提取方法在分析和大批量数据操作中速度更快,而标准游标方法则更适合结果集较小的事务型模式。
何时使用 Arrow 而非标准 fetch
| Scenario | 建议的方法 |
|---|---|
| 获取几行用于显示 | fetchone() / fetchall() |
| 将数据加载到Panda或Polar | cursor.arrow() |
| 将大型数据集分块处理 | cursor.arrow_reader() |
| 单行查找或小型结果集 | fetchone() / fetchval() |
| 分析管道或聚合管道 |
cursor.arrow() + Polars/DuckDB |
| 将结果写入Parquet或Arrow IPC |
cursor.arrow_reader() + PyArrow I/O |
大型数据集的内存管理
对于可能超出可用内存的结果集,请使用 arrow_reader(),并设置合理的 batch_size。
cursor.execute("SELECT * FROM Production.TransactionHistory")
# Process in batches of 100K rows
reader = cursor.arrow_reader(batch_size=100000)
total_rows = 0
for batch in reader:
# Work with each batch individually
total_rows += batch.num_rows
# batch goes out of scope and memory is freed
print(f"Processed {total_rows} rows")
调整批量大小
该 batch_size 参数控制每批获取的行数。 最佳大小取决于你的行宽度和可用内存。 包含 nvarchar(max) 或 varbinary(max) 等较大列的宽行更适合较小的批处理大小,而窄行则更适合较大的批处理大小。
- 默认(8192):大多数工作负载的平衡良好。
- 较小(1000-5000):用于宽大表格和大列。
- 更大(50000-100000):用于窄表或吞吐量比内存更重要时。
# Narrow table with many rows - use larger batches
cursor.execute("SELECT ProductID, ListPrice FROM Production.Product")
table = cursor.arrow(batch_size=100000)
# Wide table with LOB columns - use smaller batches
cursor.execute("SELECT * FROM Production.Document")
table = cursor.arrow(batch_size=1000)