ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

使用 Dapr 运行 Semantic Kernel Processes:FastAPI/Flask 集成实战指南

使用 Dapr 运行 Semantic Kernel Processes:FastAPI/Flask 集成实战指南 使用 Dapr 运行 Semantic Kernel ProcessesFastAPI/Flask 集成实战指南【免费下载链接】semantic-kernelIntegrate cutting-edge LLM technology quickly and easily into your apps项目地址: https://gitcode.com/GitHub_Trending/se/semantic-kernel本指南以python/samples/demos/process_with_dapr/演示为蓝本系统讲解如何在 FastAPI 与 Flask 应用中接入 Dapr 运行时来托管 Semantic Kernel Processes实现流程的弹性伸缩、状态持久化与可靠执行。读完本文你将掌握 Dapr Actor 的注册机制、进程的启动方式、状态恢复行为以及如何在本地一键运行这一 Demo。为什么用 Dapr 托管 Semantic Kernel ProcessesSemantic Kernel 的 Process Framework 允许开发者将多步骤业务编排定义为有状态、事件驱动的工作流其中每个步骤Step都是一段可复用逻辑步骤之间通过事件传递数据。而 DaprDistributed Application Runtime是一个可移植、事件驱动的运行时能够简化在云端和边缘环境中构建弹性、有状态应用的过程天然适合托管 Semantic Kernel Processes让你在不牺牲性能与可靠性的前提下从规模和数量上横向扩展你的流程实例。两者的结合点非常清晰进程实例有状态Dapr 的 Actor 模型为每个进程实例和步骤实例提供了持久化状态存储步骤间通信天然事件化Dapr 的 Actor 消息通道与 Semantic Kernel Process 的emit_event事件机制高度契合弹性伸缩Dapr 运行时负责 Actor 的激活、休眠与恢复多个进程实例可以并行运行而互不干扰。Dapr 扩展支持矩阵当前仓库的 Dapr 运行时python/semantic_kernel/processes/dapr_runtime/支持以下 Web 框架集成方式扩展支持状态FastAPI✅Flask✅gRPC❌Dapr Workflow❌从源码看dapr_actor_registration.py中分别提供了register_fastapi_dapr_actors异步注册与register_flask_dapr_actors同步注册两个入口而 gRPC 与 Dapr Workflow 两种扩展在当前版本中尚未实现。前置条件配置 Dapr 本地开发环境在运行 Demo 之前必须先完成 Dapr 的本地开发环境配置并且Dapr 容器必须处于运行状态否则演示应用无法正常工作。建议按 Dapr 官方文档完成自托管self-hosted安装通常包含dapr init初始化运行时以及启动相关的 sidecar 容器。开发调试时推荐两种方式Dapr CLI直接使用命令行工具启动服务Dapr VS Code 扩展若希望在代码运行过程中调试这是推荐方式——从 Run and Debug 下拉列表中选择Python FastAPI App with Dapr或Python Flask API App with Dapr即可。运行 Demo构建并启动服务按以下步骤运行示例构建并运行示例程序见上节的 Dapr CLI 或 VS Code 扩展两种方式服务启动后会在localhost 端口 5001暴露一个 API 端点。以 FastAPI 应用为例服务入口在 fastapi_app.py其__main__块通过uvicorn.run(app, host0.0.0.0, port5001, log_levelerror)监听 5001 端口Flask 版本则在 flask_app.py 中以app.run(host0.0.0.0, port5001)监听同一端口。调用流程实例打开浏览器访问http://localhost:5001/processes/1234即可发起一个Id 1234的新流程实例流程执行完成后浏览器中应显示{processId:1234}同时服务控制台会输出与下面格式一致的日志。第一次运行时的日志##### Kickoff ran. ##### AStep ran. ##### BStep ran. ##### CStep activated with Cycle 1. ##### CStep run cycle 2. ##### Kickoff ran. ##### AStep ran. ##### BStep ran. ##### CStep run cycle 3 - exiting.刷新页面再次运行同一实例后的日志##### Kickoff ran. ##### AStep ran. ##### BStep ran. ##### CStep run cycle 4 - exiting.理解状态持久化带来的日志差异两次运行的日志并不相同这正是 Dapr 状态持久化的直接体现第一次运行CState以Cycle 1初始化该初始值来自构建流程时指定的初始状态见 process.py 中的process.add_step(step_typeCStep, initial_stateCStepState(current_cycle1))在达到终止条件Cycle 3之前CState共被调用了两次。第二次运行CState以Cycle 3初始化该值正是第一次运行的最终状态——Dapr 已将流程状态持久化到 Actor 状态存储中并成功恢复CState只被调用一次因为它已处于Cycle 3的终止条件。新实例验证将浏览器指向http://localhost:5001/processes/ABCD创建一个Id ABCD的新实例你会看到它如预期从初始状态开始执行——因为这是全新的进程 IDDapr 中不存在对应的持久化状态。演示流程结构该 Demo 展示的是一个带扇出fan-out、扇入fan-in和循环的流程其拓扑如下Kickoff启动后同时向A与B发出事件扇出A与B都完成后汇聚到C扇入C需要同时收到astepdata与bstepdata两个参数C根据内部计数决定是回到Kickoff开启新一轮循环还是触发ExitRequested事件结束整个流程。理解代码FastAPI 集成详解项目文件结构python/samples/demos/process_with_dapr/ ├── README.md # 本文档 ├── fastapi_app.py # FastAPI 服务入口 ├── flask_app.py # Flask 服务入口 └── process/ ├── __init__.py ├── process.py # 流程定义ProcessBuilder └── steps.py # 步骤与事件定义FastAPI 应用导入与 Dapr 包import logging from contextlib import asynccontextmanager import uvicorn from dapr.ext.fastapi import DaprActor from fastapi import FastAPI from fastapi.responses import JSONResponseSemantic Kernel Process 相关导入from samples.demos.process_with_dapr.process.process import get_process from samples.demos.process_with_dapr.process.steps import CommonEvents from semantic_kernel import Kernel from semantic_kernel.processes.dapr_runtime import ( register_fastapi_dapr_actors, start, )关键点register_fastapi_dapr_actors与start均从semantic_kernel.processes.dapr_runtime包导出见 dapr_runtime/init.py其中还导出了ProcessActor、StepActor、EventBufferActor、MessageBufferActor、ExternalEventBufferActor五个 Dapr Actor 类型。定义 FastAPI App、DaprActor 与 actor 注册# Define the kernel that is used throughout the process kernel Kernel() # Define a lifespan method that registers the actors with the Dapr runtime asynccontextmanager async def lifespan(app: FastAPI): print(## actor startup ##) await register_fastapi_dapr_actors(actor, kernel, process.factories) yield # Define the FastAPI app along with the DaprActor app FastAPI(titleSKProcess, lifespanlifespan) actor DaprActor(app)这里有几个值得注意的实现细节全局kernel整个流程执行期间共享同一个Kernel实例lifespan 异步注册FastAPI 通过asynccontextmanager定义的 lifespan 钩子在应用启动时调用register_fastapi_dapr_actors将各 Actor 注册进 Dapr 运行时。process.factories是get_process()返回的KernelProcess对象上挂载的步骤工厂字典它会把bstep_factory这样的自定义工厂传递给 Actor第三个参数process.factories实际代码中区别于 README 的简化版本显式传入了步骤工厂用于在 Actor 创建BStep时注入其依赖。在 Dapr 侧register_fastapi_dapr_actors的底层实现见 dapr_actor_registration.py会依次注册ProcessActor使用process_actor_factory注入 kernel 与 factoriesStepActor使用step_actor_factory同样注入 kernel 与 factoriesEventBufferActorMessageBufferActorExternalEventBufferActor。其中ProcessActor与StepActor通过工厂函数注入Kernel依赖——这正是注释中所说的 The ProcessActor and the StepActor require a kernel dependency to be injected during creation。触发流程的 API 端点app.get(/processes/{process_id}) async def start_process(process_id: str): try: context: DaprKernelProcessContext await start( processprocess, kernelkernel, initial_eventCommonEvents.StartProcess, process_idprocess_id, ) kernel_process await context.get_state() c_step_state: KernelProcessStepState[CStepState] next( (s.state for s in kernel_process.steps if s.state.name CStep), None ) c_step_state_validated CStepState.model_validate(c_step_state.state) print(f[FINAL STEP STATE]: CStepState current cycle: {c_step_state_validated.current_cycle}) return JSONResponse(content{processId: process_id}, status_code200) except Exception: return JSONResponse(content{error: Error starting process}, status_code500)start函数见 dapr_kernel_process.py的完整签名与行为如下experimental async def start( process: KernelProcess, initial_event: KernelProcessEvent | str | Enum, process_id: str | None None, max_supersteps: int | None None, **kwargs, ) - DaprKernelProcessContext:process要启动的KernelProcess对象不能为 None且其 state 不能为空initial_event启动进程的初始事件可以是KernelProcessEvent、字符串或枚举本例传入CommonEvents.StartProcess字符串枚举process_id流程实例 ID若为 None 则自动生成新 ID本例显式传入路径参数max_supersteps最大超级步superstep数即所有步骤总共运行的最大次数默认为 None 时由DaprKernelProcessContext取默认值100见 dapr_kernel_process_context.py。start内部会创建DaprKernelProcessContext并通过ActorProxy代理到ProcessActor随后调用start_with_event(initial_event)真正启动流程。start返回的DaprKernelProcessContext可用于查询流程状态——示例端点通过context.get_state()取回执行后的KernelProcess再从步骤状态中读取CStep的current_cycle打印到控制台例如[FINAL STEP STATE]: CStepState current cycle: 3。理解代码Flask 集成详解若使用 Flask定义方式略有不同见 flask_app.pykernel Kernel() app Flask(SKProcess) # Enable DaprActor Flask extension actor DaprActor(app) # Synchronously register actors print(## actor startup ##) register_flask_dapr_actors(actor, kernel) # Create the global event loop loop asyncio.new_event_loop() asyncio.set_event_loop(loop)与 FastAPI 版本的关键差异导入来源不同Flask 使用from flask_dapr.actor import DaprActorFastAPI 使用from dapr.ext.fastapi import DaprActor注册是同步的register_flask_dapr_actors(actor, kernel)在模块加载时同步执行不需要 lifespan 钩子需要手动管理事件循环由于 Flask 的请求处理是同步的代码创建了全局asyncio事件循环并在路由处理函数中用asyncio.run(...)运行start协程依赖导入Flask 版本从semantic_kernel.processes.dapr_runtime导入register_flask_dapr_actors与start。流程定义ProcessBuilder 与事件流流程本体定义在 process/process.pydef get_process() - KernelProcess: # Define the process builder process ProcessBuilder(nameProcessWithDapr) # Add the step types to the builder kickoff_step process.add_step(step_typeKickOffStep) myAStep process.add_step(step_typeAStep) myBStep process.add_step(step_typeBStep, factory_functionbstep_factory) # Initialize the CStep with an initial state and the states current cycle set to 1 myCStep process.add_step(step_typeCStep, initial_stateCStepState(current_cycle1)) # Define the input event and where to send it to process.on_input_event(event_idCommonEvents.StartProcess).send_event_to(targetkickoff_step) # Define the process flow kickoff_step.on_event(event_idCommonEvents.StartARequested).send_event_to(targetmyAStep) kickoff_step.on_event(event_idCommonEvents.StartBRequested).send_event_to(targetmyBStep) myAStep.on_event(event_idCommonEvents.AStepDone).send_event_to(targetmyCStep, parameter_nameastepdata) # Define the fan in behavior once both AStep and BStep are done myBStep.on_event(event_idCommonEvents.BStepDone).send_event_to(targetmyCStep, parameter_namebstepdata) myCStep.on_event(event_idCommonEvents.CStepDone).send_event_to(targetkickoff_step) myCStep.on_event(event_idCommonEvents.ExitRequested).stop_process() # Build the process return process.build()流程拓扑要点on_input_event(StartProcess)声明流程的外部输入事件入口外部调用start(initial_eventCommonEvents.StartProcess)即触发此入口KickOffStep同时发出StartARequested与StartBRequested实现扇出fan-outAStepDone与BStepDone均汇聚到myCStep并分别携带参数名astepdata与bstepdata实现扇入fan-in——CStep.do_it的签名async def do_it(self, context, astepdata: str, bstepdata: str)正好对应这两个参数CStepDone事件将流程送回KickoffStep形成循环ExitRequested事件则调用stop_process()终止流程BStep通过factory_functionbstep_factory指定了自定义工厂用于在运行时构造带依赖的步骤实例。步骤实现与状态建模事件枚举所有流程事件统一定义在 process/steps.py 的CommonEvents枚举中包括StartProcess、StartARequested、StartBRequested、AStepDone、BStepDone、CStepDone、ExitRequested等。KickOffStep 与 AStepclass KickOffStep(KernelProcessStep): KICK_OFF_FUNCTION: ClassVar[str] kick_off kernel_function(nameKICK_OFF_FUNCTION) async def print_welcome_message(self, context: KernelProcessStepContext): print(##### Kickoff ran.) await context.emit_event(process_eventCommonEvents.StartARequested, dataGet Going A) await context.emit_event(process_eventCommonEvents.StartBRequested, dataGet Going B)AStep则模拟一个耗时操作睡眠 1 秒后发出AStepDone事件并携带数据I did A。BStep 与依赖注入工厂BStep是理解 Dapr 状态序列化边界的关键示例它可选地持有ChatCompletionAgent与ChatHistoryAgentThread但这两个字段不会被持久化到 Dapr因为作者刻意没有将它们放进步骤状态模型step state model中——它们属于临时引用ephemeral references通过工厂函数在每次创建步骤实例时注入async def bstep_factory(): Creates a BStep instance with ephemeral references like ChatCompletionAgent. agent ChatCompletionAgent( serviceAzureChatCompletion(credentialAzureCliCredential()), nameecho, instructionsrepeat the input back ) step_instance BStep() step_instance.agent agent step_instance.thread ChatHistoryAgentThread() return step_instance这里使用AzureCliCredential()走 Azure CLI 登录凭据链因此运行该步骤前需要先完成 Azure CLI 登录。BStep.do_it中若self.agent存在会调用agent.get_response(messagesHello from BStep!)并打印响应。设计启示凡是需要持久化的数据必须放进步骤状态模型像 LLM Agent 这类运行时才能构造、无法序列化的对象应通过工厂函数注入并标记为临时引用。CStep 与有状态循环CStep演示了带状态的步骤与终止条件class CStepState(KernelBaseModel): current_cycle: int 1 class CStep(KernelProcessStep[CStepState]): state: CStepState Field(default_factoryCStepState) # The activate method overrides the base class method to set the state in the step. async def activate(self, state: KernelProcessStepState[CStepState]): Activates the step and sets the state. self.state state.state print(f##### CStep activated with Cycle {self.state.current_cycle}.) kernel_function() async def do_it(self, context: KernelProcessStepContext, astepdata: str, bstepdata: str): self.state.current_cycle 1 if self.state.current_cycle 3: print(##### CStep run cycle 3 - exiting.) await context.emit_event(process_eventCommonEvents.ExitRequested) return print(f##### CStep run cycle {self.state.current_cycle}) await context.emit_event(process_eventCommonEvents.CStepDone)CStepState继承自KernelBaseModelPydantic 模型字段current_cycle默认值为 1可被序列化后存入 Dapr 状态activate钩子在步骤每次被激活时执行从KernelProcessStepState[CStepState]中恢复状态并打印当前 Cycle——这就是日志中##### CStep activated with Cycle 1的来源do_it每次执行递增current_cycle达到 3时发出ExitRequested结束流程否则发出CStepDone回到Kickoff继续循环。状态持久化与 Actor 模型原理从 dapr_kernel_process_context.py 可以看到流程启动时会将进程 ID 包装为ActorId通过ActorProxy.create创建ProcessActor的代理。整个运行期间ProcessActor代表整个流程实例负责维护流程状态与调度StepActor代表每个步骤实例负责执行步骤函数、维护步骤状态其构造函数接收kernel、factories以及allowed_module_prefixes默认仅允许从semantic_kernel.开头的模块加载步骤类生产环境不建议放宽EventBufferActor、MessageBufferActor、ExternalEventBufferActor分别负责事件缓冲、消息缓冲与外部事件缓冲保障事件在步骤间的可靠传递。StepActor在激活步骤时会通过KernelProcessStepState恢复步骤状态并调用activate这正是第二次运行从Cycle 3开始这一行为的底层机制——步骤状态在每次执行后被 Dapr Actor 持久化下次激活时自动还原。同理ProcessActor持久化整个流程的编排状态使得同一process_id的多次调用表现为继续执行而非从头开始。常见问题排查服务启动但请求返回 500最常见的原因是 Dapr 容器未运行。启动服务前务必先dapr init并确保 sidecar 容器处于运行状态。BStep 报认证错误bstep_factory使用AzureCliCredential()需要先在环境里执行az login完成 Azure CLI 认证若未配置 Azure OpenAI 资源可自行替换为其他ChatCompletion服务或移除 agent 逻辑。看不到 CStep 的 Cycle 日志确认请求路径中的process_id与期望的实例一致复用同一 ID 会从持久化状态继续执行。想限制流程总步数可通过start(..., max_superstepsN)或DaprKernelProcessContext(max_superstepsN)控制最大超级步数默认 100。延伸阅读Getting Started with Processes 示例不依赖 Dapr 的本地进程框架入门示例含食物准备、订单处理等完整场景Semantic Kernel Dapr Runtime 源码包含五个 Actor 实现、注册逻辑与启动上下文Processes 相关概念文档进程框架的架构决策记录。【免费下载链接】semantic-kernelIntegrate cutting-edge LLM technology quickly and easily into your apps项目地址: https://gitcode.com/GitHub_Trending/se/semantic-kernel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表