Azure Cosmos DB×Python:AllowBulkExecutionは使える?2025年最新動向と“実質バルク”実装レシピ徹底解説

「.NETやJavaのAllowBulkExecutionのように、Pythonでも“勝手に”クロスパーティションを束ねて超高速に書けないの?」——この疑問に2025年11月時点の正確な答えと、現場で成果を出すための代替実装・設計指針をまとめました。単なる機能の有無に留まらず、Python SDKで“実質バルク”を実現するレシピ、RU/秒の使い切り方、スロットリング対策、運用チューニングまで一気通貫で解説します。

目次

Azure Cosmos DB Python SDKでのバルク実行(AllowBulkExecution)対応

質問概要

  • .NET/Java SDKにある AllowBulkExecution(複数パーティションをまたぐ操作をSDKが自動でまとめて送信し、スループットを最適化する機能)が、Python SDK(azure‑cosmos 4.x)にも存在するか?
  • 将来的にPython SDKへ実装される予定はあるか?

最初に結論(2025年11月時点)

現状解説
サポート状況Python SDKにAllowBulkExecution相当のネイティブ機能は未実装。正式な「バルク要求(クロスPK一括送信)」機構は提供されていません。公式ドキュメントでも「Python SDKはまだバルク要求をサポートしていない」と明記されています。
公式情報/ロードマップ公開情報の範囲では、導入時期はアナウンスされていません。現状はTransactional Batch(同一パーティションキー内の一括処理)と非同期I/Oでの高並列化が推奨されます。
コミュニティ要望GitHubやQ&Aで要望は継続的に上がっていますが、優先度はユーザー需要・エンジニアリング判断に依存します。

要点:「.NET/JavaのAllowBulkExecution=クライアントSDKが送信バッチングと混雑制御を最適化する」のに対し、Pythonは同等の自動最適化は未提供。ただし非同期I/O+並列化とTransactional Batch(同一PK)を組み合わせれば、実務で十分“バルク級”のスループットを引き出せます。


なぜPython SDKにAllowBulkExecutionがないのか(背景と理解)

  • サービス機能ではなくSDK実装:BulkはREST APIの「一枚岩のエンドポイント」ではなく、SDK側でのリクエスト集約・並列送信・再試行制御の実装です。言語ごとのSDKで成熟度・優先度が異なり、.NET/Javaに先行機能が集約されています。
  • Python SDKの仕様:PythonはDirect TCPモードが未サポート(Gateway over HTTPSのみ)。ネットワーク面の差分や言語特性(GILやイベントループ)もあって、.NET/Javaと同一の内部最適化がそのまま提供されていない経緯があります。
  • 代替パターンが整備:公式は「asyncクライアントでの高い同時実行」「Transactional Batch」の活用を明示的に案内。運用チューニングと組み合わせれば高スループットを達成可能です。

現時点での代替アプローチ(実運用で効く順)

1) asyncio+azure.cosmos.aio で高並列I/O(最も手軽で効果大)

非同期クライアント azure.cosmos.aio.CosmosClient と asyncio.gather()、セマフォで同時実行数を制御しながら、大量の書き込みを“自前のバルク”として畳み掛ける方式です。429(スロットリング)時は x-ms-retry-after-ms を読み、指数バックオフ+ジッタでリトライします。

# pip install azure-cosmos
import os, asyncio, random, time
from uuid import uuid4
from typing import Iterable, Dict, Any, List
from azure.cosmos.aio import CosmosClient
from azure.cosmos import exceptions

RETRY_BASE = 0.1  # 秒
MAX_BACKOFF = 5.0 # 秒

