跳到主要内容

数据截至 (上游 commit 5e1f1fb87d9a)

第 06 章 · 另外两个子系统:流程框架与数据检索层

本章讲什么: SK 里除了 agent(见 第 04 章第 05 章),还有两块独立成体系、常被忽略的东西:一个把业务流程画成事件驱动图的 Process Framework,一个把向量库统一成一套接口的数据检索层。两块都能脱离 agent 单独用,原理也完全不同,所以合成一章讲。


6.0 为什么把这两块放在一起

它们的共同点只有一个:都不属于"模型对话"这条主线,却都是 SK 作为"企业级 SDK"而非"聊天玩具"的支撑件。

除此之外,两者毫无关系,读的时候当成两篇独立文章:

子系统它替你管什么入口符号能不能不用 Kernel
Process Framework多步骤业务流程的编排与恢复ProcessBuilder不能,步骤靠 Kernel 插件机制注册
数据检索层向量库的 schema、读写、过滤、检索vectorstoremodel / VectorStore能,VectorStore 完全独立

两块唯一的交汇点在最后:数据层的检索结果可以用 create_search_function 包成一个 KernelFunction 挂回 Kernel(见 6.14),从而被自动函数调用循环(第 03 章)选中。

全部 Process / 数据代码都带 @experimental@release_candidate 标记(python/semantic_kernel/utils/feature_stage_decorator.py),接口随时可能变。


Part A · Process Framework:事件驱动的业务流程图

6.1 一句话:它是什么

把一段多步骤业务,写成"一堆只会收参数、发事件的步骤" + "一张事件该送给谁的边表",然后交给运行时批量跑。

它要解决的问题不是"让模型更聪明",而是:

  • 一个业务流有 5 个环节,环节之间要传数据、要有分支、要能循环;
  • 某个环节调 LLM、某个环节调数据库,你不想手写 await a(); await b(); if ... 那一坨;
  • 流程跑到一半挂了,重启后想从存下来的状态接着跑;
  • 同一张流程图,本地跑用协程,线上跑要换成分布式 actor。

一句话直觉: 像 Kubernetes 的声明式编排,但编排对象是你的业务函数——你只声明"谁的哪个事件流向谁的哪个参数",不写调度代码。

6.2 用起来什么样(最小真实例子)

下面这段是官方样例 python/samples/getting_started_with_processes/step01/step01_processes.py:166-199 的骨架,已删去无关行:

process = ProcessBuilder(name="ChatBot")

# 1. 注册步骤:传"类型"而不是实例,运行时才实例化
intro_step = process.add_step(IntroStep)
user_input_step = process.add_step(ScriptedInputStep)
response_step = process.add_step(ChatBotResponseStep)

# 2. 外部事件进入流程的入口
process.on_input_event(event_id=ChatBotEvents.StartProcess).send_event_to(target=intro_step)

# 3. 步骤之间连线:某个函数跑完 → 触发下一个步骤
intro_step.on_function_result(function_name="print_intro_message").send_event_to(target=user_input_step)

# 4. 带参数名的连线:事件里的数据填给目标函数的 user_message 参数
user_input_step.on_event(event_id=CommonEvents.UserInputReceived).send_event_to(
target=response_step, parameter_name="user_message"
)

kernel_process = process.build() # 编译成纯数据
await start(process=kernel_process, kernel=kernel, initial_event=...) # 跑

注意第 3、4 行的形态:没有一处写"接下来调用谁",全是"谁的什么事件,流到谁的什么参数"。

6.3 三层结构:构建期 → 编译产物 → 运行期

这是理解整个框架的骨架,先看图(从左到右是时间顺序):

构建期(你写的链式代码) 编译产物(纯数据,可 JSON) 运行期(真正跑)
┌────────────────────────┐ ┌────────────────────────┐ ┌────────────────────────┐
│ ProcessBuilder │ │ KernelProcess │ │ LocalProcess │
│ .add_step(类型) │build│ state: 名字/版本/id │ 装载│ 超步循环(默认 100轮) │
│ .on_input_event(...) │────>│ steps: [StepInfo] │────>│ LocalStep × N │
│ .send_event_to(...) │ │ edges: {事件id:[边]} │ │ 每步一个事件队列 │
└────────────────────────┘ └────────────────────────┘ └────────────────────────┘
攒一张边表 一张边表 + 一堆步骤状态 也可整体换成 Dapr actor

三层各由谁负责:

主类文件
构建期ProcessBuilder / ProcessStepBuilderpython/semantic_kernel/processes/process_builder.pyprocess_step_builder.py
编译产物KernelProcess / KernelProcessStepInfo / KernelProcessEdgeprocesses/kernel_process/
运行期(本地)LocalProcess / LocalStepprocesses/local_runtime/
运行期(分布式)ProcessActor / StepActorprocesses/dapr_runtime/actors/

关键设计:中间那层是纯数据。 ProcessBuilder.build()(process_builder.py:164-175)把 builder 全部塌缩成 KernelProcess(state, steps, edges, factories),里面没有任何可执行对象——步骤只以"类型 + 状态"的形式存在。正因为如此,同一份编译产物既能喂给本地协程运行时,也能喂给 Dapr actor 运行时。

