工作流引擎:Cubelet 如何编排插件
Cubelet 使用插件式工作流引擎来编排沙箱创建过程。这种设计将复杂的沙箱创建流程分解为可组合、可并行的步骤,每个步骤由独立的插件实现,实现了高度的可扩展性和可维护性。
1. 工作流引擎的架构概览
工作流引擎的核心是 Engine 结构体,它管理多个工作流(Workflow),每个工作流包含多个步骤(Step),每个步骤包含多个动作(Flow)。这种层次化的结构使得复杂的编排逻辑变得清晰。
1.1 核心数据结构
- Engine: 工作流引擎,管理所有工作流
- Workflow: 工作流,代表一个完整的流程(如创建、销毁)
- Step: 步骤,工作流中的一个阶段
- Flow: 动作,实际执行业务逻辑的插件
2. 工作流引擎的初始化
工作流引擎通过 containerd 的插件系统进行初始化。在 plugins/workflow/plugin/plugin.go 中,使用 registry.Register 注册工作流插件:
go
func init() {
registry.Register(&plugin.Registration{
Type: constants.WorkflowPlugin,
ID: constants.WorkflowID.ID(),
Config: &Config{Flows: map[string]flow{}},
Requires: []plugin.Type{
constants.InternalPlugin,
constants.CubeStorePlugin,
},
InitFn: func(ic *plugin.InitContext) (_ interface{}, err error) {
// 初始化逻辑
},
})
}2.1 配置驱动的工作流定义
工作流通过配置文件定义,配置格式如下:
toml
[flows]
[flows.create]
concurrent = 10
actions = [
["network", "storage"], # 第一步:网络和存储并行
["cgroup", "mount"], # 第二步:cgroup和挂载并行
["container"] # 第三步:创建容器
]2.2 插件解析与注册
初始化过程中,引擎会解析配置并将插件名称映射到实际的 Flow 实现:
3. 插件注册机制
3.1 插件接口定义
每个插件必须实现 Flow 接口:
go
type Flow interface {
ID() string
Init(context.Context, *InitInfo) error
Create(context.Context, *CreateContext) error
Destroy(context.Context, *DestroyContext) error
CleanUp(context.Context, *CleanContext) error
}3.2 插件注册流程
插件通过 containerd 的插件系统注册,需要在 init() 函数中调用 registry.Register:
- 声明插件类型和 ID
- 定义依赖的其他插件类型
- 实现初始化函数,返回实现
Flow接口的对象
3.3 插件发现与依赖管理
工作流插件依赖 InternalPlugin 类型,在初始化时通过 ic.GetByType(constants.InternalPlugin) 获取所有已注册的内部插件,然后根据配置中的插件名称进行匹配。
4. 执行顺序与并行控制
4.1 执行流程
工作流的执行遵循严格的顺序:
- 步骤顺序执行:Step 按照配置中的顺序依次执行
- 动作并行执行:同一个 Step 中的 Flow 并行执行
- 并发限制:通过信号量控制工作流的并发数
4.2 并行执行机制
使用 errgroup 实现并行执行:
go
func (e *Engine) parallelRunSteps(do string, ctx context.Context, opts ReqContext, step *Step) error {
eg, ctxWithCancel := errgroup.WithContext(ctx)
for i := range step.Actions {
flow := step.Actions[i]
eg.Go(func() (err error) {
// 执行单个 Flow
})
}
return eg.Wait()
}4.3 并发控制
每个工作流都有并发限制,通过信号量实现:
go
type Workflow struct {
Name string
MaxConcurrent int64
Limiter *semaphore.Limiter
Steps []*Step
}5. 错误处理与恢复机制
5.1 错误分类与处理策略
引擎对不同类型的错误采取不同的处理策略:
- PreConditionFailed:前置条件失败,直接返回错误
- 其他错误:触发 failover 机制
- 上下文取消:触发 failover 机制
5.2 Failover 机制
当创建过程中发生错误时,引擎会触发 failover 机制,自动销毁已创建的资源:
go
func (e *Engine) failover(ctx context.Context, opts ReqContext) {
// 创建销毁上下文
destroyOpt := &DestroyContext{
DestroyInfo: &cubebox.DestroyCubeSandboxRequest{
SandboxID: sandboxID,
},
}
destroyOpt.IsRollBack = true
// 执行销毁
if err = e.Destroy(ctx, destroyOpt); err != nil {
log.G(ctx).Fatalf("fail over Destroy fatal error:%v", err)
}
}5.3 清理机制
当销毁过程中发生错误时,引擎会触发清理机制:
go
func (e *Engine) cleanUp(ctx context.Context, opts ReqContext) {
if e.cleanupFlow == nil {
return
}
// 执行清理工作流
if err := e.cleanupFlow.Create(ctxTmp, rOpts); err != nil {
CubeLog.WithContext(ctxTmp).Fatalf("cleanupFlow.Create fatal error:%v", err)
}
}6. 监控与指标收集
6.1 性能指标
引擎会收集每个步骤的执行时间:
go
type Metric struct {
id string
err error
duration time.Duration
}6.2 并发监控
可以实时监控工作流的并发状态:
GetOnFlyingRequest():当前正在执行的请求数GetPeakRequest():历史峰值请求数
7. 实际工作流示例
7.1 沙箱创建工作流
典型的沙箱创建工作流可能包含以下步骤:
- 网络初始化:分配 IP 地址、配置网络
- 存储准备:创建 rootfs、准备存储层
- 容器创建:创建容器、配置 cgroup
- 挂载配置:挂载文件系统、配置卷
7.2 沙箱销毁工作流
沙箱销毁工作流通常包含:
- 停止容器:停止所有进程
- 卸载存储:卸载文件系统
- 释放网络:释放 IP 地址
- 清理资源:删除临时文件
8. 设计优势与最佳实践
8.1 设计优势
- 可扩展性:新功能只需实现新的 Flow 插件
- 可组合性:通过配置文件灵活组合工作流
- 并行性:自动并行无依赖的步骤
- 容错性:完善的错误处理和恢复机制
8.2 最佳实践
插件设计原则:
- 单一职责:每个插件只负责一个功能
- 幂等性:插件操作应支持重试
- 可观测性:提供详细的日志和指标
工作流配置:
- 合理设置并发限制
- 避免循环依赖
- 考虑错误处理路径
延伸阅读
- 沙箱在节点上的一生 - 了解沙箱的完整生命周期
- 存储与快照:rootfs 怎么来的 - 深入存储层实现
- Cubelet 贡献指南 - 参与 Cubelet 开发