本文介紹阿里雲 Milvus AI Function 的兩項降本能力:Embedding Cache 讓重複內容複用已有向量,AI Batch 把可以等待的大批量任務離線處理。文中給出兩項能力的配置方法、適用邊界,以及判斷緩衝是否真正命中的可靠判據。
方案概述
AI 應用上線後,成本壓力通常來自兩類浪費:
相同內容被反覆向量化:商品標題會反覆同步,客服 FAQ 會多次發布,知識庫段落也會因為重試或增量匯入再次進入 Embedding 模型。文本沒變,向量通常也不變,但模型調用和等待時間又實實在在發生了一次。
可以等待的任務卻走了即時介面:商品描述補全、歷史會話摘要、內容翻譯、知識庫初始化,今晚處理、明早交付完全不影響使用者體驗。這類任務走即時介面既按即時價格付費,又會擠佔在線搜尋與問答的模型額度。
對應的兩條降本路徑是:相同內容只做一次向量化,以及可以等待的任務離線批量處理。阿里雲 Milvus AI Function 分別用 Embedding Cache 與 AI Batch 覆蓋這兩類情境:
能力 | 作用 | 解決的浪費 |
Embedding Cache | 按文本與調用上下文精確匹配已有向量,命中即複用。 | 相同內容只做一次向量化,省下重複的 Token 與等待時間。 |
AI Batch | 請求寫入 JSONL 後非同步提交,由平台離線批量執行。 | 可等待的大批量任務以更低單價完成,且不擠佔在線額度。 |
兩者面向的是兩種不同的浪費,選擇依據很清晰:
業務特徵 | Embedding Cache | AI Batch |
資料特點 | 相同文本反覆出現 | 首次處理的大批量資料 |
返回要求 | 線上返回 | 允許延後完成 |
節省來源 | 減少重複的模型調用 | 離線批量執行,單價更低 |
常見情境 | FAQ 與標準答案、熱門商品標題與屬性、知識庫重複匯入的段落、高頻搜尋字詞 | 歷史文檔產生摘要或標籤、存量商品資料批量加工、素材批量產生描述、模型評測與資料標註、周期性重建與夜間任務 |
聊天回複、搜尋補全等即時互動仍應使用即時調用——只要使用者正在等待結果,就不該為了折扣把請求放進 Batch。
前提條件
已建立 Milvus 2.6 版本執行個體。AI Function 依賴 2.6 版本核心,建立後無需單獨綁定模型服務。
如需從公網訪問執行個體,已在執行個體詳情頁的 安全配置 頁簽開啟 公網訪問 並將用戶端出口 IP 加入公網訪問白名單。
已安裝 pymilvus,本文樣本基於 pymilvus 3.0.0 驗證。
RESTful 介面與 gRPC 共用 19530 連接埠,調用時必須顯式帶連接埠,例如 http://c-xxx.milvus.aliyuncs.com:19530;省略連接埠會預設訪問 80 連接埠並導致連線逾時。
Embedding Cache:讓相同內容只向量化一次
工作方式與適用範圍
Embedding Cache 會先按文本與調用上下文尋找已有向量,命中後直接複用,只有首次出現的新內容才真正調用模型。緩衝採用精確匹配——兩段文字只要寫法不同,就會分別計算。因此它的收益完全來自內容重複率:重複越多,省下的模型請求越多。
以日均 10 萬次向量化、平均每次 50 Token、快取命中率 50% 估算,命中的那一半請求不再產生任何模型調用與等待,整體 Token 消耗大約下降一半。
緩衝不可用或讀取逾時時,Milvus 會繼續調用模型,業務寫入與查詢不會因此中斷。
開啟緩衝
在 Embedding Function 的 params 中增加 cache 配置即可。ttl_hours 按內容更新頻率設定:商品標題、FAQ 這類更新不頻繁的內容可適當延長,內容變化快時縮短,讓新版本更快重新計算。
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 小時
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)如何確認緩衝真正命中
這一步容易踩坑,需要特別說明。比較兩次寫入得到的向量是否一致,無法判斷緩衝是否命中——Embedding 模型對相同輸入本身就是確定性輸出,實測在完全不開啟緩衝的情況下,兩次寫入同一文本得到的向量逐元素差值同樣為 0。同理,寫入耗時也不能作為判據:一次 insert 加 flush 的固有開銷遠大於單次模型調用,實測第二次寫入甚至可能比第一次更慢。
可靠的判據是響應中的 usage.total_tokens 與 request_id:命中緩衝時本次沒有真正調用模型,因此 Token 消耗歸零、且不會產生模型側的請求 ID。
# ==================== 驗證緩衝是否真正命中 ====================
# 判據: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"快取命中驗證 {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 '未命中,調用了模型'}")對同一段新文本連續請求三次的實測結果:
次序 |
|
| 判定 |
第 1 次 | 20 | 有值 | 未命中,真實調用了模型 |
第 2 次 | 0 | 空 | 命中緩衝 |
第 3 次 | 0 | 空 | 命中緩衝 |
這兩個欄位在響應體中直接可得,無需額外許可權。如果還想從服務端側交叉驗證,可在控制台觀察 Embedding 模型的調用次數變化。
AI Batch:把可以等待的任務離線處理
典型情境
一家電商白天把模型額度留給正在搜尋和諮詢的使用者,夜裡再集中處理當天積累的商品資料與客服會話:批量產生商品描述、補齊標籤、把長對話整理成摘要。第二天營運上班時結果已寫回業務系統——使用者沒有多等一秒,也不必按即時價格處理這批資料。
準備輸入檔案
輸入為 JSONL,一行對應一條任務,用 custom_id 關聯未經處理資料,便於結果回寫時對應到具體記錄:
{"custom_id":"article-001","method":"POST","url":"/v1/chat/completions","body":{"model":"qwen3.7-max","messages":[{"role":"user","content":"請為這篇知識庫文章產生一句話摘要……"}],"enable_thinking":false}}
{"custom_id":"article-002","method":"POST","url":"/v1/chat/completions","body":{"model":"qwen3.7-max","messages":[{"role":"user","content":"請為這篇知識庫文章產生一句話摘要……"}],"enable_thinking":false}}每行 body.model 必須與上傳時聲明的 model_name 一致,否則任務會失敗。
提交與回查
完整流程為四步:上傳 JSONL、建立任務、回查狀態、下載結果。
import json
import shutil
import tempfile
import time
import uuid
from pathlib import Path
from urllib.error import HTTPError
from urllib.request import Request, urlopen
# AI Batch 走 RESTful 介面,base_url 用執行個體地址(連接埠 19530)
MILVUS_BASE_URL = "http://c-xxxx.milvus.aliyuncs.com:19530"
MILVUS_TOKEN = "root:xxx"
MODEL_NAME = "qwen3.7-max" # 必須與 input.jsonl 裡每行 body.model 一致
PROVIDER = "aliyun_milvus"
ENDPOINT = "/v1/chat/completions"
def post_json(path, body, timeout=180):
request = Request(
f"{MILVUS_BASE_URL.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"))
def upload_input_file(input_file):
"""multipart 上傳 input.jsonl"""
boundary = f"----milvus-ai-batch-{uuid.uuid4().hex}"
fields = {"provider": PROVIDER, "model_name": MODEL_NAME,
"endpoint": ENDPOINT, "purpose": "batch"}
with tempfile.TemporaryFile(mode="w+b") as payload:
for name, value in fields.items():
payload.write(f"--{boundary}\r\n".encode())
payload.write(f'Content-Disposition: form-data; name="{name}"\r\n\r\n'.encode())
payload.write(value.encode())
payload.write(b"\r\n")
payload.write(f"--{boundary}\r\n".encode())
payload.write(b'Content-Disposition: form-data; name="file"; filename="input.jsonl"\r\n')
payload.write(b"Content-Type: application/jsonl\r\n\r\n")
with open(input_file, "rb") as fh:
shutil.copyfileobj(fh, payload)
payload.write(b"\r\n")
payload.write(f"--{boundary}--\r\n".encode())
length = payload.tell()
payload.seek(0)
request = Request(
f"{MILVUS_BASE_URL.rstrip('/')}/v2/vectordb/ai/batch/files/upload",
data=payload,
headers={"Authorization": f"Bearer {MILVUS_TOKEN}",
"Content-Type": f"multipart/form-data; boundary={boundary}",
"Content-Length": str(length)},
method="POST",
)
try:
with urlopen(request, timeout=600) 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"))
def download_batch_file(batch_id, file_type, output_file):
"""下載結果檔案。注意用 batch_id + file_type,不是 file_id"""
request = Request(
f"{MILVUS_BASE_URL.rstrip('/')}/v2/vectordb/ai/batch/files/content",
data=json.dumps({"provider": PROVIDER, "batch_id": batch_id,
"file_type": file_type}, ensure_ascii=False).encode("utf-8"),
headers={"Authorization": f"Bearer {MILVUS_TOKEN}",
"Content-Type": "application/json"},
method="POST",
)
with urlopen(request, timeout=600) as response, output_file.open("wb") as out:
shutil.copyfileobj(response, out)
# 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) 建立 Batch 任務,聲明完成視窗
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) 回查任務狀態。Batch 為小時級非同步任務,生產環境建議記錄 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 分鐘回查一次即可
# 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"))四個介面的實測返回:
步驟 | 介面 | 返回 |
上傳 |
| 返回 |
建立 |
| 返回 |
回查 |
|
|
下載 |
| 結果檔案流;任務未完成時報 |
下載介面用 batch_id 加 file_type 定位檔案,而不是用檔案 ID,這一點容易弄錯。
Batch 是小時級非同步任務,請不要在前台密集輪詢。實測提交後任務會持續處於 in_progress 狀態較長時間,同類任務從建立到完成可能需要數小時,completion_window 通常設為 24h。生產環境建議提交後記錄 batch_id,由定時任務稍後回查,而不是讓業務線程阻塞等待。
任務完成後結果檔案仍帶回原來的 custom_id,業務可以把摘要準確寫回對應記錄。部分資料失敗時,只需下載 error_file_id 對應的錯誤檔案並重試失敗行,已完成的結果直接保留,無需整批重做。
成本估算與選型建議
兩項能力的節省來源不同,估算方式也不同:
能力 | 估算口徑 | 效果 |
Embedding Cache | 日均 10 萬次向量化、平均每次 50 Token、命中率 50% | 命中部分不產生模型調用,整體 Token 消耗約降一半 |
AI Batch | 日均 10 萬條、輸入 50 / 輸出 500 Token | 離線批量執行,模型調用費用按較低單價計算 |
以上僅用於說明節省來源與量級關係。實際費用取決於模型、地區與當期價格,且輸入與輸出通常單價不同,請以官方價格頁與實際賬單為準。
選型可以按兩個問題判斷:
使用者是否正在等待結果?在等——走即時調用;不在等且量大——考慮 Batch。
內容是否會重複出現?會重複——開啟 Embedding Cache;全是首次出現的新內容——緩衝收益有限,重點放在 Batch 與調用量控制上。
兩者並不互斥。落到電商情境:新商品上架繼續即時向量化,熱門商品重複同步由 Embedding Cache 複用向量,幾百萬條歷史商品的摘要與標籤交給 Batch 在夜間完成。