本节深入探讨 Dapr 运行时的内部工作原理。面向维护者、贡献者以及对 Dapr 构建块底层实现细节感兴趣的任何人。
与面向用户的 API 参考不同,这些文档侧重于:
- 运行时如何处理请求。
- 内部状态转换。
- 与组件接口的交互。
- 协议层面的细节(gRPC 和 HTTP)。
构建块
选择一个构建块以探索其内部协议和机制:
This is the multi-page printable view of this section. Click here to print.
本节深入探讨 Dapr 运行时的内部工作原理。面向维护者、贡献者以及对 Dapr 构建块底层实现细节感兴趣的任何人。
与面向用户的 API 参考不同,这些文档侧重于:
选择一个构建块以探索其内部协议和机制:
本文档从底层详细说明 Dapr 工作流协议与运行时约定。面向的读者是构建工作流 Worker 的 SDK 作者以及演进 Dapr 边车工作流引擎的运行时维护者。
Dapr 工作流采用"边车即调度器"模式:Dapr 运行时(边车)作为工作流引擎,应用 SDK 作为工作流 Worker。所有控制与执行流量均通过 gRPC 传输。
协议界面分为两类:
TaskHubSidecarService)。核心组件
工作流引擎(Dapr 边车)
管理工作流状态转换、历史持久化、编排与活动任务的调度,以及可靠交付语义。默认情况下,利用 Dapr Actors 作为后端,实现可持久化、分区的执行。
工作流 Worker(应用 SDK)
连接到边车,轮询编排与活动工作项,执行用户定义的逻辑,并将结果、失败与心跳返回给引擎。编排逻辑必须是确定性的;活动逻辑无需是确定性的。
编排
定义工作流的确定性协调器。引擎通过重放历史来驱动编排,重建状态并调度出站任务(活动、子编排、定时器、外部事件)。
活动
原子工作单元。活动执行保证至少一次,并将结果或失败报告回引擎。建议幂等,并在上下文中提供任务执行标识符以辅助实现幂等。
状态存储与后端
工作流历史与状态持久化保存。引擎通常在所选持久化层上实现任务中心模式,并使用 Dapr Actors 作为默认的可靠性底层。
Dapr 工作流基于持久任务框架(DTFx)执行语义:
编排器从其事件历史重放,以重建确定性状态。所有非确定性操作(时间、随机值、I/O)必须由引擎中介(例如通过定时器、活动调用、外部事件)。
除引擎中介效果外,编排器代码必须无副作用。控制流在重放期间必须可复现。
活动可能被多次传递。引擎确保工作流状态提交是幂等的,并恰好一次应用。
边车拥有调度权,并在向 Worker 分发工作前持久化所有历史/事件。从引擎视角看,Worker 是无状态执行器。
注意:关于确切的 RPC 形状、错误码和语义,参见 管理 API 规范。
注意:参见定义了
TaskHubSidecarService约定、负载架构与排序规则的 执行 API 规范。
StartWorkflow。ExecutionStarted)并实例化实例。Workflow 管理 API 允许 Dapr 客户端控制 workflow 实例的生命周期。这些 API 通过标准 Dapr gRPC 端点公开,通常通过 SDK 提供。
管理 API 是 dapr.proto.runtime.v1 中 Dapr 服务的一部分。虽然可能存在多个版本(Alpha1、Beta1),以下描述的是当前的实现逻辑。
启动一个新的 workflow 实例。
请求(StartWorkflowRequest):
| 字段 | 类型 | 描述 |
|---|---|---|
instance_id | string | 可选。workflow 实例的唯一标识符。如果未提供,Dapr 将生成一个随机的 UUID。 |
workflow_component | string | 要使用的 workflow 组件的名称。目前,Dapr 使用内置引擎。 |
workflow_name | string | 要执行的 workflow 定义的名称。 |
options | map<string, string> | 可选。组件特定选项。 |
input | bytes | 可选。workflow 实例的输入数据,通常是 JSON 序列化字符串。 |
响应(StartWorkflowResponse):
| 字段 | 类型 | 描述 |
|---|---|---|
instance_id | string | 已启动的 workflow 实例的 ID。 |
检索 workflow 实例的当前状态和元数据。
请求(GetWorkflowRequest):
| 字段 | 类型 | 描述 |
|---|---|---|
instance_id | string | 要查询的 workflow 实例的 ID。 |
workflow_component | string | workflow 组件的名称。 |
响应(GetWorkflowResponse):
| 字段 | 类型 | 描述 |
|---|---|---|
instance_id | string | workflow 实例的 ID。 |
workflow_name | string | workflow 的名称。 |
created_at | Timestamp | 实例创建的时间。 |
last_updated_at | Timestamp | 实例最后更新的时间。 |
runtime_status | string | 状态(例如 RUNNING、COMPLETED、FAILED、TERMINATED、PENDING)。 |
properties | map<string, string> | 额外的组件特定元数据。 |
强制终止正在运行的 workflow 实例。
请求(TerminateWorkflowRequest):
| 字段 | 类型 | 描述 |
|---|---|---|
instance_id | string | 要终止的 workflow 实例的 ID。 |
workflow_component | string | workflow 组件的名称。 |
向正在运行的 workflow 实例发送事件。
请求(RaiseEventWorkflowRequest):
| 字段 | 类型 | 描述 |
|---|---|---|
instance_id | string | workflow 实例的 ID。 |
workflow_component | string | workflow 组件的名称。 |
event_name | string | 要引发的事件的名称。 |
event_data | bytes | 与事件关联的数据。 |
暂停或恢复 workflow 实例。
请求(PauseWorkflowRequest / ResumeWorkflowRequest):
| 字段 | 类型 | 描述 |
|---|---|---|
instance_id | string | workflow 实例的 ID。 |
workflow_component | string | workflow 组件的名称。 |
移除与 workflow 实例关联的所有状态和历史记录。这通常只能对已完成、失败或已终止的实例执行。
请求(PurgeWorkflowRequest):
| 字段 | 类型 | 描述 |
|---|---|---|
instance_id | string | workflow 实例的 ID。 |
workflow_component | string | workflow 组件的名称。 |
检索 workflow 实例 ID 列表,可选择按状态或名称筛选。这目前是内部 Task Hub 协议的一部分,用于管理工具中的分页。
请求(ListInstanceIDsRequest):
| 字段 | 类型 | 描述 |
|---|---|---|
page_size | int32 | 要返回的 ID 的最大数量。 |
continuation_token | string | 用于检索下一页结果的不透明 token。 |
响应(ListInstanceIDsResponse):
| 字段 | 类型 | 描述 |
|---|---|---|
instance_ids | repeated string | 实例 ID 列表。 |
continuation_token | string | 用于下一页结果的 token。 |
continuation_token 是由 Dapr 运行时(以及底层状态存储)生成的不透明字符串。其目的是允许客户端可靠地对大量 workflow 实例进行分页,而无需一次性将所有 ID 加载到内存中。
SDK 对继续标记的要求:
ListInstanceIDsResponse 时,如果打算获取更多结果,应存储 continuation_token。continuation_token 传递到下一个 ListInstanceIDsRequest。continuation_token 表示没有更多页面可检索。运行时行为:
运行时从状态存储的 KeysLike 操作派生此 token。由于它与底层数据库的分页机制相关联,因此该 token 可能具有过期时间或与初始请求中使用的特定查询参数(如 page_size)相关联。
边车接收这些请求并将其转换为底层 durabletask-go 客户端的操作。例如,StartWorkflow 调用后端以创建新的业务流程实例并排队 ExecutionStarted 事件。
Workflow Execution API 是一个低级 gRPC 协议,Dapr Workflow SDK 通过它充当"Worker"。SDK 通过此协议连接 Dapr sidecar 以轮询任务并报告完成状态。
该服务名为 TaskHubSidecarService。
service TaskHubSidecarService {
rpc GetWorkItems(GetWorkItemsRequest) returns (stream WorkItem);
rpc CompleteOrchestratorTask(OrchestratorResponse) returns (CompleteBatchResponse);
rpc CompleteActivityTask(ActivityResponse) returns (CompleteBatchResponse);
// ... other management methods
}
GetWorkItems 的长连接双向流。WorkItem 消息。CompleteOrchestratorTask 并附带要执行的操作列表。CompleteActivityTask 并附带结果或失败信息。TaskHubSidecarService打开流以接收 orchestration 和 activity 的工作项。
Request (GetWorkItemsRequest):
通常为空或包含 worker 元数据。
Response (stream WorkItem):
WorkItem 可以是以下之一:
orchestrator_item:包含某个 orchestration 的历史和新事件。activity_item:包含单个 activity 任务的详情。报告 orchestration 执行的结果。
Request (OrchestratorResponse):
instance_id:workflow 实例的 ID。actions:OrchestratorAction 消息列表。custom_status:可选的用户定义状态字符串。OrchestratorAction 类型:
ScheduleTask:调度一个新的 activity。CreateTimer:调度一个持久化计时器。CreateSubOrchestration:启动一个子 workflow。CompleteOrchestration:将 workflow 标记为完成(成功或失败)。TerminateOrchestration:强制终止实例。SendEvent:向另一个 workflow 发送事件。报告 activity 执行的结果。
Request (ActivityResponse):
instance_id:workflow 实例的 ID。task_id:activity 任务的唯一 ID。completion_token:ActivityWorkItem 中收到的 opaque token。result:activity 的序列化输出(如果成功)。failure_details:错误详情(如果失败)。Dapr 中的 workflow 是事件溯源的。Orchestration 的状态通过重放一系列 HistoryEvent 消息来重建。
常见事件类型:
ExecutionStarted:初始事件,包含 workflow 名称和输入。TaskScheduled:一个 activity 被调度。TaskCompleted:一个 activity 成功完成。TaskFailed:一个 activity 失败。TimerCreated:一个计时器被调度。TimerFired:一个计时器到期。OrchestrationCompleted:workflow 完成。用于报告来自 activity 或 orchestration 的错误。
error_type:标识错误类型的字符串。error_message:人类可读的错误消息。stack_trace:可选的堆栈跟踪。is_non_retriable:布尔标志。GetWorkItems 是服务端到客户端的流。Dapr 在工作可用时将工作推送到 SDK。OrchestratorWorkItem 中提供的历史记录,以避免重新执行已记录的操作。本文档从协议层面描述编排的生命周期,特别是 Dapr 引擎与 SDK 如何交互以可靠地执行工作流逻辑。
Dapr Workflow 使用 事件溯源 和 重放 来维护状态。Dapr 不是保存整个 worker 进程的状态(栈、变量等),而是保存已发生事件的历史记录。
GetWorkItems 流向 SDK 发送一个 OrchestratorWorkItem。该工作项包含工作流实例的完整历史记录以及任何新事件(例如,activity 完成或外部事件)。CompleteOrchestratorTask 请求。该请求包含引擎应执行的一系列 Actions(例如 ScheduleTask、CreateTimer)。假设一个工作流:Activity A -> Activity B。
ExecutionStarted 事件加入队列。[ExecutionStarted] 的 OrchestratorWorkItem。Activity A。Activity A 不在其中。[ScheduleTask(Activity A)] 的 CompleteOrchestratorTask。TaskScheduled(Activity A)。TaskCompleted(Activity A, result="foo")。[ExecutionStarted, TaskScheduled(A), TaskCompleted(A)] 的 OrchestratorWorkItem。Activity A。SDK 在历史记录中找到 TaskCompleted(A)。返回 "foo"。Activity B。Activity B 不在其中。[ScheduleTask(Activity B)] 的 CompleteOrchestratorTask。TaskScheduled(Activity B)。[CompleteOrchestration(result="final")] 的 CompleteOrchestratorTask。OrchestrationCompleted 并将实例标记为 COMPLETED。编排函数必须是确定性的。它不能使用:
CurrentUtcDateTime)。当工作流已经在运行时,您可能需要更新其逻辑。然而,由于工作流是基于重放的,直接更改逻辑会破坏进行中实例的确定性。
Dapr 提供了 Patching 机制(例如 ctx.IsPatched("patch-id"))来安全地引入更改:
Dapr 还提供了 命名版本控制 机制,其中 SDK 维护可用命名工作流版本的注册表。当它收到通过名称初始化新工作流的请求时,它将查询注册表以确定该名称是否匹配与指定工作流名称不同的工作流版本,并负责将请求重定向到预期的"最新"版本。
当引擎检测到需要手动干预或代码修复才能继续的不可恢复条件时,工作流实例进入 STALLED 状态。常见原因包括:
当停滞时,实例停止执行但保留在系统中。一旦根本问题得到解决(例如,部署了正确的代码版本),实例就可以恢复或将在下一个事件时自动恢复。
SDK 必须高效地搜索历史记录。通常,这是通过维护执行期间遇到的任务计数器并将它们与历史中的事件序列进行匹配来完成的。
SDK 需要一种机制,在任务已调度但尚未完成时停止编排函数的执行,同时不丢失稍后重新启动它的能力。
Activity 是 Dapr Workflow 中的基本工作单元。与编排不同,Activity 不会重放,也不需要具有确定性。它们每次"调度"仅执行一次(尽管可能会发生重试)。
ScheduleTask 操作来请求一个 Activity。GetWorkItems 流发送一个 ActivityWorkItem。ActivityWorkItem,其中包含:name:要执行的 Activity 的名称。input:Activity 的输入数据。instance_id:调度该 Activity 的工作流实例的 ID。task_id:此特定 Activity 执行的唯一标识符。task_execution_id:此特定 Activity 的特定尝试的唯一标识符。这对于在 Activity 逻辑中实现幂等性很有用。completion_token:一个不透明 token,用于将响应与此特定工作项关联。CompleteActivityTask 请求。result 字段中提供序列化输出。failure_details(错误消息、类型、堆栈跟踪)。task_execution_id(也称为 Task Execution Key)是一个唯一的、运行时生成的字符串(通常是 UUID),用于标识执行 Activity 任务的特定尝试。
虽然 Workflow SDK 在工作项之间通常是无状态的,但 task_execution_id 为 Activity Worker 提供了关键的上下文:
task_execution_id 作为幂等性键。task_id(在工作流中特定步骤保持不变)不同,task_execution_id 在引擎每次重试 Activity 时都会变化(例如,由于超时或 worker 崩溃)。task_execution_id,worker 可以确定它是否是一个不再需要其结果的"zombie"。task_execution_id 公开给 Activity 实现逻辑(例如,通过 ActivityContext)。completion_token 是由 Dapr 运行时生成的不透明字符串,并作为 ActivityWorkItem 的一部分传递给 SDK。
completion_token 将 ActivityResponse(来自 CompleteActivityTask)可靠地匹配到它分发的原始任务。completion_token。如果原始"zombie" worker 最终使用旧 token 响应,边车可以轻松识别并忽略延迟的响应。ActivityWorkItem 中捕获 completion_token。CompleteActivityTask 发送的 ActivityResponse 中包含完全相同的 completion_token。ActivityContext 中)。在 Dapr 运行时中(特别是使用 Actors 后端时),Activity 被表示为 actors。每个 Activity 执行都有一个唯一的 Task Activity ID(也称为 Activity Actor ID)。
ID 遵循特定模式:
{workflowInstanceID}::{taskID}::{generation}
这个唯一的 ID 确保 Activity 执行被隔离,并且可以在重试和重启期间可靠地跟踪。
Dapr 根据编排中定义的策略处理 Activity 重试(如果 SDK 支持在 ScheduleTask 操作中定义重试策略)。如果 Activity 失败并且存在重试策略,引擎将在指定的延迟后重新将 Activity 任务加入队列。
从 Activity worker 的角度来看,重试只是一个具有相同名称和输入的新 ActivityWorkItem,但可能具有不同的 task_id(或相同的,取决于后端实现)。
因为 Activity 可能被执行多次(例如,如果 worker 在执行之后但在报告完成之前崩溃),所以建议 Activity 逻辑尽可能具有幂等性。
| 功能 | 编排 | Activity |
|---|---|---|
| 执行方式 | 基于重放(确定性) | 直接执行 |
| 状态 | 通过历史事件管理 | 无内部工作流状态 |
| 副作用 | 禁止(必须使用 Activity) | 允许(IO、数据库等) |
| 生命周期 | 可以长时间运行(天/月) | 通常短期 |
| 连接性 | 通过 GetWorkItems 连接 | 通过 GetWorkItems 连接 |
Dapr Workflows 采用事件溯源模式,这意味着工作流的状态是从一系列事件中推导出来的。本文档描述了 Dapr 如何存储和管理这些历史记录和状态。
默认情况下,Dapr Workflow 引擎使用 Dapr Actors 作为其存储后端。每个工作流实例都映射到一个唯一的 actor 实例。这提供了以下特性:
工作流实例(actor)的状态由以下几个组件构成:
存储有关该实例的高级信息:
instance_id:工作流的唯一 ID。name:工作流的名称。status:当前运行时状态(如 Running、Completed、Stalled 等)。version:工作流版本的名称及任何活跃的补丁。created_at:创建时间戳。last_updated_at:最后活动时间戳。input:原始输入数据。output:最终输出数据(如已完成的)。一系列 HistoryEvent 对象,记录工作流中发生的所有事件。为了优化大型历史记录,Dapr 通常将历史事件以分块或单独键值的形式存储在状态存储中:
wf-history-<instance_id>-<index>TaskScheduled、TaskCompleted)。已发生但尚未被编排器处理(重放)的事件集合。包括:
当编排器下次运行时,它会"排空"收件箱,将这些事件移入历史记录,然后重放逻辑。
当 Worker(SDK)收到工作项时,Dapr 提供历史事件。SDK 通过按顺序重放这些事件来重建编排的内部状态(例如本地变量、当前执行点)。
历史记录是"真相的来源"。如果编排代码以非确定性的方式发生变化(例如在现有代码中间添加新的 activity 调用),重放将失败,因为代码的请求将与记录的历史记录不匹配。
当清除工作流时,Dapr 会从状态存储中删除元数据及所有相关的历史事件键。这通常是为了在完成工作流后进行清理。
Dapr 工作流支持工作流定义的版本控制,允许你在现有实例继续运行其原始逻辑的同时更新工作流逻辑。
向 Dapr 引擎注册工作流时,你可以提供版本名称。这允许同一工作流的多个版本共存。
SDK 使用其任务注册表中的 AddVersionedOrchestrator(或类似)方法来注册版本化工作流。
// 示例(内部注册表 API)
registry.AddVersionedOrchestrator("MyWorkflow", "v2", true, MyWorkflowV2)
registry.AddVersionedOrchestrator("MyWorkflow", "v1", false, MyWorkflowV1)
Dapr 边车会跟踪实例正在运行的工作流命名版本,以及在该工作流执行过程中已应用的补丁列表。此信息存储在工作流历史记录中,位于 OrchestratorStarted 事件的 version 字段内。
message OrchestrationVersion {
string name = 1;
repeated string patches = 2;
}
当实例恢复时(例如,在某个活动完成后),Dapr 引擎负责处理因客户端版本不匹配而导致的停滞。它通过监控 placement 表的变化来实现这一点,这表明有新的 SDK 客户端已连接(可能代表应用程序的较新副本实例),并调度重放最后一个事件以尝试重试该操作。这一次,如果 SDK 能够成功完成任务,运行时将把工作流状态从 Stalled 更改回 Running。