本文介绍阿里云 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 在夜间完成。