async def upsert_with_retry(container, item: Dict[str, Any], sem: asyncio.Semaphore) -> float:
    # RU消費量を返す
    async with sem:
        attempt = 0
        while True:
            try:
                # レスポンスヘッダからRUを取るためresponse_hookを渡す
                ru_holder = {"ru": 0.0}
                def _hook(headers, _):
                    try:
                        ru_holder["ru"] = float(headers.get("x-ms-request-charge", "0"))
                    except Exception:
                        ru_holder["ru"] = 0.0

                await container.upsert_item(item, response_hook=_hook)
                return ru_holder["ru"]

            except exceptions.CosmosHttpResponseError as e:
                if getattr(e, "status_code", None) == 429:
                    # Retry-After-msヘッダーを優先利用
                    ms =  e.headers.get("x-ms-retry-after-ms") if hasattr(e, "headers") and e.headers else None
                    delay = float(ms)/1000.0 if ms else min(MAX_BACKOFF, (RETRY_BASE * (2 ** attempt)) + random.random()*0.2)
                    await asyncio.sleep(delay)
                    attempt += 1
                    continue
                raise  # 429以外はそのまま伝播

async def bulk_upsert_async(endpoint: str, key: str, db: str, container_name: str, items: Iterable[Dict[str, Any]],
                            max_concurrency: int = 256) -> Dict[str, Any]:
    async with CosmosClient(endpoint, credential=key) as client:
        database = client.get_database_client(db)
        container = database.get_container_client(container_name)

        sem = asyncio.Semaphore(max_concurrency)
        tasks = [asyncio.create_task(upsert_with_retry(container, it, sem)) for it in items]
        rus = await asyncio.gather(*tasks, return_exceptions=False)

    return {
        "count": len(rus),
        "total_ru": sum(rus),
        "avg_ru": (sum(rus) / len(rus)) if rus else 0.0
    }

# 例)100,000件の合成データを投入
def gen_items(n: int) -> List[Dict[str, Any]]:
    out = []
    for _ in range(n):
        pk = f"user#{random.randint(1, 10000)}"  # パーティション分散例
        out.append({
            "id": str(uuid4()),
            "pk": pk,
            "type": "event",
            "ts": int(time.time()*1000),
            "payload": {"v": random.random()}
        })
    return out

# 実行イメージ
# asyncio.run(bulk_upsert_async(os.environ["COSMOS_URI"], os.environ["COSMOS_KEY"],
#                              "appdb", "events", gen_items(100_000), max_concurrency=512))

ポイント

  • SDKの再利用:CosmosClientはアプリ内でシングルトン(初期化は重い)。毎リクエスト作り直しは厳禁。
  • 同時実行数:RU/秒とアイテム当たりRUコストを基に算出(後述のチューニング式を参照)。
  • 429対策:x-ms-retry-after-msに従うのが最優先。指数バックオフ+ジッタで輻輳を回避。
  • 冪等:id をクライアント生成し、upsert_item を使う。重複書き込みの整合が取りやすい。

2) Transactional Batch(同一パーティションキー)

同一PK(論理パーティション)内の複数操作を1往復・同一トランザクションで実行。100操作/1.2MBの上限(目安)があります。同一PKに偏るケース(例えば「ユーザーIDごとに一括作成・更新」)は劇的に効きます。Python SDKは同期・非同期いずれも対応しています。

# 非同期版:同一パーティションキー "user#123" に対する複数操作を一括
from azure.cosmos.aio import CosmosClient

async def exec_batch(endpoint, key, db, ctn, pk_value, docs):
    async with CosmosClient(endpoint, credential=key) as client:
        container = client.get_database_client(db).get_container_client(ctn)
        ops = []
        # PythonのTransactional Batchは「操作名+引数タプル」のリスト
        for d in docs:
            ops.append(("upsert", (d,), {}))  # 例: upsertを多数

        # patch例(条件付きの部分更新)
        ops.append(("patch", (
            docs[0]["id"],
            [{"op": "add", "path": "/labels/-", "value": "bulk"}]
        ), {}))

        # 実行(pkはバッチ全体で同一必須)
        resp = await container.execute_item_batch(ops, partition_key=pk_value)
        # 必要なら resp を検査(失敗時は CosmosBatchOperationError)
        return resp