6.4 构建期:链式 DSL 如何变成一张边表

6.4.1 四个入口 API

链式调用的起点只有四个,区别在于事件 ID 怎么算:

API定义位置事件 ID 形态用途
ProcessBuilder.on_input_event(id)process_builder.py:132-140原样(不加前缀)流程外部事件的入口
ProcessStepBuilder.on_event(id)process_step_builder.py:84-91{步骤名}_{步骤uuid}.{id}接某步骤主动 emit 的事件
ProcessStepBuilder.on_function_result(fn)process_step_builder.py:224-239on_event(f"{fn}.OnResult")接某函数的返回值
ProcessStepEdgeBuilder.stop_process()process_step_edge_builder.py:57-65固定登记为 "END"终止流程(有坑,见 6.16)

on_event 的加前缀动作在 get_scoped_event_id(process_step_builder.py:141-143),前缀 event_namespace 是构造时算好的 f"{name}_{id}"(:59)。为什么要加命名空间? 因为同一个步骤类可以 add_step 多次,若不加前缀,两个实例发的 OnResult 会串线。

6.4.2 send_event_to 做了什么

链条的第二环 ProcessStepEdgeBuilder.send_event_to(process_step_edge_builder.py:28-55)只做三件事:

  1. 若传进来的是步骤(不是目标),包成 ProcessFunctionTargetBuilder;
  2. self.source.link_to(event_id, self) —— 把这条边挂进源步骤edges 字典(process_step_builder.py:255-264);
  3. 返回一个新的 edge builder,所以同一个事件可以继续 .send_event_to(...) 连多个目标(扇出)。

6.4.3 目标解析:函数名和参数名可以不写

ProcessFunctionTargetBuilder.__init__(process_function_target_builder.py:21-39)会调 resolve_function_target 帮你猜:

  • 目标步骤只有一个 kernel function → 函数名可省;多于一个则抛 KernelException(process_step_builder.py:104-108);
  • 该函数除 KernelProcessStepContext 外只剩一个参数 → 参数名可省(:120-135)。

这就是为什么 6.2 的例子里,连 intro_step 时什么都不用写,连 response_step 时才要写 parameter_name="user_message"

6.4.4 编译:build()

ProcessBuilder.build()(process_builder.py:164-175)把每条 edge builder .build()KernelProcessEdge(process_step_edge_builder.py:67-73),每个 step builder .build_step()KernelProcessStepInfo(process_step_builder.py:156-222,顺带把 @kernel_process_step_metadata 的版本号写进状态)。

结果就是一张字典:事件 ID → 边列表,边里记着"目标步骤 id / 目标函数名 / 目标参数名"。运行期只查这张表。

6.5 步骤契约:一个步骤要长什么样

步骤基类薄得惊人——只有一个可选的 activate 钩子(processes/kernel_process/kernel_process_step.py:16-23):

class KernelProcessStep(ABC, KernelBaseModel, Generic[TState]):
state: TState | None = None
async def activate(self, state: "KernelProcessStepState[TState]"): ...

真正的契约是约定,共三条:

  1. 业务方法用 @kernel_function 标注——步骤实例会被 kernel.add_plugin 注册成一个插件(local_step.py:203-209),所以步骤的函数就是普通 kernel function(第 01 章)。
  2. 要发事件就声明一个 KernelProcessStepContext 参数,调 context.emit_event(...)(kernel_process_step_context.py:22-42)。它接受 KernelProcessEvent、字符串或 Enum,内部统一包成事件对象往消息通道扔。
  3. 要有状态就写成 KernelProcessStep[MyState],activate 里接住框架给的状态对象。

事件有可见性之分(kernel_process_event.py:13-21):Internal 只在本流程内流转,Public 会被冒泡出流程边界(嵌套子流程时有用)。

6.6 运行期(核心):超步批处理循环

6.6.1 直觉先行

运行时不是"A 跑完立刻跑 B"。它是一轮一轮的:

一轮里,先把上一轮攒下的所有事件一次性收齐、算出该投递的消息、然后并发执行;这一轮里新产生的事件一律不在本轮生效,要等下一轮。

这就是"超步"(superstep)——批量同步并行的经典做法。好处是并发天然无竞态,坏处是"同一轮内 A 写了 B 立刻读"这种期待会落空。

6.6.2 一轮的四个动作

LocalProcess.internal_execute(local_runtime/local_process.py:191-222)整个循环只有 30 行:

┌─────────── 一个超步(最多 max_supersteps=100 轮)───────────┐
外部│ ① enqueue_external_messages: 抽干外部事件队列 → 查边表 → 造消息 │
事件│ ② 逐个步骤 enqueue_step_messages: 抽干它的事件队列 → 造消息 │
──> │ ③ 把 message_channel 一次性取空;取不到 且 没有外部事件 → 结束 │
│ ④ asyncio.gather:所有消息并发投递给目标步骤 handle_message │
└───────────────────────────┬───────────────────────────────────┘
│ 步骤里 emit_event 只是 put 进自己队列
└──> 要到下一轮的 ② 才被看见 ← 这就是"栅栏"

