Apache Arrowでmssql-pythonを使う

mssql-pythonドライバーは、Microsoft SQLおよびAzure SQL Databaseからの高性能な列形式データ取得のためのApache Arrowフェッチ手法を提供します。

Apache Arrowは、メモリ内のカラムデータのためのクロス言語開発プラットフォームです。 このドライバーは、ODBCの結果セットを直接C++のArrow形式に変換し、Pythonオブジェクト作成を省略してパフォーマンスを向上させます。

Arrow 統合により可能なこと:

  • Polars、pandas、DuckDBへのゼロコピーデータ転送。 「ゼロコピー」とは、データがドライバが書き込み、使用するライブラリが直接読み込む単一のメモリバッファに留まるため、行が中間のPythonオブジェクトに重複しないようにすることを意味します。
  • RecordBatchReader を介して、結果セット全体をメモリに読み込むことなくストリーミングします。
  • 分析や機械学習のワークロードに最適なカラム形式データ形式。
  • 行ごとのPythonオブジェクト作成と比べてメモリ使用量が削減されました。

カーソル メソッド

pyarrowパッケージはArrowフェッチメソッドを使用する必要があります。 pip install pyarrowと共にインストールします。 pyarrowインストールされていなければ、Arrowメソッドを呼び出すとImportErrorが上がります。

mssql-pythonドライバーは、Arrowのデータアクセスのためにカーソルオブジェクトに3つのメソッドを追加します。 これら3つの方法すべてが、ODBCの結果セットをドライバーのC++レイヤーでArrow形式に変換するため、中間のPythonオブジェクトの作成を回避できます。

  • arrow() 結果セット全体を1つのメモリ内テーブルとして返します。 使い方が最も単純です。
  • 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

Note

もし接続文字列がAuthentication=ActiveDirectoryDefaultを使っているなら、ドライバーはDefaultAzureCredentialを使い、複数の認証情報提供者を連続して試します。 最初の接続は遅くなることがあります。なぜならSDKが動作するプロバイダーを見つけるまでチェーンを歩くからです。 本番環境では、環境がどの認証情報タイプを使っているか分かっているなら、チェーンウォークを避けるために直接指定してください(例えばマネージドIDの 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) の使用

結果セットが尽きるまで 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")

リーダーは接続を通じて結果をストリーミングするため、未読リーダーが開いている間は別の文を開始できません。 試みても Connection is busy with results for another command エラーで失敗します。

リーダーを解放するのは3つのことです:最後まで繰り返しること、親カーソルを閉じること、リーダーを閉じることです。 結果セットが使い果たされる前に読み終えてカーソルを使い続けた場合は、リーダーを閉じてください。 閉じると親カーソルもリセットされるので、別の文を実行できます。

リーダーをコンテキストマネージャーとして使い、例外がループを中断しても閉じられるようにします:

cursor.execute("SELECT * FROM Production.TransactionHistory")

rows_seen = 0
with cursor.arrow_reader(batch_size=50000) as reader:
    for batch in reader:
        rows_seen += batch.num_rows
        if rows_seen >= 100000:
            break

# The reader is closed here, and the cursor is ready for the next statement.
cursor.execute("SELECT COUNT(*) FROM Production.TransactionHistory")

直接 reader.close() に電話することもできます。 複数回呼び出しても安全です。また、reader.closed プロパティは、それを閉じたかどうかを示します。

一般的なパターン

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 データを SQL Server にロードする

cursor.bulkcopy_arrow()メソッドは、Arrow データをまず Python の行タプルに変換することなく、テーブルに書き込みます。 sourceの議論は以下のいずれかを受け入れます。

  • pyarrow.Table です。

  • pyarrow.RecordBatch です。

  • pyarrow.RecordBatchReader によって返されるリーダーを含む cursor.arrow_reader()

  • __arrow_c_stream____arrow_c_array__を通じてArrow Cのデータインターフェースを公開するオブジェクトです。

import mssql_python
import pyarrow as pa

conn = mssql_python.connect(connection_string)

# bulkcopy_arrow() opens its own connection, so commit the table creation first.
conn.autocommit = True
cursor = conn.cursor()

cursor.execute("""
    CREATE TABLE ##SensorArchive (
        SensorID int NOT NULL,
        Reading float NULL,
        Location nvarchar(50) NULL
    )
""")

table = pa.table({
    "SensorID": pa.array([1, 2, 3], type=pa.int32()),
    "Reading": pa.array([20.5, None, 22.1], type=pa.float64()),
    "Location": pa.array(["Plant A", "Plant B", None], type=pa.string()),
})

