全部产品
Search
文档中心

向量检索服务 Milvus 版:通过Embedding Cache与AI Batch降低阿里云Milvus的AI调用成本

更新时间:Aug 13, 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 在夜间完成。