作者:周弘懿(锦琛)
业务背景
在线客服与电话客服每天沉淀的最大一笔数据资产,就是海量的通话录音和语音留言。一个中等规模的客服中心,单日通话动辄数万通、录音时长以千小时计。这些语音里藏着最真实的客户诉求、坐席服务质量以及潜在的舆情风险,但绝大多数团队至今仍把它们当作「合规留档」的冷数据,很少真正用起来。
要把语音资产盘活,业务上有几件事必须做:
- 录音转文字:把通话、语音留言批量转写成可检索、可分析的文本,这是后续一切工作的前提。
- 建可检索知识库:把历史优质应答、产品手册、FAQ 沉淀成语义可检索的知识库,供机器人和坐席实时调用。
- FAQ 智能匹配:用户一句口语化的问题,要能准确命中标准 FAQ,而不是靠关键词碰运气。
- 情感与意图识别:识别客户在通话中的情绪波动(愤怒、不满、满意),实时预警、事后质检。
- 质检与合规脱敏:录音转写后往往包含手机号、身份证号、银行卡号等 PII,入库、分析、共享前必须脱敏。
技术挑战
传统方案要凑齐上面这套能力,通常得串联好几套异构系统,痛点集中在以下几条:
- 多系统串联,链路复杂:ASR 服务负责转文字,向量库负责语义检索,还要外接 NLP 平台做情感、分类、脱敏,再加一层重排服务。数据在多个系统间来回搬运,接口、鉴权、格式转换的胶水代码越写越厚,任何一环抖动都会拖垮整条链路。
- 录音数据合规风险高(PII):原始转写文本里混着手机号、证件号、银行卡号。数据一旦流出到外部 NLP 服务或分析库,就存在 PII 泄露风险,而脱敏往往被放到链路末端甚至遗漏。
- FAQ 召回不准,需要重排:纯向量召回(ANN)在语义相近但意图不同的问题上容易「召回但不精准」,Top-1 命中率不理想。要提升精度必须引入 rerank 阶段,又是一套额外的模型服务。
- 情感质检靠人工:传统质检靠人工抽检,覆盖率低、标准不统一、滞后严重,恶性投诉往往在客户已经升级到公开渠道后才被发现。
- 运维成本高:ASR、Embedding、Rerank、NLP 各自是独立服务,各有各的扩缩容、版本、监控和账单,团队疲于维护「系统之间的系统」。
解决方案
先介绍一下文中几个 AI Function 的作用:
函数 |
作用 |
在客服语音链路里 |
AI_AUDIO_TRANSCRIBE |
写入时自动把录音/语音留言转写成文本(无需应用侧先调 ASR) |
把海量通话录音变成可检索、可分析的文本 |
AI_PII_MASK |
自动掩码文本里的手机号、证件号、银行卡号等 PII |
让进入知识库/分析库的都是已脱敏文本,合规前置 |
AI_EMBEDDING |
写入时自动把 FAQ/问题文本转成向量(无需应用侧先调模型) |
让「实例怎么连外网」这类语义需求能被向量检索命中 |
AI_RERANK |
对召回候选按与查询的相关性重新排序(精排) |
把「最贴合意图」的 FAQ 顶到应答前面 |
AI_SENTIMENT |
判断每条会话是好评/差评/中性 |
聚合出情感标签,支撑全量质检与舆情预警 |
AI_CLASSIFY |
把会话自动归类到工单类目(账号/咨询/故障/计费) |
支撑工单自动分类与路由 |
具体AI Function用法参考官方文档
整体架构如下:
实战操作
下面把整条客服语音链路的可运行代码合并在一起。
- 录音转写(AI_AUDIO_TRANSCRIBE):把客服录音(真实可下载的签名 URL)用
qwen3-asr-flash转写为文本,enable_itn规范化口语数字。 - PII 脱敏(AI_PII_MASK):在文本进入知识库/分析库前,用
qwen3.7-max掩码手机号、证件号等敏感信息;提供即时调用与写入型 Collection两种方式。 - 建 FAQ 知识库(AI_EMBEDDING):用
text-embedding-v4把 FAQ 自动向量化入库,HNSW/COSINE索引。 - 语义检索召回:把用户口语化问题向量化,在 FAQ 库做 ANN 召回 Top-N 候选。
- 候选重排(AI_RERANK):用
qwen3-rerank精排候选,把最贴合意图的 FAQ 顶到 Top-1;提供即时调用与 search 挂载 ranker 两种方式。 - 质检与路由(AI_SENTIMENT + AI_CLASSIFY):对脱敏后的会话文本做情感识别与工单分类,支撑全量质检与自动分流。
from __future__ import annotations import json import time from typing import Any 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" # RESTful 接口与 gRPC 同在 19530 端口,需显式带端口(否则默认 80 端口会连接超时) MILVUS_REST_BASE_URL = MILVUS_URI # TEXTTRANSFORM 在部分 pymilvus 版本里没有枚举常量,这里做一次兼容兜底 TEXTTRANSFORM_FUNCTION_TYPE = 9 client = MilvusClient(uri=MILVUS_URI, token=MILVUS_TOKEN) def texttransform_function_type() -> Any: """取 TEXTTRANSFORM 的 FunctionType;老版本枚举缺失时动态补一个成员。""" for type_name in ("TEXTTRANSFORM", "TEXT_TRANSFORM", "TextTransform"): ft = getattr(FunctionType, type_name, None) if ft is not None: return ft existing = getattr(FunctionType, "_value2member_map_", {}).get(TEXTTRANSFORM_FUNCTION_TYPE) if existing is not None: return existing extension = int.__new__(FunctionType, TEXTTRANSFORM_FUNCTION_TYPE) extension._name_ = "TEXTTRANSFORM" extension._value_ = TEXTTRANSFORM_FUNCTION_TYPE FunctionType._value2member_map_[TEXTTRANSFORM_FUNCTION_TYPE] = extension FunctionType._member_map_["TEXTTRANSFORM"] = extension return extension def post_json(path: str, body: dict[str, Any], timeout: int = 120, retries: int = 3) -> tuple[int, dict[str, Any]]: """AI_PII_MASK / AI_RERANK 等 REST 即时接口统一封装(走大模型,加重试更稳)。""" last: tuple[int, dict[str, Any]] | None = None for _ in range(retries): request = Request( f"{MILVUS_REST_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: status, data = response.status, json.loads(response.read().decode("utf-8")) except HTTPError as exc: status, data = exc.code, json.loads(exc.read().decode("utf-8")) last = (status, data) if status == 200 and data.get("code") == 0: return status, data time.sleep(1) return last # type: ignore[return-value] # ==================== 步骤 1:AI_AUDIO_TRANSCRIBE 录音转写 ==================== def step_1_transcribe() -> None: print("\n" + "=" * 64) print("步骤 1 | AI_AUDIO_TRANSCRIBE 录音转写(qwen3-asr-flash)") print("=" * 64) collection_name = "cs_transcribe" 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("audio_url", DataType.VARCHAR, max_length=4096) # 音频输入字段 schema.add_field("transcript", DataType.VARCHAR, max_length=4096) # 转写输出字段 schema.add_field("dummy_vector", DataType.FLOAT_VECTOR, dim=2) schema.add_function( Function( name="transcribe_audio", function_type=texttransform_function_type(), input_field_names=["audio_url"], output_field_names=["transcript"], params={ "provider": "aliyun_milvus", "model_name": "qwen3-asr-flash", "task": "ai_audio_transcribe", "language": "zh", "enable_itn": "true", }, ) ) index_params = client.prepare_index_params() index_params.add_index(field_name="dummy_vector", index_type="AUTOINDEX", metric_type="COSINE") client.create_collection(collection_name=collection_name, schema=schema, index_params=index_params) # audio_url 必须是服务端可下载的真实地址(生产用你自己 OSS 的短时效签名 URL); # 占位符 ?Signature=... 无法下载,会报 "Failed to download multimodal content"。 # dummy_vector 需在 insert 时显式赋值(本版本不会自动补 2 维占位向量)。 audio_urls = [ "https://dashscope.oss-cn-beijing.aliyuncs.com/audios/welcome.mp3", "https://dashscope.oss-cn-beijing.aliyuncs.com/samples/audio/paraformer/realtime_asr_example.wav", ] client.insert(collection_name, [{"audio_url": u, "dummy_vector": [0.0, 0.0]} for u in audio_urls]) client.flush(collection_name) print("转写结果:") for row in client.query(collection_name, filter="", output_fields=["audio_url", "transcript"], limit=10): print(f" · {row['audio_url'].split('/')[-1]}") print(f" → {row['transcript']}") # ==================== 步骤 2:AI_PII_MASK 脱敏(即时调用 + 写入型 Collection) ==================== def step_2_pii_mask() -> None: print("\n" + "=" * 64) print("步骤 2 | AI_PII_MASK 敏感信息脱敏(qwen3.7-max)") print("=" * 64) # 2.1 REST 即时接口:文本进库前先脱敏 status, data = post_json( "/v2/vectordb/ai/pii_mask", { "model_name": "qwen3.7-max", "texts": ["您好,我的手机号是 13812345678,身份证 110101199003071234,麻烦帮我查下工单。"], "params": { "pii_types": ["PERSON", "PHONE", "ID_CARD"], "mask_char": "*", "preserve_length": True, "temperature": 0, }, }, ) assert status == 200 and data.get("code") == 0, data print("[2.1] REST 即时接口 /v2/vectordb/ai/pii_mask:") for item in data["data"]["output"]["outputs"]: print(f" · {item}") # 2.2 写入型 Collection:content 入库自动生成脱敏字段 masked collection_name = "cs_pii_mask" 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("masked", DataType.VARCHAR, max_length=4096) schema.add_field("dummy_vector", DataType.FLOAT_VECTOR, dim=2) schema.add_function( Function( name="mask_pii", function_type=texttransform_function_type(), input_field_names=["content"], output_field_names=["masked"], params={ "provider": "aliyun_milvus", "model_name": "qwen3.7-max", "task": "ai_pii_mask", "pii_types": "PERSON,PHONE,ID_CARD", "mask_char": "*", "preserve_length": "true", "temperature": "0", }, ) ) index_params = client.prepare_index_params() index_params.add_index(field_name="dummy_vector", index_type="AUTOINDEX", metric_type="COSINE") client.create_collection(collection_name=collection_name, schema=schema, index_params=index_params) client.insert(collection_name, [{"content": "您好,我的手机号是 13812345678,麻烦帮我查下工单。", "dummy_vector": [0.0, 0.0]}]) client.flush(collection_name) print("[2.2] 写入型 Collection:content 入库即自动脱敏为 masked:") for row in client.query(collection_name, filter="", output_fields=["masked"], limit=1): print(f" · {row['masked']}") # ==================== 步骤 3:AI_EMBEDDING 建 FAQ 知识库并入库 ==================== def step_3_build_faq_kb() -> None: print("\n" + "=" * 64) print("步骤 3 | AI_EMBEDDING 建 FAQ 知识库(text-embedding-v4)") print("=" * 64) collection_name = "cs_faq_kb" 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) # FAQ 问题文本 schema.add_field("answer", DataType.VARCHAR, max_length=4096) # 标准答案 schema.add_field("embedding", DataType.FLOAT_VECTOR, dim=1024) schema.add_function( Function( name="embed_content", function_type=FunctionType.TEXTEMBEDDING, input_field_names=["content"], output_field_names=["embedding"], params={ "provider": "aliyun_milvus", "model_name": "text-embedding-v4", "dim": 1024, "max_client_batch_size": 10, "max_concurrency": 1, }, ) ) index_params = client.prepare_index_params() index_params.add_index( field_name="embedding", index_type="HNSW", metric_type="COSINE", params={"M": 16, "efConstruction": 200}, ) client.create_collection(collection_name=collection_name, schema=schema, index_params=index_params) faqs = [ {"content": "如何开启 Serverless Milvus 实例的公网访问?", "answer": "在控制台实例详情页开启公网并配置白名单。"}, {"content": "忘记控制台登录密码怎么办?", "answer": "通过账号中心的找回密码流程重置。"}, {"content": "账单为什么比预期高?", "answer": "查看用量明细,重点关注计算与存储用量。"}, {"content": "如何创建 Collection 并写入向量?", "answer": "使用 create_collection 定义 schema 后 insert。"}, {"content": "实例扩容会影响在线业务吗?", "answer": "扩容为在线操作,通常不中断服务。"}, ] client.insert(collection_name, faqs) client.flush(collection_name) client.load_collection(collection_name) print(f"FAQ 知识库入库完成,共 {len(faqs)} 条,写入即自动向量化:") for f in faqs: print(f" · {f['content']}") # ==================== 步骤 4:用户问题语义检索召回 FAQ ==================== def step_4_semantic_recall() -> None: print("\n" + "=" * 64) print("步骤 4 | 语义检索召回 FAQ(ANN Top-N)") print("=" * 64) query = "我的实例怎么才能让外网连上?" result = client.search( collection_name="cs_faq_kb", data=[query], anns_field="embedding", limit=3, output_fields=["content", "answer"], ) print(f"查询语句:{query}") print("-" * 64) for rank, hit in enumerate(result[0], 1): e = hit["entity"] print(f" {rank}. [相似度 {hit['distance']:.4f}] {e['content']}") # ==================== 步骤 5:AI_RERANK 重排(即时调用 + search 挂载 ranker) ==================== def step_5_rerank() -> None: print("\n" + "=" * 64) print("步骤 5 | AI_RERANK 大模型重排(qwen3-rerank)") print("=" * 64) query = "我的实例怎么才能让外网连上?" # 5.1 REST 即时接口 status, data = post_json( "/v2/vectordb/ai/rerank", { "model_name": "qwen3-rerank", "query": query, "documents": [ "如何开启 Serverless Milvus 实例的公网访问?", "实例扩容会影响在线业务吗?", "如何创建 Collection 并写入向量?", ], "params": {"is_multimodal": True, "max_concurrency": 2, "timeout_sec": 10}, }, ) assert status == 200 and data.get("code") == 0, data print("[5.1] REST 即时接口 /v2/vectordb/ai/rerank:") ranked = sorted(data["data"]["output"]["results"], key=lambda x: x["relevance_score"], reverse=True) for rank, item in enumerate(ranked, 1): print(f" {rank}. [相关度 {item['relevance_score']:.4f}] 候选 index={item['index']}") # 5.2 在 search 时挂载 ranker,召回 + 重排一次完成 reranker = Function( name="rerank_faq", function_type=FunctionType.RERANK, input_field_names=["content"], params={ "reranker": "model", "provider": "aliyun_milvus", "model_name": "qwen3-rerank", "queries": [query], "is_multimodal": "true", "max_concurrency": 2, "timeout_sec": 10, }, ) result = client.search( collection_name="cs_faq_kb", data=[query], anns_field="embedding", limit=3, output_fields=["content", "answer"], ranker=reranker, ) print("[5.2] search() 内置 ranker:ANN 召回后由大模型二次精排:") for rank, hit in enumerate(result[0], 1): e = hit["entity"] print(f" {rank}. [重排分 {hit['distance']:.4f}] {e['content']} | 答:{e['answer']}") # ==================== 步骤 6:AI_SENTIMENT 情感识别 + AI_CLASSIFY 工单分类 ==================== def step_6_qc_and_route() -> None: print("\n" + "=" * 64) print("步骤 6 | AI_SENTIMENT 情感识别 + AI_CLASSIFY 工单分类(qwen3.7-max)") print("=" * 64) collection_name = "cs_qc" 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("sentiment", DataType.VARCHAR, max_length=64) # 情感输出 schema.add_field("category", DataType.VARCHAR, max_length=64) # 工单类目输出 schema.add_field("dummy_vector", DataType.FLOAT_VECTOR, dim=2) schema.add_function( Function( name="analyze_sentiment", function_type=texttransform_function_type(), input_field_names=["content"], output_field_names=["sentiment"], params={ "provider": "aliyun_milvus", "model_name": "qwen3.7-max", "task": "ai_sentiment", "categories": "positive,negative,neutral", "temperature": "0", }, ) ) schema.add_function( Function( name="classify_ticket", function_type=texttransform_function_type(), input_field_names=["content"], output_field_names=["category"], params={ "provider": "aliyun_milvus", "model_name": "qwen3.7-max", "task": "ai_classify", "labels": "账号,咨询,故障,计费", "prompt": "根据客户问题所属主题分类。", "temperature": "0", }, ) ) index_params = client.prepare_index_params() index_params.add_index(field_name="dummy_vector", index_type="AUTOINDEX", metric_type="COSINE") client.create_collection(collection_name=collection_name, schema=schema, index_params=index_params) client.insert(collection_name, [ {"content": "你们这个问题我已经打了三次电话了,到现在还没解决,太让人失望了!", "dummy_vector": [0.0, 0.0]}, {"content": "请问 Serverless Milvus 实例怎么开启公网访问?", "dummy_vector": [0.0, 0.0]}, ]) client.flush(collection_name) print("会话质检与工单分类结果:") for row in client.query(collection_name, filter="", output_fields=["content", "sentiment", "category"], limit=10): print(f" · 情感={row['sentiment']:<10} 类目={row['category']:<6} | {row['content'][:28]}") # ==================== 主流程:逐步执行,单步失败不影响其余步骤 ==================== def main() -> None: print("=" * 64) print("阿里云 Milvus AI Function 客服语音链路一体化演示") print("转写 → 脱敏 → 建库 → 召回 → 重排 → 质检") print("=" * 64) steps = [ ("步骤 1 录音转写", step_1_transcribe), ("步骤 2 PII 脱敏", step_2_pii_mask), ("步骤 3 建 FAQ 知识库", step_3_build_faq_kb), ("步骤 4 语义检索召回", step_4_semantic_recall), ("步骤 5 大模型重排", step_5_rerank), ("步骤 6 情感/分类质检", step_6_qc_and_route), ] passed, failed = [], [] for name, fn in steps: try: fn() passed.append(name) except Exception as exc: # noqa: BLE001 - 逐步隔离,保证所有功能都被跑到 failed.append(name) print(f"\n[!] {name} 执行失败:{type(exc).__name__}: {exc}") print("\n" + "=" * 64) print("执行汇总") print("=" * 64) print(f"成功 {len(passed)}/{len(steps)}:{', '.join(passed) if passed else '无'}") if failed: print(f"失败 {len(failed)}/{len(steps)}:{', '.join(failed)}") if __name__ == "__main__": main()
实战效果
依次跑通 转写 → 脱敏 → 建库 → 召回 → 重排 → 质检 六个步骤。下图是一次完整运行的真实控制台输出:
与传统多系统方案对比,简化非常直观:
维度 |
传统方案(ASR + 向量库 + NLP + Rerank 服务) |
阿里云 Milvus 一站式 |
系统数量 |
4+ 套,数据多处搬运 |
1 套,数据全程不出库 |
录音转写 |
应用侧先调 ASR 再写库 |
写入即转写(Collection Function) |
向量化链路 |
应用侧调 embedding 再写入 |
写入即向量化 |
PII 脱敏 |
放在链路末端,易遗漏 |
入库即脱敏,合规前置 |
召回 + 重排 |
额外部署 rerank 模型服务 |
|
情感 + 分类 |
外接 NLP 平台 |
|
总结
把在线/电话客服的语音链路收敛到阿里云 Milvus AI Function 之后,最直接的收益是工程形态大幅简化:AI_AUDIO_TRANSCRIBE、AI_PII_MASK、AI_EMBEDDING、AI_RERANK、AI_SENTIMENT、AI_CLASSIFY 六个函数在一套系统里串起了「转写 → 脱敏 → 建库 → 召回 → 重排 → 质检」的全流程,数据全程不出库,凭据不落地,脱敏前置,团队不必再维护「系统之间的系统」。
可延伸的场景包括:
- 坐席实时辅助:把召回 + 重排接入坐席工作台,通话进行中就推送最贴切的话术与知识。
- 批量质检:历史录音可结合
AI_BATCH做离线全量转写、脱敏与情感/分类打标,把质检覆盖率从抽检提升到 100%。 - 舆情预警:对
AI_SENTIMENT输出为negative的工单实时触发预警与升级流程。
最后再强调几条合规注意事项:录音处理前务必完成录音告知与客户授权;媒体一律使用短时效、最小权限的签名 URL,不要在 URL 中嵌入长期凭据;为原始音频、转写文本、情感标签、复核记录设置明确的保留期限,到期删除或匿名化,并保证派生数据随源数据的删除/授权撤回而同步处置;情感与分类结果是模型判断而非事实认定,在下架、升级等高影响动作前应保留人工复核环节。
让 AI 落地,不必再从复杂链路开始
阿里云 Milvus 是目前唯一支持 AI Native 的 Milvus 云服务。它不仅提供高性能、全托管的向量检索能力,更将 AI Function 与 AI-Gateway 原生嵌入数据链路,让模型能力真正成为数据库的一部分。
从数据进入阿里云Milvus 的那一刻起,理解、加工、检索和生成便可以在一条原生链路中完成。阿里云 Milvus 正在把“接入 AI”变成“原生拥有 AI”。