对应源码位置:

动作符号位置
收外部事件enqueue_external_messageslocal_process.py:237-245
收各步骤事件enqueue_step_messageslocal_process.py:247-260
批量取消息 / 判停internal_execute 循环体local_process.py:203-209
并发执行asyncio.gather(*message_tasks)local_process.py:218

max_supersteps 默认 100(local_process.py:47-49),start() 可以覆盖(local_kernel_process.py:17-52)。它是防死循环的闸门,不是"步骤最多跑 100 个"。

6.6.3 事件是怎么变成函数参数的

这条链是整个框架最值得看懂的部分:

步骤 A 的函数 run() 正常返回
│ LocalStep 自动 emit 一个 "run.OnResult" 事件 (local_step.py:163-170)
v
LocalEvent(id = "A_<uuid>.run.OnResult") (local_event.py:20-23)
│ get_edge_for_event(id) 查 A 的 output_edges (local_step.py:270-275)
v
KernelProcessEdge{ source: A, target: {step:B, fn:"handle", param:"text"} }
│ LocalMessageFactory.create_from_edge(edge, data) (local_message_factory.py:15-30)
v
LocalMessage{ destination:B, function_name:"handle", values:{"text": data} }
│ B.handle_message: 把 text 塞进 inputs["handle"] (local_step.py:104-115)
v
参数攒齐 → kernel.invoke(B.handle, text=...) → 再 emit "handle.OnResult"

注意 scoped_event(local_step.py:292-297):事件被放进队列前,namespace 被强制改写成发射者自己的 {name}_{id}。所以事件 ID 一定和构建期 get_scoped_event_id 算出来的前缀对得上——这是两端能对接的原因。

6.6.4 "参数攒齐才触发"

LocalStep.handle_message(local_step.py:81-173)是步骤侧的核心。它不是"来一条消息跑一次函数",而是:

  1. 把消息里的键值填进 self.inputs[函数名],重复填会打日志覆盖(:104-115);
  2. 扫出所有"参数全非 None"的函数,叫 invocable_functions(:117-121);
  3. 一个都没有 → 记日志、直接返回,等下一条消息(:129-131);
  4. 可调用的函数如果和消息指定的函数对不上 → 抛 ProcessTargetFunctionNameMismatchException(:135-139);
  5. 调用,成功发 {fn}.OnResult、失败发 {fn}.OnError,finally 里必发一个(:156-170);
  6. 把该函数的输入重置回初始值(:173),下一轮重新攒。

这套 join 语义意味着:一个需要两个参数的函数,可以由两个不同步骤各喂一个,谁后到谁触发。

初始输入表由 find_input_channels(step_utils.py:19-40)算出:Kernel 类型的参数跳过(Kernel 后面自动注入),非必填参数跳过,KernelProcessStepContext 参数直接预填好实例——所以 context 参数不会卡住 invocable 判定。

6.6.5 步骤实例化与工厂

步骤真正被 new 出来是在第一条消息到达时,initialize_step(local_step.py:188-261)里:

  • 有工厂函数就调工厂(支持 async),否则 step_cls()(:191-201);
  • 注册成插件、抽出所有 kernel function(:203-209);
  • 用泛型参数反推状态类型 KernelProcessStepState[TState],没有则用裸的基类(:219-253);
  • 最后 await step_instance.activate(state_object)(:261)。

工厂机制(ProcessBuilder.add_step(..., factory_function=...),process_builder.py:73-75)存在的理由写在 docstring 里:有些步骤依赖不能 JSON 序列化(比如数据库连接),不能靠"存类型再反序列化"复原。

6.7 状态快照与版本化恢复

6.7.1 状态存哪

运行中随时可以要一份快照:LocalKernelProcessContext.get_state()LocalProcess.to_kernel_process()(local_process.py:224-231),它把每个 LocalStep 反向压回 KernelProcessStepInfo,得到一个和 build() 产物同构的 KernelProcess

注意:这是主动拉取,不是每步自动落盘。 本地运行时没有内置 checkpoint。

6.7.2 序列化格式

落盘用的是另一组类型(kernel_process_step_state_metadata.py:18-41),字段名带 $type / versionInfo / stepsState 这类别名,明显是为了和 .NET 版对齐:

class KernelProcessStepStateMetadata(KernelBaseModel, Generic[TState]):
type_: Literal["Step", "Process"] = Field("Step", alias="$type")
version_info: str | None = Field(None, alias="versionInfo")
state: TState | None = Field(None, alias="state")

转换函数在 process_state_metadata_utils.py:kernel_process_to_process_state_metadata(:25-39)递归展开子流程,step_info_to_process_state_metadata(:54-67)处理单步。读回来用 KernelProcessStateMetadata.load_from_file(kernel_process_step_state_metadata.py:43-74)。

6.7.3 版本号与迁移

步骤类可以打版本标记(kernel_process_step_metadata.py:17-36):

@kernel_process_step_metadata("CutFoodStep.V2")
class CutFoodWithSharpeningStep(KernelProcessStep[MyState]): ...

恢复时,_sanitize_process_state_metadata(process_builder.py:201-254)做两件事:

情况处理
步骤改了名,但在 aliases 里能匹配上,且版本相同把旧 key 的状态挪到新名字下(:226-229)
匹配上但版本不同,且该节点是子流程递归下钻,逐层 sanitize(:232-243)
匹配上但版本不同,普通步骤放弃旧状态,只迁移 id(:245-249)

也就是说:改名可以无损恢复,改版本号则默认丢状态。这个取舍是显式的——版本号的语义就是"行为/状态 schema 不兼容"(见 kernel_process_step_metadata.py:22-25 的 docstring)。

6.8 分布式替身:Dapr runtime

processes/dapr_runtime/ 是本地运行时的同构替换:相同的 KernelProcess 编译产物,换一套运行时跑。

本地Dapr差别
LocalProcessProcessActor(actors/process_actor.py:47)进程 = 一个 actor
LocalStepStepActor(actors/step_actor.py)每个步骤 = 一个 actor
Queue 内存队列MessageBufferActor / EventBufferActor / ExternalEventBufferActor队列本身也是 actor,状态持久
快照靠主动拉取每次收发都写 actor stateActorStateKeys(actors/actor_state_key.py:10-30)

超步循环的形状几乎一样(process_actor.py:334-364),但多了两处分布式必需的动作:

  • 每轮开头先问"END 消息发出去了吗"(_is_end_message_sent,:408),因为 actor 之间没有共享内存;
  • 收消息拆成两阶段:prepare_incoming_messages() 返回条数用于判停,再统一 process_incoming_messages()(:347-357)。本地版本因为共享内存,一次抽干队列就够了。

持久化在 StepActor 里到处可见:收到消息先 try_add_state(StepIncomingMessagesState, ...)save_state()(step_actor.py:165-169),函数跑完把状态 JSON 落盘(:371-372)。这是 Dapr 版才有的可恢复性,本地版没有。

启动入口也是平行的:本地 local_runtime/local_kernel_process.py:17start,Dapr dapr_runtime/dapr_kernel_process.py:16start,签名几乎相同,只是后者多一个 process_id

6.9 和 LangGraph 超步模型的异同

和 index 的分工: index §7 比的是四个框架各自把什么当核心抽象;本节只比一件事——SK Process 与 LangGraph Pregel 的超步触发机制差在哪。两处不重复,可以对照着读。

两者都叫"超步",但触发机制完全不同。先看相同点:

  • 都是"一轮批量执行 + 栅栏"的 BSP 形状;
  • 都有超步数上限防死循环(SK 的 max_supersteps=100,LangGraph 的 recursion_limit);
  • 一轮内产生的写入/事件都要下一轮才可见。

不同点是本质性的:

维度SK Process FrameworkLangGraph Pregel
谁决定下一步跑谁事件 ID 命中静态边表(get_edge_for_event)channel 版本号 > 节点已见版本 → 唤醒
数据怎么传消息直接填目标函数的具名参数节点读写共享 state channel
并发写冲突不存在共享状态,谁写谁的参数槽reducer 合并同一 channel 的多份写入
判停条件本轮没有任何消息可投递没有任务被唤醒
持久化本地版靠主动拉快照;Dapr 版每步落 actor state每个超步末尾写 checkpoint,内建
恢复语义靠步骤名/别名 + 版本号对齐(见 6.7.3)靠 channel 版本号与 versions_seen 对齐

一句话概括差异: LangGraph 是"共享状态 + 版本号驱动",SK Process 是"无共享状态 + 事件路由驱动"。前者更像数据流图,后者更像消息总线。细节见 ../langgraph/01-pregel-bsp.md


Part B · 数据检索层:向量库的统一外壳

6.10 一句话 + 最小例子

把"一个 Python 类"当成向量库的表结构声明,剩下的建表、写入、生成 embedding、过滤、检索,全部由统一接口代劳,底下换哪家数据库不改业务代码。

样例 python/samples/concepts/memory/simple_memory.py:30-39:

@vectorstoremodel(collection_name="test")
@dataclass
class DataModel:
content: Annotated[str, VectorStoreField("data")]
id: Annotated[str, VectorStoreField("key")] = field(default_factory=lambda: str(uuid4()))
vector: Annotated[list[float] | str | None, VectorStoreField("vector", dimensions=1536)] = None
title: Annotated[str, VectorStoreField("data", is_full_text_indexed=True)] = "title"
tag: Annotated[str, VectorStoreField("data", is_indexed=True)] = "tag"

重点看:类型注解里塞了 VectorStoreField 一个普通 dataclass,加一个装饰器,就同时是"业务对象"和"库表 schema"。

三条主管线长这样:

写入: 你的对象 ──成 dict──> 补 embedding ──> 存储层格式 ──> _inner_upsert
检索: query ──(可选)本地算向量──> _inner_search ──> 异步结果流 ──> 反序列化 ──> 你的对象
过滤: lambda x: x.tag == "a" ──getsource + ast.parse──> AST ──_lambda_parser──> 厂商查询语法

6.11 @vectorstoremodel:把类解析成 schema

6.11.1 三种字段