result = cursor.bulkcopy_arrow("##SensorArchive", table)
print(f"Copied {result['rows_copied']} rows in {result['batch_count']} batches")

矢印のnull値はSQLのNULL値として表記されます。

結果セットを別のテーブルにストリーミングする

bulkcopy_arrow()リーダーを受け入れているため、大きな結果セットをメモリに具現化せずにテーブル間で移動させることができます。

cursor.execute("""
    CREATE TABLE ##ProductArchive (
        ProductID int NOT NULL,
        Name nvarchar(50) NOT NULL,
        ListPrice money NOT NULL
    )
""")

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

with cursor.arrow_reader(batch_size=100000) as reader:
    result = cursor.bulkcopy_arrow("##ProductArchive", reader, batch_size=100000)

print(f"Copied {result['rows_copied']} rows")

矢印の種類を宛先の列に合わせる

Arrowライターは、各Arrow列タイプが宛先のSQLカラムタイプと互換性があることを求めています。 ファミリー間で変換されないため、不一致があっても行が書かれる前に ValueError が上がります:

ValueError: Cannot map Arrow column 'ListPrice' (Float64) to SQL column 'ListPrice'
(Money): Usage Error: type combination is not supported by the Arrow row-major writer

データ型のマッピングを逆に使い、矢印タイプを選択します。 貨幣小数数字 の列は decimal128が必要で、 float64ではありません。 cursor.arrow()で読み戻すデータはすでに正しい型を持っているため、SQL Serverからのテーブル読み取りは変換なしで対応テーブルに読み込まれます。

列を名前で対応付ける

Arrowの列順序が宛先テーブルと一致しない場合は、宛先列名をArrowの列順序で渡す column_mappings :

from decimal import Decimal

table = pa.table({
    "Name": pa.array(["Widget"], type=pa.string()),
    "ProductID": pa.array([9001], type=pa.int32()),
    "ListPrice": pa.array([Decimal("12.34")], type=pa.decimal128(19, 4)),
})

cursor.bulkcopy_arrow(
    "##ProductArchive",
    table,
    column_mappings=["Name", "ProductID", "ListPrice"],
)

この方法は cursor.bulkcopy()と同じ選択肢、すなわち batch_sizetimeoutkeep_identitytable_lockkeep_nullsを受け入れます。 これらのオプションの詳細については、一括コピーをご覧ください。

Note

Arrow ソースを cursor.bulkcopy() に渡すと、TypeError が発生し、cursor.bulkcopy_arrow() に誘導されます。

データ型マッピング

Arrowのフェッチ手法は、MicrosoftのSQL型をC++レベルでArrow型にマッピングします。

Microsoft SQL タイプ 矢印の種類
int, smallint, tinyint, bigint int32int16int8int64
floatreal float64float32
十進法数値 decimal128
bit bool
char, varchar, nchar, nvarchar utf8
テキストnテキスト large_utf8
バイナリ可変長バイナリ binarylarge_binary
date date32
time time64[us]
デートタイムデートタイム2スモールデイトタイム timestamp[us]
datetimeoffset timestamp[us, tz=UTC]
uniqueidentifier utf8 (大文字の文字列)
xml utf8

Note

ドライバーは datetimeoffset タイプをUTCに変換します。なぜなら、矢印列は固定のタイムゾーンを必要とするからです。 ドライバーは変換時にMicrosoft SQLからUTCへのセル単位のタイムゾーン情報を正規化します。

sql_variant型はArrowのfetchメソッドではサポートされておらず、サポートされていないデータ型例外を発生させます。 fetchone()列を返すクエリには標準のfetchmany()fetchall()、またはsql_variantを使用します。

パフォーマンスに関する考慮事項

矢印フェッチ手法は分析や大量データ処理に最も高速であり、標準カーソル方式は小規模な結果セットを持つトランザクションパターンに適しています。

Arrow と標準の fetch のどちらを使うべきか

Scenario 推奨される方法
表示用に数行を取得する fetchone() / fetchall()
pandasやPolarsにデータを読み込む cursor.arrow()
大規模なデータセットをチャンク単位で処理します cursor.arrow_reader()
単一行の検索または結果セットが小さい場合 fetchone() / fetchval()
分析または集約パイプライン cursor.arrow() + Polars/DuckDB
結果をParquetまたはArrow IPCに書き込む cursor.arrow_reader() + PyArrow I/O

大規模データセットのメモリ管理

使用可能なメモリを超える可能性がある結果セットには、妥当なbatch_sizeを指定してarrow_reader()を使用してください。

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)