PythonにおけるPostgreSQLのパフォーマンス最適化:asyncpgでI/Oボトルネックを解消し秒間10万件を投入する

Python tutorial - IT technology blog
Python tutorial - IT technology blog

1. 本番障害:DBクエリ急増によるバックエンドのダウン

先週火曜日の午後2時頃、自社のIoTシステムでPagerDutyのアラートが一斉に鳴り響きました。APIレイテンシはわずか5分で40msから1.8sへと急上昇。15,000台のセンサーデバイスからデータを受信するゲートウェイでタイムアウトが多発し、Redisキューには20万件以上のメッセージが滞留してしまいました。

チームの初期対応としてまずデータベースを確認しました。しかし奇妙なことに、PostgreSQLクラスター(8 vCPU、32GB RAM)にSSH接続してhtopを実行し、pg_stat_activityをチェックしたところ、DBのCPU使用率はわずか22%、RAMは半分近く空いており、ディスク書き込みも15MB/s前後でした。データベース自体は極めて健全だったのです。真の原因はPythonアプリケーション層にありました。Postgresへのクエリ実行時にコネクションプールが枯渇し、数百ものワーカーがネットワークI/O待ちでブロックされていました。

2. 従来のpsycopg2がボトルネックになった理由

このプロジェクトは当初、同期型(ブロッキングI/O)のpsycopg2を採用していました。数百req/s程度の低トラフィックでは順調でしたが、5,000req/sに達した際、以下の3つの根本的な弱点により破綻しました。

  • スレッドのブロッキング(Thread-blocking): cursor.execute()を呼び出すたびに、ワーカープロセスはデータベースの応答を待って停止します。PythonはこのI/O待機時間を他のリクエスト処理に活用できません。
  • 接続ハンドシェイクのオーバーヘッド: PostgreSQLの新規接続には、TLSハンドシェイクとユーザー認証で30〜50msかかります。アプリが頻繁に接続を作成・破棄すると、無駄な接続処理だけでサーバーのCPUリソースを消費し尽くしてしまいます。
  • テキストプロトコルのオーバーヘッド: 旧来のドライバーの多くはプレーンテキスト形式でデータをやり取りします。JSONB、Timestamp、UUIDなどの複雑なデータ型をデコードする際、アプリ側のCPUが文字列パースを継続的に行う必要があり、スループットの大幅な低下を招きます。

3. 3つの解決策を徹底比較

システムを復旧させるため、チームでは以下の3つの選択肢を検証しました。

選択肢1:Gunicorn / Uvicornのワーカー数を増やす

最も手軽なアプローチですが、コストパフォーマンスが極めて悪いです。Pythonの各ワーカーは約110MBのRAMを消費します。ワーカー数を8から32に増やすと約4GBのRAMを消費する上、低速クエリが発生した際の詰まりは解消されません。これは単なる一時しのぎであり、I/O問題の根本的な解決にはなりません。

選択肢2:psycopg3の非同期(async)モードを採用する

psycopg3はasync/await構文をサポートしつつ、DB-API 2.0標準を維持しています。大規模な既存コードベースをリファクタリングする場合には安全な選択肢です。ただし、下位互換性を維持しているため、AsyncIO専用にスクラッチから開発されたドライバーと比較すると、生データの処理速度で劣ります。

選択肢3:asyncpgへ完全移行する

asyncpgはMagicStack社がCythonでゼロから開発した、PostgreSQLのパフォーマンスを極限まで引き出すためのライブラリです。DB-API 2.0を意図的に排除し、バイナリプロトコルを介してPostgresと直接通信します。実環境テストの結果、psycopg2と比較してスループットは3.5倍に向上し、アプリサーバーのCPU使用率は約40%削減されました。

4. asyncpgの実装:基本接続から大規模データ処理まで

インストール

pip経由でライブラリをインストールします:

pip install asyncpg

共有コネクションプールの初期化

各エンドポイントで個別に接続を開くのは避けてください。アプリケーション起動時に単一のプールを作成し、全リクエストで共有します:

import asyncio
import asyncpg

DATABASE_URL = "postgresql://postgres:[email protected]:5432/iot_db"

async def create_db_pool():
    return await asyncpg.create_pool(
        dsn=DATABASE_URL,
        min_size=10,        # アイドル状態の接続を常時10個維持
        max_size=30,        # インスタンスあたりの最大接続数を30に設定
        max_queries=50000,  # メモリリークを防ぐため5万クエリ後に接続を自動リフレッシュ
        timeout=10.0        # プールからの接続取得が10秒を超えた場合はタイムアウト
    )

async def get_device_by_id(pool: asyncpg.Pool, device_id: str):
    # プールから接続を1つ取得し、async withブロックを抜けた際に自動返却
    async with pool.acquire() as conn:
        row = await conn.fetchrow(
            "SELECT id, name, firmware_version, created_at FROM devices WHERE id = $1",
            device_id
        )
        return dict(row) if row else None

注意点: asyncpgでは%sや?ではなく、位置指定プレースホルダー$1, $2, $3...を使用します。

大容量データの投入:forループを避け、copy_records_to_tableを活用する

10万行のセンサーログデータをDBに書き込む際、コードの書き方次第でシステムの可用性が左右されます。forループを回して1行ずつawait conn.execute()を実行すると45秒以上かかります。executemany()を使用しても約4.2秒かかります。しかし、PostgresのバイナリCOPYコマンドを活用するcopy_records_to_tableなら、わずか0.82秒で完了します:

import time

async def bulk_insert_sensor_logs(pool: asyncpg.Pool, logs_data: list[tuple]):
    async with pool.acquire() as conn:
        start_time = time.perf_counter()
        
        # logs_dataの形式: [('dev-001', 28.5, 65.2, 1718000000), ...]
        await conn.copy_records_to_table(
            table_name="sensor_logs",
            records=logs_data,
            columns=["device_id", "temperature", "humidity", "recorded_at"]
        )
        
        elapsed = time.perf_counter() - start_time
        print(f"{len(logs_data):,} 件のレコードを {elapsed:.2f}秒 で投入しました")

安全なトランザクション管理

決済処理や在庫引き当てなどでは、データの整合性(ACID特性)の保証が不可欠です。asyncpgでは、コンテキストマネージャーを用いて簡潔にトランザクションを管理できます:

async def deduct_wallet_balance(pool: asyncpg.Pool, user_id: int, amount: float):
    async with pool.acquire() as conn:
        async with conn.transaction():
            # レースコンディションを防ぐため行をロック
            balance = await conn.fetchval(
                "SELECT balance FROM wallets WHERE user_id = $1 FOR UPDATE",
                user_id
            )
            if balance is None or balance < amount:
                raise ValueError("残高不足です")
            
            await conn.execute(
                "UPDATE wallets SET balance = balance - $1 WHERE user_id = $2",
                amount, user_id
            )
            await conn.execute(
                "INSERT INTO audit_logs (user_id, amount, action) VALUES ($1, $2, 'WITHDRAW')",
                user_id, amount
            )

5. 本番運用のベストプラクティス

  • プールサイズの算出式: 根拠なしにmax_size = 100と設定するのは避けてください。推奨される実践的な計算式は、プールサイズ = ((CPUコア数 * 2) + ディスク数) / アプリインスタンス数です。8コアNVMeのPostgresサーバーに対して4つのワーカーアプリを稼働させる場合、各ワーカーのプールは5〜10接続で十分です。
  • JSONBでjson.dumps()は不要: asyncpgはバイナリプロトコルを介して、JSONBカラムとPythonのdict/listを相互に自動デコードします。手動でのjson.loads()処理を省くことで、CPUサイクルを節約できます。
  • Graceful Shutdown(安全な終了処理): FastAPIやSanicのlifespanイベント内でSIGTERMシグナルを適切にハンドルし、await pool.close()を呼び出してください。これにより全てのソケット接続がクリーンに切断され、DBのハングを引き起こす孤立プロセス(idle in transaction)の発生を防止できます。
Share: