Alibaba Cloud Milvus AI Function は、コスト削減のための 2 つの機能を提供します:Embedding Cache は重複コンテンツに対して既存のベクターを再利用し、AI Batch はオフラインで待機可能な大規模ジョブを処理します。このチュートリアルでは、両方の機能を設定し、それぞれの適用箇所を判断し、信頼性の高い基準を用いてキャッシュヒットを確認します。
ソリューション概要
AI アプリケーションが本番稼働すると、通常、以下の 2 種類の無駄からコストプレッシャーが生じます:
同一コンテンツが何度も埋め込まれる — プロダクトのタイトルが繰り返し同期されたり、サポートのよくある質問が再公開されたり、ナレッジベースの段落がリトライや増分インポートによって埋め込みモデルに再入力されたりします。テキストは変更されていないため、通常ベクターも変更されませんが、モデルの呼び出しとその待機時間は同様に発生します。
待機可能なジョブがリアルタイム呼び出しを使用する — リアルタイム呼び出しは、AI Batch が実行するオフラインバッチとは対照的に、各結果をレスポンスで返す同期パスです。プロダクト説明の入力、過去の会話の要約、コンテンツの翻訳、ナレッジベースの初期化などは、すべて夜間に実行し、翌朝に配信してもユーザーエクスペリエンスに影響はありません。これらのジョブをリアルタイム呼び出しで送信すると、オンラインの検索や Q&A が必要とするモデルクォータを消費しながら、リアルタイムの料金を支払うことになります。
これらに対応する 2 つのコスト削減パスは、同一コンテンツの埋め込みを一度だけにすることと、待機可能なジョブをオフラインバッチで処理することです。Alibaba Cloud Milvus AI Function は、これら 2 つのケースをそれぞれ Embedding Cache と AI Batch でカバーします:
| 機能 | 機能内容 | 削減される無駄 |
| Embedding Cache | テキストと呼び出しコンテキストの完全一致により既存のベクターを検索し、ヒットした場合に再利用します。 | 同一コンテンツの埋め込みが一度だけになり、重複するトークンと待機時間を節約します。 |
| AI Batch | リクエストは JSONL に書き込まれて非同期で送信され、プラットフォームがオフラインでバッチ処理します。 | 待機可能な大規模ジョブがより低い単価で完了し、オンラインクォータを消費しません。 |
この 2 つの機能は、異なる種類の無駄に対応します。以下の表は、ビジネス特性と各機能の対応関係を示しています:
| ビジネス特性 | Embedding Cache | AI Batch |
| データパターン | 同一テキストが繰り返し出現する | 大量のデータが初めて処理される |
| 応答要件 | オンラインで返される | 後で完了してもよい |
| コスト削減の源泉 | 重複するモデル呼び出しの削減 | より低い単価でのオフラインバッチ実行 |
| 一般的なシナリオ | よくある質問とグラウンドトゥルースの回答、人気プロダクトのタイトルと属性、繰り返しインポートされるナレッジベースの段落、高頻度の検索語 | アーカイブされたドキュメントの要約やタグの生成、既存プロダクトデータのバッチ処理、アセット説明のバッチ生成、モデル評価とデータラベリング、定期的な再構築と夜間ジョブ |
どちらのパスを選択するかは、次の 2 つの質問に答えることで判断できます:
ユーザーは結果を待っていますか? はいの場合、リアルタイム呼び出しを使用します。いいえの場合で、かつボリュームが大きい場合は、AI Batch を検討してください。
コンテンツは繰り返し出現しますか? はいの場合、Embedding Cache を有効にします。すべてのコンテンツが初めて出現する場合、キャッシュのメリットは限定的なので、AI Batch と呼び出し回数の制御に重点を置きます。
チャットの返信や検索のオートコンプリートなどのインタラクティブなシナリオでは、デフォルトでリアルタイム呼び出しが引き続き使用されます。ユーザーが結果を待っている限り、リクエストを AI Batch に移行しないでください。
この 2 つの機能は相互排他的ではなく、本番環境では両方を組み合わせるのが一般的な方法です。e コマースのシナリオを考えてみましょう。新しいプロダクトはリアルタイムで埋め込まれ続け、人気プロダクトの繰り返し同期では Embedding Cache を通じてベクターが再利用されます。数百万件のアーカイブされたプロダクトの要約とタグは AI Batch に渡され、夜間に完了します。
前提条件
Milvus 2.6 インスタンス。AI Function は 2.6 カーネルに依存しており、インスタンス作成後に別途モデルサービスをバインドする必要はありません。
インターネット経由でインスタンスにアクセスするには、インスタンス詳細ページの Security Configuration タブで Public Access を有効にし、クライアントの出口 IP アドレスをパブリックアクセスホワイトリストに追加します。
Embedding Cache の例を実行するために pymilvus がインストールされていること。このトピックの例は、pymilvus 3.0.0 で検証されています。AI Batch の例は RESTful API を呼び出し、Python 標準ライブラリのみを使用します。
RESTful API は gRPC とポート 19530 を共有するため、呼び出す際には明示的にポートを指定してください。例:http://c-xxx.milvus.aliyuncs.com:19530。ポートを省略すると、リクエストはデフォルトでポート 80 に送信され、接続がタイムアウトします。
Embedding Cache:同一コンテンツの埋め込みを一度に
Embedding Cache の仕組みと適用範囲
Embedding Cache は、まずテキストと呼び出しコンテキストによって既存のベクターを検索し、ヒットした場合は直接再利用し、初めて出現するコンテンツに対してのみモデルを呼び出します。キャッシュは完全一致を使用するため、記述が異なる 2 つのテキストは別々に計算されます。したがって、メリットは完全にコンテンツの重複率に依存します。重複が多ければ多いほど、より多くのモデルリクエストを節約できます。おおよその見積もりについては、「コスト見積もり」をご参照ください。
キャッシュが利用できない場合や読み取りがタイムアウトした場合でも、Milvus はモデルの呼び出しを続行するため、書き込みやクエリが中断されることはありません。
Embedding Cache の有効化
Embedding Function の params に cache 構成を追加します。コンテンツの変更頻度に応じて ttl_hours を設定します。プロダクトのタイトルやよくある質問など、変更頻度の低いコンテンツには長く設定し、変更の速いコンテンツには短く設定して、新しいバージョンがより早く再計算されるようにします。
import json
import uuid
from urllib.error import HTTPError
from urllib.request import Request, urlopen
from pymilvus import DataType, Function, FunctionType, MilvusClient
MILVUS_URI = "http://c-xxx.milvus.aliyuncs.com:19530" # ポートは 19530 である必要があります
MILVUS_TOKEN = "root:xxx"
MODEL_NAME = "text-embedding-v4"
VECTOR_DIM = 1024
# キャッシュ構成:完全一致、Redis バックエンド、24 時間の TTL
CACHE_CONFIG = json.dumps({
"enabled": True,
"exact_cache": {
"enabled": True,
"backend": "redis",
"ttl_hours": 24,
},
})
client = MilvusClient(uri=MILVUS_URI, token=MILVUS_TOKEN)
collection_name = "ai_embedding_cache_demo"
if client.has_collection(collection_name):
client.drop_collection(collection_name)
schema = MilvusClient.create_schema(auto_id=True, enable_dynamic_field=False)
schema.add_field("id", DataType.INT64, is_primary=True)
schema.add_field("content", DataType.VARCHAR, max_length=4096)
schema.add_field("embedding", DataType.FLOAT_VECTOR, dim=VECTOR_DIM)
schema.add_function(
Function(
name="embed_content_with_cache",
function_type=FunctionType.TEXTEMBEDDING,
input_field_names=["content"],
output_field_names=["embedding"],
params={
"provider": "aliyun_milvus",
"model_name": MODEL_NAME,
"dim": VECTOR_DIM,
"cache": CACHE_CONFIG, # ← Embedding Cache を有効化
},
)
)
index_params = client.prepare_index_params()
index_params.add_index(field_name="embedding", index_type="AUTOINDEX", metric_type="COSINE")
client.create_collection(collection_name=collection_name, schema=schema,
index_params=index_params)
# 書き込み中にキャッシュにヒットしたコンテンツは、モデル呼び出しを生成しません
row = {"content": "Milvus is an open-source vector database."}
client.insert(collection_name, [row])
client.flush(collection_name)書き込み中にキャッシュにヒットしたコンテンツは、モデル呼び出しを生成しません。キャッシュが有効になっていることを確認するには、「キャッシュヒットの確認」をご参照ください。
キャッシュヒットの確認
キャッシュヒットは、RESTful 埋め込みエンドポイントのレスポンスから確認します。このレスポンスは usage.total_tokens と request_id を返します。キャッシュヒットの場合、実際にはモデル呼び出しは発生しないため、トークン消費量はゼロになり、モデル側でリクエスト ID は生成されません。
ベクター比較も書き込みレイテンシも、キャッシュヒットの有効な基準ではありません。このステップは間違いやすいので注意してください。
直感的に考えられる 2 つの確認方法は機能しません:
ベクター比較 — モデルは同じ入力に対して同じベクターを返します。このトピックのテストでは、キャッシュを完全に無効にして同じテキストを 2 回書き込んでも、要素ごとの差は 0 でした。
書き込みレイテンシ — 1 回の
insertとflushの固有のオーバーヘッドは、1 回のモデル呼び出しをはるかに超えており、テストでは 2 回目の書き込みが 1 回目よりも遅くなることさえありました。
以下のコードは、同じ新しいテキストを 3 回送信し、各リクエストの結果を報告します。前のステップで定義されたインポートと定数を再利用するため、同じスクリプトまたはセッションで実行してください。
# ==================== キャッシュが実際にヒットしたかどうかの確認 ====================
# 基準:usage.total_tokens がゼロになるか、request_id が空になるか。
# 両方が真の場合、モデル呼び出しは実際には発生しておらず、キャッシュヒットを意味します。
def post_json(path, body, timeout=180):
request = Request(
f"{MILVUS_URI.rstrip('/')}{path}",
data=json.dumps(body, ensure_ascii=False).encode("utf-8"),
headers={"Authorization": f"Bearer {MILVUS_TOKEN}",
"Content-Type": "application/json"},
method="POST",
)
try:
with urlopen(request, timeout=timeout) as response:
return response.status, json.loads(response.read().decode("utf-8"))
except HTTPError as exc:
return exc.code, json.loads(exc.read().decode("utf-8"))
# 最初のクエストが確実にミスするように、全く新しいテキストを使用します
text = f"cache hit verification {uuid.uuid4().hex[:12]}"
body = {
"model_name": MODEL_NAME,
"texts": [text],
"params": {"dim": VECTOR_DIM, "cache": CACHE_CONFIG},
}
for i in (1, 2, 3):
status, data = post_json("/v2/vectordb/ai/embedding", body)
assert status == 200 and data.get("code") == 0, data
usage = data["data"].get("usage", {})
request_id = data["data"].get("request_id", "")
total_tokens = usage.get("total_tokens")
hit = (total_tokens == 0) and not request_id
print(f"リクエスト {i}:total_tokens={total_tokens} "
f"request_id={'(空)' if not request_id else request_id} "
f"-> {'キャッシュヒット' if hit else 'キャッシュミス、モデルが呼び出されました'}")以下の表は、同じ新しいテキストに対する 3 回連続のリクエストの測定結果です:
| 注文 | usage.total_tokens | request_id | 結果 |
| リクエスト 1 | 20 | 値あり | キャッシュミス、モデルが実際に呼び出されました |
| リクエスト 2 | 0 | 空 | キャッシュヒット |
| リクエスト 3 | 0 | 空 | キャッシュヒット |
両方のフィールドはレスポンスボディで直接利用でき、追加の権限は必要ありません。サーバー側で相互検証するには、コンソールで埋め込みモデルの呼び出し回数がどのように変化するかを監視します。
AI Batch:待機可能なジョブのオフライン処理
AI Batch の仕組みと適用範囲
AI Batch は、JSONL に書き込まれたリクエストを受け取り、非同期で受理し、オフラインでバッチ処理します。e コマースにおける典型的な運用モデルでは、日中のモデルクォータは検索や質問を行うユーザーのために確保されます。そして、その日に蓄積されたプロダクトデータやサポートの会話は夜間に処理されます。例えば、プロダクト説明の一括生成、タグの入力、長い会話の要約などです。このデータに対する結果は、リアルタイム料金を支払うことなく、翌営業日までに業務システムに書き戻されます。
入力ファイルの準備とバッチジョブの実行
入力は JSONL 形式で、1 行が 1 つのジョブに対応します。custom_id を使用して各ジョブを元のデータにリンクさせることで、結果を特定のレコードにマッピングできます:
{"custom_id":"article-001","method":"POST","url":"/v1/chat/completions","body":{"model":"qwen3.7-max","messages":[{"role":"user","content":"Generate a one-sentence summary for this knowledge base article..."}],"enable_thinking":false}}
{"custom_id":"article-002","method":"POST","url":"/v1/chat/completions","body":{"model":"qwen3.7-max","messages":[{"role":"user","content":"Generate a one-sentence summary for this knowledge base article..."}],"enable_thinking":false}}各行の body.model の値は、アップロード時に宣言された model_name と一致している必要があります。一致しない場合、ジョブは失敗します。
全プロセスは、JSONL ファイルのアップロード、ジョブの作成、ステータスのポーリング、結果のダウンロードの 4 つのステップで構成されます。AI Batch は RESTful API を使用するため、以下のブロックでは、認証済みリクエスト、マルチパートアップロード、結果ダウンロードのための定数とヘルパー関数を定義します。4 つのステップそれぞれでこれらの関数を再利用するため、最初にこのブロックを実行してください。
ステップ 1:入力ファイルのアップロード
input.jsonl をアップロードし、返された input_file_id を保持します。これは、次のステップで作成するジョブの入力を識別します。
# 1) JSONL 入力ファイルをアップロード
status, data = upload_input_file("input.jsonl")
input_file_id = (data.get("data") or {}).get("id")
if status != 200 or not input_file_id:
raise SystemExit(f"アップロードに失敗しました:HTTP={status} message={data.get('message')}")
print(f"input_file_id = {input_file_id}")ステップ 2:バッチジョブの作成
input_file_id からジョブを作成し、完了ウィンドウを宣言します。テストでは、この種のジョブは作成から完了まで数時間かかることがあるため、completion_window は通常 24h に設定されます。
# 2) バッチジョブを作成し、完了ウィンドウを宣言
status, data = post_json("/v2/vectordb/ai/batch/jobs/create", {
"provider": PROVIDER,
"input_file_id": input_file_id,
"endpoint": ENDPOINT,
"completion_window": "24h",
})
batch_id = (data.get("data") or {}).get("id")
if status != 200 or not batch_id:
raise SystemExit(f"作成に失敗しました:HTTP={status} message={data.get('message')}")
print(f"batch_id = {batch_id}")ステップ 3:ジョブステータスのポーリング
バッチジョブは時間単位の非同期ジョブであるため、フォアグラウンドで集中的にポーリングしないでください。テストでは、送信されたジョブは長時間 in_progress 状態のままです。本番環境では、送信後に batch_id を記録し、ビジネススレッドをブロックするのではなく、定期タスクから後でポーリングします。
# 3) ジョブステータスをポーリングします。バッチジョブは時間単位で非同期に実行されます。本番環境では、
# batch_id を記録し、フォアグラウンドで集中的にポーリングするのではなく、定期タスクから後でポーリングします。
while True:
status, data = post_json("/v2/vectordb/ai/batch/jobs/describe",
{"provider": PROVIDER, "batch_id": batch_id})
batch = data.get("data") or {}
batch_status = batch.get("status")
print(f"{batch_status} {batch.get('request_counts')}")
if batch_status in {"completed", "failed", "expired", "cancelled"}:
break
time.sleep(300) # 5 分に 1 回のポーリングで十分ですステップ 4:結果のダウンロード
status が completed になるまで結果をダウンロードしないでください。失敗した行は error_file_id に収集され、別途リトライできます。
# 4) 完了後に結果をダウンロードします。失敗した行は error_file_id にあり、別途リトライできます
if batch_status == "completed":
download_batch_file(batch_id, "output", Path("output.jsonl"))
if batch.get("error_file_id"):
download_batch_file(batch_id, "error", Path("error.jsonl"))結果の確認と利用
以下の表は、4 つの操作の測定されたレスポンスを示しています:
| ステップ | 操作 | レスポンス |
| アップロード | /v2/vectordb/ai/batch/files/upload | file-batch-xxx 形式の input_file_id |
| 作成 | /v2/vectordb/ai/batch/jobs/create | batch_xxx 形式の batch_id |
| ポーリング | /v2/vectordb/ai/batch/jobs/describe | status と request_counts(合計/完了/失敗カウントを含む) |
| ダウンロード | /v2/vectordb/ai/batch/files/content | 結果ファイルのストリーム。ジョブが完了していない場合、output_file_id is empty を報告します |
ダウンロード操作は、ファイル ID ではなく batch_id と file_type でファイルを特定します。これは間違いやすい点です。
ジョブが完了した後も、結果ファイルには元の custom_id が保持されているため、アプリケーションは各要約を対応するレコードに書き戻すことができます。データの一部が失敗した場合は、error_file_id に対応するエラーファイルのみをダウンロードし、失敗した行をリトライします。完了した結果はそのまま保持され、バッチ全体を再実行する必要はありません。
コスト見積もり
Embedding Cache と AI Batch は異なる方法でコストを削減するため、見積もり方法も異なります:
| 機能 | 見積もり基準 | 効果 |
| Embedding Cache | 1 日あたり 100,000 回の埋め込み呼び出し、1 呼び出しあたり平均 50 トークン、キャッシュヒット率 50% | キャッシュにヒットした部分はモデル呼び出しを生成せず、全体のトークン消費量が約半分に減少します |
| AI Batch | 1 日あたり 100,000 エントリ、入力 50 トークン/出力 500 トークン | オフラインバッチ実行のため、モデル呼び出し料金がより低い単価で計算されます |
これらの数値は、コスト削減の源泉とその規模を示すためのものです。実際の料金は、モデル、リージョン、現在の価格に依存し、通常、入力と出力では単価が異なるため、公式の料金ページと実際の請求書をご参照ください。
トラブルシューティング
以下の表は、Embedding Cache と AI Batch を設定する際に最も発生しやすいエラーをまとめたものです:
| 症状 | 原因 | 対処法 |
| RESTful API を呼び出すと接続がタイムアウトする。 | ポートが省略されているため、リクエストがデフォルトでポート 80 に送信される。 | ポート 19530 を明示的に指定します。例:http://c-xxx.milvus.aliyuncs.com:19530。 |
| バッチジョブが失敗する。 | input.jsonl の行の body.model の値が、アップロード時に宣言された model_name と一致しない。 | すべての行の body.model を、ファイルのアップロード時に宣言した model_name と一致させます。 |
ダウンロードリクエストが output_file_id is empty を報告する。 | ジョブが完了していない。 | status が completed になるまで /v2/vectordb/ai/batch/jobs/describe をポーリングし、その後で結果をダウンロードします。 |
| ダウンロードリクエストが結果ファイルを見つけられない。 | リクエストがファイル ID でファイルを識別している。 | batch_id と file_type を使用して /v2/vectordb/ai/batch/files/content を呼び出します。 |
| 同じテキストを 2 回書き込んでも、返されるベクターが同一であるか、2 回目の書き込みが 1 回目より遅い。 | ベクター比較と書き込みレイテンシは、キャッシュヒットの基準ではない。 | RESTful 埋め込みエンドポイントのレスポンスで usage.total_tokens と request_id を確認します。 |