VectorStoreField(data/vector.py:276-417)是一个带三个 @overload 的 dataclass,靠第一个位置参数分流:

field_type必填可选项语义
"key"storage_name / type主键,有且只能有一个
"data"is_indexed / is_full_text_indexed普通字段,可建普通索引或全文索引
"vector"dimensionsindex_kind / distance_function / embedding_generator向量字段

向量字段的两个枚举值得单独记:

  • IndexKind(vector.py:170-216):HNSW / FLAT / IVF_FLAT / DISK_ANN / QUANTIZED_FLAT / DYNAMIC / DEFAULT,docstring 里逐个解释了精度与代价的取舍;
  • DistanceFunction(vector.py:219-262):余弦相似度、余弦距离、点积、欧氏、欧氏平方、曼哈顿、汉明。

一个容易忽略的巧思: DISTANCE_FUNCTION_DIRECTION_HELPER(vector.py:265-273)是一张 距离函数 → 比较运算符 的表——相似度类用 operator.gt(越大越好),距离类用 operator.le(越小越好)。这样上层代码判断"哪个结果更好"时不必写 if-else。

6.11.2 解析流程

装饰器本身只有十几行(vector.py:639-684):

def wrap(cls):
cls_sig = signature(cls) # 拿 __init__ 的参数表
setattr(cls, "__kernel_vectorstoremodel__", True)
setattr(cls, "__kernel_vectorstoremodel_definition__",
_parse_signature_to_definition(cls_sig.parameters, collection_name))
return cls

注意它用的是 inspect.signature(cls) 而不是 __annotations__ 这就是为什么它对 dataclass、pydantic model、普通类都通用——只要构造函数参数带注解就行。

逐字段解析在 _parse_parameter_to_field(vector.py:602-616):

  1. 参数注解有 __metadata__(即 Annotated[...]),且里面有 VectorStoreField 实例 → 收下;
  2. 否则,该参数必须有默认值,不然直接抛 VectorStoreModelException(:611-614)。理由写在注释里:这个字段不会被存,取回来时构造对象会失败;
  3. 有默认值的无注解字段则被静默忽略(打 debug 日志)。

字段的 type__parse_vector_store_record_field_instance(vector.py:569-599)从注解里扒:对向量字段,会剥掉 Optional 和外层容器,取出内层元素类型——这就是为什么样例里写 list[float] | str | None 也能被正确识别。

最后 VectorStoreCollectionDefinition.model_post_init(vector.py:543-563)做三项校验:至少一个字段、字段名非空、key 字段有且仅有一个

6.11.3 装饰器顺序的坑

_parse_signature_to_definition(vector.py:619-626)的报错信息本身就是文档:

参数为空时提示:如果你和 @dataclass 一起用,可能把装饰器顺序写反了,vectorstoremodel 必须在最上面

原因很直白:@dataclass 要先生成 __init__,signature(cls) 才有东西可读。

6.12 把 Python lambda 翻译成厂商查询语法

6.12.1 问题与思路

问题: 你想写 filter=lambda x: x.tag == "general" and x.year > 2020,但 Azure AI Search 要的是 OData 字符串、Postgres 要的是 SQL、Qdrant 要的是 Filter 对象树、Redis 要的是 redisvl 的 FilterExpression

思路: 不去执行这个 lambda,而是读它的源码、解析成 AST、再由每家 connector 各自翻译

6.12.2 三步实现

第一步,_build_filter(vector.py:1964-1999)拿到源码并解析:

visitor = LambdaVisitor(self._lambda_parser)
for filter_ in filters:
tree = parse(filter_ if isinstance(filter_, str) else getsource(filter_).strip())
visitor.visit(tree)

getsource 是关键:直接把 lambda 的源文本抠出来。所以过滤条件既可以传真的 lambda,也可以传一个字符串 "lambda x: x.tag == 'a'" ——两条路走同一套解析。

第二步,LambdaVisitor(vector.py:717-727)只重写了一个方法:

def visit_Lambda(self, node: Lambda) -> None:
self.output_filters.append(self.lambda_parser(node.body))

它只取 lambda 的函数体(丢掉参数名),交给厂商解析器。

第三步,_lambda_parser 是抽象方法(vector.py:2001-2009),每家 connector 用 match node: 模式匹配自己实现。同一个 AST 节点,四家的产出:

Connector_lambda_parser 位置x.tag == "a" 翻成产物类型
Azure AI Searchconnectors/azure_ai_search.py:636tag eq 'a'OData 字符串
Postgresconnectors/postgres.py:852"tag" = 'a'SQL 片段字符串
Qdrantconnectors/qdrant.py:367FieldCondition(key="tag", match=MatchValue(value="a"))qdrant 模型对象
Redisconnectors/redis.py:366Tag("tag") == "a"redisvl FilterExpression
In-Memoryconnectors/in_memory.py:864—(空实现)直接跑原 lambda

