「.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グルーピング+並列バッチ
- レコードをパーティションキー値でグループ化。
- 各グループごとにTransactional Batch(最大100件・1.2MB)を作成。
- バッチを非同期で一斉送信(セマフォで同時実行を制御)。
# 例: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),
}
このように、測る→上げる(あるいは下げる)→再測を繰り返すことで、環境・データ特性に合った最適解に到達できます。

コメント