適用のコツ

  • 「同一PKで100操作・1.2MB」の制限に収まるように分割。複数バッチを非同期で並行送信すればよい。
  • 置換(replace)は全文書を送る必要あり。更新差分はpatchが有効。

3) ThreadPool/ProcessPool を使ったI/O並列化(同期API派)

CPU負荷が小さいI/O中心の処理では、ThreadPoolExecutorで十分なスループットが得られます。マルチプロセスはプロセスごとの接続・メモリを消費するため、まずはスレッドから。

from concurrent.futures import ThreadPoolExecutor, as_completed
from azure.cosmos import CosmosClient, exceptions
import time

def upsert_one(container, item):
    attempt = 0
    while True:
        try:
            container.upsert_item(item)
            return float(container.client_connection.last_response_headers.get("x-ms-request-charge", "0"))
        except exceptions.CosmosHttpResponseError as e:
            if getattr(e, "status_code", None) == 429:
                ms = e.headers.get("x-ms-retry-after-ms") if hasattr(e, "headers") and e.headers else "100"
                time.sleep(float(ms)/1000.0)
                attempt += 1
                continue
            raise

def bulk_upsert_sync(endpoint, key, db, ctn, items, workers=64):
    client = CosmosClient(endpoint, credential=key)
    container = client.get_database_client(db).get_container_client(ctn)
    rus = []
    with ThreadPoolExecutor(max_workers=workers) as ex:
        futs = [ex.submit(upsert_one, container, it) for it in items]
        for f in as_completed(futs):
            rus.append(f.result())
    return {"count": len(rus), "total_ru": sum(rus), "avg_ru": (sum(rus)/len(rus)) if rus else 0.0}

4) 外部ツール・サービスの併用

  • Data Migration Tool(初期ロードに有効)
  • Azure Data Factory(Copyアクティビティ)でCosmos DBへコピー
  • プロキシ方式:Azure Functions/Container Appsに.NETまたはJava SDKを載せ、AllowBulkExecutionをONにした“バルク挿入API”を用意。Python側はそのAPIを叩く
  • 一時的なRU/sスケールアップ:短時間だけプロビジョンドRU/s(またはオートスケール上限)を引き上げ、投入完了後に戻す

性能チューニング実践:RU/秒を“使い切る”ための設計

1) 同時実行数の決め方(簡易式)

概算として、

target_concurrency ≈ (container_RU_per_sec / avg_item_RU) × 利用率(0.6〜0.8)
  • 例:コンテナ 50,000 RU/s、アイテム平均 10 RU → 5,000 ops/s。利用率0.7なら3,500並列相当ですが、物理パーティション数に分散される点に注意。
  • 物理パーティションはRU/sに応じ自動で増えます。PKの偏りはスロットリングの主因。pk設計を見直すか、投入順序をシャッフルしてホットパーティションを避けます。

2) アイテム当たりRUの測り方

直近の操作のヘッダ x-ms-request-charge を読めば1リクエストのRUが取れます。Python SDKは response_hook または container.client_connection.last_response_headers から取得可能です。これを平均化すると設計の精度が上がります。

3) スロットリング(429)とバックオフ

  • まず x-ms-retry-after-ms に従う。
  • 429が続く場合は同時実行数を段階的に下げ、送出側の輻輳制御を行う(AIMDのイメージ)。
  • 429の多発はホットパーティションやRU不足のサイン。PK設計/RU増強/投入順序の分散で改善。

4) インデックス/ドキュメント設計

  • インデックス最適化:不要な深いパスや巨大配列は除外し、書き込みRUを削減。includedPaths/excludedPaths を適切化。
  • ドキュメントサイズ:大きいほどRU高騰。ペイロードは圧縮・分割(履歴は別コンテナなど)。
  • Patchの活用:部分更新でネットワーク&RUのミニマイズ。