6.12.3 几处共同的讲究

  • 链式比较自动展开。 1 < x.year < 3 这种 Python 特有写法,各家都把它拆成多个二元比较再 AND 起来(如 postgres.py:856-864qdrant.py:374-380)。
  • 字段名对着 schema 校验。 Postgres 在 ast.Attribute 分支直接查 self.definition.storage_names,不在里面就抛 VectorStoreOperationException(postgres.py:900-905)。这既防笔误,也顺带防注入。
  • 字符串常量做转义。 Postgres 把 ' 变成 ''(postgres.py:918)。
  • Redis 要看字段元数据才知道用哪种表达式。 get_field_expr(redis.py:369-385)按字段是全文索引 / 数值 / 其它,分别产出 Text / Num / Tag;向量字段上过滤直接报错。
  • 内存实现选择不翻译。 InMemoryCollection._run_filter(in_memory.py:856-861)把记录包成支持属性访问的 dict,直接调用原 lambda(:739 处使用)。既然数据在内存里,何必翻译。

代价要说清楚: getsource 依赖源码文件可读。在 REPL、exec 出来的代码里定义的 lambda 拿不到源码——这时只能改传字符串形式。

6.13 统一接口:VectorStore 与 VectorStoreCollection

6.13.1 两层职责

VectorStore —— 连接层:一个数据库实例
├─ get_collection(record_type) —— 抽象方法,返回下面这个
├─ list_collection_names()
├─ collection_exists(name) —— 用一个只含 key 的临时 definition 探测
└─ ensure_collection_deleted(name)

v
VectorStoreCollection —— 操作层:一张表 + 一个数据模型
├─ 公共 API: upsert / get / delete / ensure_collection_exists ...
└─ 抽象钩子: _inner_upsert / _inner_get / _inner_delete

v
VectorSearch (mixin) —— 检索层
├─ 公共 API: search / hybrid_search / create_search_function
└─ 抽象钩子: _inner_search / _get_record_from_result / _get_score_from_result / _lambda_parser

VectorStorevector.py:1540-1614,VectorStoreCollection:1133,VectorSearch:1621

6.13.2 公共方法 / 抽象钩子的分工

这是整套设计的核心套路:公共方法负责序列化和异常语义,抽象钩子只管跟数据库说话。

公共方法位置它包办了什么落到哪个钩子
upsertvector.py:1282-1323单条/批量归一、序列化、异常包成 VectorStoreOperationException_inner_upsert(:1174)
getvector.py:1327 起(四个 overload)按 key 取 / 按条件取、反序列化_inner_get(:1200)
deletevector.py:1517批量归一_inner_delete(:1227)
searchvector.py:1801-1862先查 supported_search_types、组装 options、异常归类_inner_search
hybrid_searchvector.py:1866同上_inner_search

建表相关的三个方法 ensure_collection_exists(:1245)、collection_exists(:1259)、ensure_collection_deleted(:1271)是纯抽象的——各家建索引的参数差太远,没法统一。

命名值得注意: 用的是 ensure_collection_exists 而不是 create_collection,语义上就是幂等的。

6.13.3 写入管线:embedding 在哪一步补上

serialize(vector.py:853-907)是四段式,顺序有讲究:

  1. 用户自定义的 serialize 函数(definition 里可以塞)——用户优先;
  2. 对象 → dict;
  3. _add_vectors_to_records(vector.py:1038-1082)——在这里生成 embedding;
  4. dict → 该存储的模型格式。

第 3 步的逻辑:遍历所有向量字段,字段级 embedding_generator 优先于集合级,两个都没有就跳过(:1062-1064)。跳过意味着"由数据库自己算向量"——VectorStoreField 的 docstring 明说了这个分工(vector.py:355-357)。

6.13.4 连接器清单

VectorStore 的实现散在 connectors/ 下,一个文件一家:

文件Store 类支持的检索类型(supported_search_types)
connectors/azure_ai_search.py:752AzureAISearchStore向量 + 关键词混合(:308)
connectors/qdrant.py:550QdrantStore向量 + 关键词混合(:134)
connectors/mongodb.py:534MongoDBAtlasStore向量 + 关键词混合(:166)
connectors/postgres.py:941PostgresStore仅向量(:317)
connectors/redis.py:758RedisStore仅向量(:196)
connectors/in_memory.py:888InMemoryStore仅向量(:421)
connectors/faiss.py:250FaissStore继承 InMemoryStore,只换索引实现(faiss.py:67)

其余还有 weaviate.py:684sql_server.py:690oracle.py:1193pinecone.py:605chroma.py:437azure_cosmos_db.py

另外提供了两个 Protocol——VectorStoreCollectionProtocol(vector.py:2138)和 VectorSearchProtocol(:2261),都带 @runtime_checkable,给不想继承基类的第三方实现留了鸭子类型的口子。

6.14 把检索包成 Kernel 插件函数

6.14.1 一行代码接上 agent

create_search_function(vector.py:2011-2076)把一个 collection 变成一个 KernelFunction:

# 示意,非源码
search_fn = collection.create_search_function(
function_name="search_docs",
description="查内部文档",
string_mapper=lambda r: r.record.content, # 只把正文喂给模型
top=5,
)
kernel.add_function(plugin_name="memory", function=search_fn)

挂上去之后,它就和普通插件函数一样,能被自动函数调用循环选中(第 03 章)。

6.14.2 内部怎么包的

_create_kernel_function(vector.py:2078-2134)里定义了一个 search_wrapper,再用 KernelFunctionFromMethod 包成 kernel function。默认暴露给模型的参数是三个(data/_shared.py:29-53):query(必填)、top(默认 2)、skip(默认 0)。

结果映射有两档:给了 string_mapper 就用它,否则 result.model_dump_json(exclude_none=True)(vector.py:2126-2128)。

6.14.3 巧妙的一环:额外参数自动变成过滤条件

default_dynamic_filter_function(_shared.py:146-183)干了件很妙的事——把模型传来的非检索参数拼成 lambda 源码字符串:

new_filter = f"lambda x: x.{param.name} == {_format_filter_literal(kwargs[param.name])}"

于是整条链闭合了:你给 search function 声明一个 tag 参数 → 模型填 tag="general" → 自动生成 "lambda x: x.tag == 'general'" → 进 6.12 的 AST 翻译器 → 变成该数据库的过滤语法。 字面量用 repr() 格式化(_shared.py:186-188),规避引号问题。

6.14.4 TextSearch:非向量检索走同一套

data/text_search.py:55TextSearch 基类是同一个模式的另一份实现,给的是网页搜索这类没有向量的检索源:BraveSearch(connectors/brave.py:113)、GoogleSearch(connectors/google_search.py:150)。

它和向量版的差别只有输出形态(text_search.py:194-253),三选一:

output_type返回适用
"str"纯字符串直接喂给模型
"TextSearchResult"name / value / link 三元组(:43-48)要带来源链接
"Any"原始结果对象string_mapper 自己抽字段

抽象方法只有一个 search(text_search.py:318-334)——接一个新搜索源,实现它就够了。


6.15 巧妙之处(可以搬走的技术)

  • 编译产物是纯数据,运行时可换。 一份 KernelProcess 同时能喂本地协程和 Dapr actor(process_builder.py:164-175local_process.py / process_actor.py)。想支持新的运行时,不用碰 builder。
  • "参数攒齐才触发"是天然的 join。 LocalStep 用一张 inputs[函数][参数] 表实现多路汇合,执行完立刻重置(local_step.py:117-121:173),几十行就有了工作流引擎的 join 语义。
  • 事件加命名空间前缀解决实例串线。 {步骤名}_{uuid}.{事件名},构建期算一次(process_step_builder.py:59)、运行期强制改写一次(local_step.py:292-297),两端天然对齐。
  • 别名 + 版本号把"改代码"和"恢复旧状态"解耦。 改名走别名匹配,改行为改版本号,规则写死在 _sanitize_process_state_metadata(process_builder.py:201-254)。
  • signature(cls) 而非 __annotations__ 解析模型。 一套解析代码同时吃 dataclass / pydantic / 普通类(vector.py:668)。
  • 不执行 lambda,而是解析它的源码。 getsource + ast.parse + 每家一个 visitor(vector.py:1964-1999),把"一种写法 → N 种数据库方言"的翻译成本压到每家一个 match 语句。
  • 距离函数方向表。 DISTANCE_FUNCTION_DIRECTION_HELPER(vector.py:265-273)把"越大越好还是越小越好"变成查表,消灭一堆 if。
  • 函数参数自动降级为过滤条件。 default_dynamic_filter_function(_shared.py:146-183)生成 lambda 字符串再喂给 AST 翻译器,让"模型能填的参数"和"数据库能过滤的字段"自动对上。

6.16 边界与局限(诚实版)

Process Framework:

  • 整个模块打着 @experimental,接口不保证稳定。
  • 本地运行时没有持久化。 快照要自己调 get_state() 拉,崩了就没了;要自动落盘只能上 Dapr(step_actor.py:165-169)。
  • stop_process() 在 Python 本地运行时里似乎不生效。 它把边登记在 key "END" 下(process_step_edge_builder.py:65),而运行期是按真实事件 ID 查表(local_step.py:270-275);同时 EndStep 实例的 id 是随机 uuid(继承 ProcessStepBuilder.__init__,process_step_builder.py:57),对不上 local_process.py:212 判断用的 END_PROCESS_ID 常量(const.py:7)。流程实际上是靠"本轮没有任何消息"自然停下的,不是靠 END 分支 (inferred:两处 id 来源不一致是代码直接可见的,"因此 END 分支走不到"是我的推断,未实跑验证)
  • 超步内 break 会丢消息。 local_process.py:212-213 遇到目标是 END 的消息就 break 出消息循环,同一批剩下的消息不会被执行。
  • 没有条件边。 分支只能靠步骤内部 emit 不同事件来表达,边表本身不带条件判断。
  • max_supersteps 默认 100 是硬闸门,长流程要自己调大(local_kernel_process.py:30-31)。

数据检索层:

  • 过滤条件依赖能拿到 lambda 源码。 REPL / exec 里定义的 lambda 会失败,只能改传字符串。
  • 各家 _lambda_parser 支持度不齐。 Postgres 明确不支持一元 + - ~(postgres.py:898-899),Azure 不支持 Invert(azure_ai_search.py:711-712),写复杂条件前得看目标 connector 的实现。
  • in 的语义各家不同。 Azure 把 in 翻成 search.ismatch(...) 全文匹配(azure_ai_search.py:684),Postgres 翻成 SQL IN(postgres.py:869)——同一段 lambda 换库行为会变。
  • 本地生成 embedding 目前只处理字符串。 非字符串会被 json.dumps 后再向量化,源码里有 TODO 承认这点(vector.py:1957-1959)。
  • 混合检索不是都支持。 只有 Azure AI Search / Qdrant / MongoDB 声明了 KEYWORD_HYBRID,调错会抛 VectorStoreOperationNotSupportedException(vector.py:2053-2057)。

6.17 代码地图(导航索引)

路径相对克隆根,Python 侧统一在 python/semantic_kernel/ 下。

主题文件路径符号名
流程构建 DSL 主体python/semantic_kernel/processes/process_builder.pyProcessBuilder.add_step / on_input_event / build
步骤 builder / 事件命名空间python/semantic_kernel/processes/process_step_builder.pyProcessStepBuilder.on_event / on_function_result / resolve_function_target / link_to
外部事件边python/semantic_kernel/processes/process_edge_builder.pyProcessEdgeBuilder.send_event_to
步骤间边 / 终止python/semantic_kernel/processes/process_step_edge_builder.pyProcessStepEdgeBuilder.send_event_to / stop_process / build
目标(函数+参数)解析python/semantic_kernel/processes/process_function_target_builder.pyProcessFunctionTargetBuilder
步骤基类契约python/semantic_kernel/processes/kernel_process/kernel_process_step.pyKernelProcessStep.activate
步骤发事件python/semantic_kernel/processes/kernel_process/kernel_process_step_context.pyKernelProcessStepContext.emit_event
事件与可见性python/semantic_kernel/processes/kernel_process/kernel_process_event.pyKernelProcessEvent / KernelProcessEventVisibility
超步循环(本地)python/semantic_kernel/processes/local_runtime/local_process.pyLocalProcess.internal_execute / enqueue_step_messages
步骤执行与参数汇合python/semantic_kernel/processes/local_runtime/local_step.pyLocalStep.handle_message / initialize_step / get_edge_for_event / scoped_event
边 → 消息python/semantic_kernel/processes/local_runtime/local_message_factory.pyLocalMessageFactory.create_from_edge
本地启动入口python/semantic_kernel/processes/local_runtime/local_kernel_process.pystart
输入通道推断python/semantic_kernel/processes/step_utils.pyfind_input_channels / get_step_class_from_qualified_name
状态与版本python/semantic_kernel/processes/kernel_process/kernel_process_step_state.pyKernelProcessStepState
快照序列化格式python/semantic_kernel/processes/kernel_process/kernel_process_step_state_metadata.pyKernelProcessStateMetadata.load_from_file
版本标记装饰器python/semantic_kernel/processes/kernel_process/kernel_process_step_metadata.pykernel_process_step_metadata
状态转换工具python/semantic_kernel/processes/process_state_metadata_utils.pykernel_process_to_process_state_metadata / to_process_state_metadata
超步循环(Dapr)python/semantic_kernel/processes/dapr_runtime/actors/process_actor.pyProcessActor.internal_execute / _is_end_message_sent
步骤 actor 与持久化python/semantic_kernel/processes/dapr_runtime/actors/step_actor.pyStepActor.prepare_incoming_messages / process_incoming_messages
Dapr 状态键python/semantic_kernel/processes/dapr_runtime/actors/actor_state_key.pyActorStateKeys
模型装饰器python/semantic_kernel/data/vector.pyvectorstoremodel / _parse_signature_to_definition
字段与 schemapython/semantic_kernel/data/vector.pyVectorStoreField / VectorStoreCollectionDefinition / IndexKind / DistanceFunction
lambda → 查询语法python/semantic_kernel/data/vector.pyLambdaVisitor / VectorSearch._build_filter / _lambda_parser
读写与序列化管线python/semantic_kernel/data/vector.pyVectorStoreCollection.upsert / VectorStoreRecordHandler.serialize / _add_vectors_to_records
检索接口python/semantic_kernel/data/vector.pyVectorSearch.search / hybrid_search / create_search_function
存储层接口python/semantic_kernel/data/vector.pyVectorStore.get_collection / VectorStoreCollectionProtocol
检索转插件函数(通用)python/semantic_kernel/data/text_search.pyTextSearch.create_search_function
参数转过滤条件python/semantic_kernel/data/_shared.pydefault_dynamic_filter_function / DEFAULT_PARAMETER_METADATA
各厂商过滤翻译python/semantic_kernel/connectors/{azure_ai_search,postgres,redis,qdrant}.py_lambda_parser
内存/本地实现python/semantic_kernel/connectors/{in_memory,faiss}.pyInMemoryCollection._run_filter / FaissCollection
可运行样例python/samples/getting_started_with_processes/step01/step01_processes.pystep01_processes
可运行样例python/samples/concepts/memory/simple_memory.pyDataModel

回到目录: Semantic Kernel — 架构与原理 · 上一章: 多 agent 编排