全部產品
Search
文件中心

Vector Retrieval Service for Milvus:通過Embedding Cache與AI Batch降低阿里雲Milvus的AI調用成本

更新時間:Aug 14, 2026

本文介紹阿里雲 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 '未命中,調用了模型'}")

對同一段新文本連續請求三次的實測結果:

次序

usage.total_tokens

request_id

判定

第 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"))

四個介面的實測返回:

步驟

介面

返回

上傳

/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(含 total / completed / failed 計數)

下載

/v2/vectordb/ai/batch/files/content

結果檔案流;任務未完成時報 output_file_id is empty

警告

下載介面用 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

離線批量執行,模型調用費用按較低單價計算

警告

以上僅用於說明節省來源與量級關係。實際費用取決於模型、地區與當期價格,且輸入與輸出通常單價不同,請以官方價格頁與實際賬單為準。

選型可以按兩個問題判斷:

  1. 使用者是否正在等待結果?在等——走即時調用;不在等且量大——考慮 Batch。

  2. 內容是否會重複出現?會重複——開啟 Embedding Cache;全是首次出現的新內容——緩衝收益有限,重點放在 Batch 與調用量控制上。

兩者並不互斥。落到電商情境:新商品上架繼續即時向量化,熱門商品重複同步由 Embedding Cache 複用向量,幾百萬條歷史商品的摘要與標籤交給 Batch 在夜間完成。