5) ネットワーク/デプロイ

  • クライアントを同一リージョンにデプロイ(レイテンシ1〜数ms帯)。
  • Python SDKはDirect TCP未対応なので、HTTPコネクションの同時接続管理に留意(OSの nofile 上限、プロキシ制限)。

6) オートスケール/手動スケールの使い分け

  • ピークが読めない/バーストが強烈:オートスケール(上限RUは余裕多め)。
  • 短時間で確実にやり切る:事前に手動スケールアップして投入、完了したら戻す。

監視・計測と品質保証

  • RU/秒・429率・平均/P95レイテンシをダッシュボード化。SDKの response_hook でRU、Python側で処理時間を記録。
  • ActivityId(x-ms-activity-id)をログに残すと、サポート問い合わせ時に役立ちます。
  • 非同期大量投入は「短時間で膨大なRUを消費」するため、エミュレータ・検証環境でのドライランは必須。

設計パターン集(“実質バルク”を作る)

パターンA:PKグルーピング+並列バッチ

  1. レコードをパーティションキー値でグループ化。
  2. 各グループごとにTransactional Batch(最大100件・1.2MB)を作成。
  3. バッチを非同期で一斉送信(セマフォで同時実行を制御)。
# 例:PK = "tenantId" ごとに分割してバッチ実行
from collections import defaultdict

def group_by_pk(items, pk_field="tenantId"):
    g = defaultdict(list)
    for d in items:
        g[d[pk_field]].append(d)
    return g

async def bulk_by_pk_batches(client, db, ctn, items, pk_field="tenantId", batch_size=100, concurrency=128):
    database = client.get_database_client(db)
    container = database.get_container_client(ctn)

    sem = asyncio.Semaphore(concurrency)
    tasks = []
    for pk, docs in group_by_pk(items, pk_field).items():
        # 100件ずつの小分け
        for i in range(0, len(docs), batch_size):
            chunk = docs[i:i+batch_size]
            ops = [("upsert", (d,), {}) for d in chunk]
            async def _run(ops=ops, pk=pk):
                async with sem:
                    return await container.execute_item_batch(ops, partition_key=pk)
            tasks.append(asyncio.create_task(_run()))
    return await asyncio.gather(*tasks)

パターンB:.NET/JavaのバルクAPIを“プロキシ”利用

Pythonサービスから内部HTTPで呼べる専用のバルク挿入APIを用意(Azure Functions/Container Apps)。サンプル(.NET):

// .NET v3 SDK(AllowBulkExecution=true)
var options = new CosmosClientOptions { AllowBulkExecution = true };
using var client = new CosmosClient(endpoint, key, options);
var container = client.GetContainer("db", "items");
// 並列タスクでCreateItemAsync/UpsertItemAsyncを投げる(.NETはSDKが内部でバルク集約)

Pythonからは通常のHTTP POSTでペイロード(複数アイテム)を送り、.NET側がCosmosへ最適化して書き込みます。大規模な一括移行や高負荷ETLで有効です。


エラー処理・品質設計の要点

  • 429(Request rate too large):Retry-After-ms順守・指数バックオフ・同時実行の動的調整。
  • 冪等性:クライアント生成のid固定+upsert_item。再送でも二重作成しない。
  • 部分更新(Patch):必要最小の差分のみ転送・更新し、RUと帯域を節約。
  • 障害耐性:アプリ側で部分成果のチェックポイント(最後に成功したIDなど)を記録。再開可能に。
  • 観測可能性:成功・失敗件数、RU消費、アクティビティID、P95/P99レイテンシを収集。

よくある質問(FAQ)

Q1. PythonにAllowBulkExecutionはありますか?

A. 現時点ではありません。公式の記述でも「Python SDKはバルク要求をまだサポートしていない」とされ、.NET/Javaとは実装が異なります。

Q2. いつ入る予定ですか?

