Skip to content

工作流引擎: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

  1. 声明插件类型和 ID
  2. 定义依赖的其他插件类型
  3. 实现初始化函数,返回实现 Flow 接口的对象

3.3 插件发现与依赖管理

工作流插件依赖 InternalPlugin 类型,在初始化时通过 ic.GetByType(constants.InternalPlugin) 获取所有已注册的内部插件,然后根据配置中的插件名称进行匹配。

4. 执行顺序与并行控制

4.1 执行流程

工作流的执行遵循严格的顺序:

  1. 步骤顺序执行:Step 按照配置中的顺序依次执行
  2. 动作并行执行:同一个 Step 中的 Flow 并行执行
  3. 并发限制:通过信号量控制工作流的并发数

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 错误分类与处理策略

引擎对不同类型的错误采取不同的处理策略:

  1. PreConditionFailed:前置条件失败,直接返回错误
  2. 其他错误:触发 failover 机制
  3. 上下文取消:触发 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 沙箱创建工作流

典型的沙箱创建工作流可能包含以下步骤:

  1. 网络初始化:分配 IP 地址、配置网络
  2. 存储准备:创建 rootfs、准备存储层
  3. 容器创建:创建容器、配置 cgroup
  4. 挂载配置:挂载文件系统、配置卷

7.2 沙箱销毁工作流

沙箱销毁工作流通常包含:

  1. 停止容器:停止所有进程
  2. 卸载存储:卸载文件系统
  3. 释放网络:释放 IP 地址
  4. 清理资源:删除临时文件

8. 设计优势与最佳实践

8.1 设计优势

  1. 可扩展性:新功能只需实现新的 Flow 插件
  2. 可组合性:通过配置文件灵活组合工作流
  3. 并行性:自动并行无依赖的步骤
  4. 容错性:完善的错误处理和恢复机制

8.2 最佳实践

  1. 插件设计原则

    • 单一职责:每个插件只负责一个功能
    • 幂等性:插件操作应支持重试
    • 可观测性:提供详细的日志和指标
  2. 工作流配置

    • 合理设置并发限制
    • 避免循环依赖
    • 考虑错误处理路径

延伸阅读