A. 公開ロードマップ上の期日の記載はなし。コミュニティの要望は把握されていますが、正式な提供時期はアナウンスされていません。

Q3. ではどうすれば“バルク級”の性能を得られますか?

A. 非同期I/Oの高並列とTransactional Batch(同一PK)を組み合わせ、429対策・PK分散・RU計測を徹底してください。初期ロードはData Migration Tool/ADF、恒常運用は.NET/Javaプロキシ方式も選択肢です。

Q4. Python SDKの制限で注意点は?

  • Direct TCPモード非対応:Gateway(HTTPS)経由のみ。
  • Transactional Batch:同一PK限定、上限100操作/1.2MB(目安)。
  • Group Byなど一部クエリ制限:集約の継続トークン制約などがあるため、設計時に考慮。

導入・移行チェックリスト

  • ☑️ PK設計は偏っていないか(ホットパーティションの兆候:特定PKの429多発)。
  • ☑️ インデックスを最小化しているか(書き込みRU節約)。
  • ☑️ Asyncで同時実行数を段階的に上げ、429/秒・RU/秒・P95を監視しながら最適点を探ったか。
  • ☑️ バッチ(同一PK)で100件上限に留めたか。サイズ1.2MBを超えないか。
  • ☑️ エラーハンドリング(429/5xx)・再開可能設計(チェックポイント)を入れたか。
  • ☑️ 初期ロードはDMT/ADF、恒常はアプリ内バルク(Async)or .NET/Javaプロキシで整理したか。

まとめ

Python SDKにAllowBulkExecutionは未実装ですが、これは「バルクが使えない」ことを意味しません。async高並列+Transactional Batchの二枚看板に、429対策・PK分散・インデックス最適化・RU計測を重ねることで、実務では“バルク級”の挿入性能を十分に引き出せます。正式なBulk機能が不可欠な要件(例:クロスPKでの自動バッチングをSDKに全委任したい)であれば、.NET/JavaのSDKを使った専用の書き込みプロキシをアーキテクチャに組み込むのが堅実です。いずれの方針でも、まず小さく計測し、RUと429の挙動を見ながら上限に寄せていくのが最短ルートです。


付録:小さなベンチマーク雛形(Async)

以下はRU・成功件数・P95を集計する最小ベンチの一例です。検証環境(エミュレータ/開発サブスクリプション)でパラメータだけ変えながら“最適同時実行数”を探るのに使えます。

import asyncio, statistics, time
from azure.cosmos.aio import CosmosClient
from azure.cosmos import exceptions

async def run_once(container, item, sem):
    start = time.perf_counter()
    ru = await upsert_with_retry(container, item, sem)  # 前掲関数を再利用
    latency = (time.perf_counter() - start) * 1000
    return ru, latency

async def micro_bench(endpoint, key, db, ctn, items, concurrency):
    async with CosmosClient(endpoint, credential=key) as client:
        container = client.get_database_client(db).get_container_client(ctn)
        sem = asyncio.Semaphore(concurrency)
        tasks = [asyncio.create_task(run_once(container, it, sem)) for it in items]
        results = await asyncio.gather(*tasks)
    rus  = [r for r,_ in results]
    lats = [l for _,l in results]
    return {
        "count": len(results),
        "total_ru": sum(rus),
        "avg_ru": sum(rus)/len(rus),
        "p95_ms": statistics.quantiles(lats, n=20)[18] if len(lats) >= 20 else max(lats),
    }

このように、測る→上げる(あるいは下げる)→再測を繰り返すことで、環境・データ特性に合った最適解に到達できます。

この記事を書いた人

実務の現場で詰まりがちなポイントを地図にするITブログ「IT trip」を運営。Windows/Office(Teams・Excel)からSQL、サーバ運用、ガジェットまで、再現性のある手順と“なぜそうなるか”を丁寧に解説します。読んだらすぐ試せること、そして迷った人の次の一歩が見えることを大切にしています。

コメント

コメントする

目次