【问题标题】:Invoking the same activity inside a loop in cadence workflow在节奏工作流的循环内调用相同的活动
【发布时间】:2020-09-27 00:57:59
【问题描述】:

我有一个关于 cadence 工作流程的问题,我们可以在 for 循环中使用不同的输入调用相同的活动吗?该代码将是确定性的吗?如果执行工作流的工作人员在执行期间停止并稍后重新启动,cadence 是否能够在重新构建工作流时重播事件。

例如,我有以下代码。

   func init() {
    workflow.RegisterWithOptions(SampleWorkFlow, workflow.RegisterOptions{Name: "SampleWorkFlow"})
    activity.RegisterWithOptions(SampleActivity, activity.RegisterOptions{Name: "SampleActivity"})
    activity.RegisterWithOptions(SecondActivity, activity.RegisterOptions{Name: "SecondActivity"})
}

// SampleWorkFlow comment
func SampleWorkFlow(ctx workflow.Context, input string) error {

    fmt.Println("Workflow started")
    ctx = workflow.WithTaskList(ctx, sampleTaskList)
    ctx = workflow.WithActivityOptions(ctx, conf.ActivityOptions)

    var result string
    err := workflow.ExecuteActivity(ctx, "SampleActivity", input, "string-value").Get(ctx, &result)
    if err != nil {
        return err
    }

    for i := 1; i <= 10; i++ {
        value := i
        workflow.Go(ctx, func(ctx workflow.Context) {
            err := workflow.ExecuteActivity(ctx, "SecondActivity", input, value).Get(ctx, &result)
            if err != nil {
                log.Println("err=", err)
            }
        })
    }

    return nil

}

// SampleActivity comment
func SampleActivity(ctx context.Context, value, v1 string) (string, error) {
    fmt.Println("Sample activity start")
    for i := 0; i <= 10; i++ {
        fmt.Println(i)
    }
    return "Hello " + value, nil
}

// SecondActivity comment
func SecondActivity(ctx context.Context, value int) (string, error) {

    fmt.Println("Second  activity start")

    fmt.Println("value=", value)
    fmt.Println("Second activity going to end")
    return "Hello " + fmt.Sprintf("%d", value), nil
}

这里,第二个活动在 for 循环内并行调用。 我的第一个问题是这段代码是确定性的吗?

假设在循环 5 次迭代后,当 i = 5 时,执行此工作流的工作人员终止,如果工作流在另一个中启动,cadence 将能够重播事件 工人?

你能回答我的问题吗?

【问题讨论】:

    标签: go workflow cadence-workflow temporal-workflow


    【解决方案1】:

    是的,此代码是确定性的。它不调用任何非确定性操作(如随机或 UUID 生成)并使用workflow.Go 来启动一个 goroutine。所以它是确定性的。代码的复杂性在定义其确定性方面没有作用。

    不相关的 nit。 在您的示例中无需使用 goroutine,因为 ExecuteActivity 调用通过返回 Future 已经是非阻塞的。 所以样本可以简化为:

    func SampleWorkFlow(ctx workflow.Context, input string) error {
    
        fmt.Println("Workflow started")
        ctx = workflow.WithTaskList(ctx, sampleTaskList)
        ctx = workflow.WithActivityOptions(ctx, conf.ActivityOptions)
    
        var result string
        err := workflow.ExecuteActivity(ctx, "SampleActivity", input, "string-value").Get(ctx, &result)
        if err != nil {
            return err
        }
    
        for i := 1; i <= 10; i++ {
           workflow.ExecuteActivity(ctx, "SecondActivity", input, i)
        }
        return nil
    }
    

    请注意,此示例仍然可能不会以您期望的方式执行,因为它无需等待活动完成即可完成工作流。所以这些活动甚至都不会开始。

    这是等待活动完成的代码:

    func SampleWorkFlow(ctx workflow.Context, input string) error {
    
        fmt.Println("Workflow started")
        ctx = workflow.WithTaskList(ctx, sampleTaskList)
        ctx = workflow.WithActivityOptions(ctx, conf.ActivityOptions)
    
        var result string
        err := workflow.ExecuteActivity(ctx, "SampleActivity", input, "string-value").Get(ctx, &result)
        if err != nil {
            return err
        }
        var results []workflow.Future
        for i := 1; i <= 10; i++ {
            future := workflow.ExecuteActivity(ctx, "SecondActivity", input, i)
            results = append(results, future)
        }
        for i := 0; i < 10; i++ {
            var result string
            err := results[i].Get(ctx, &result)
            if err != nil {
                log.Println("err=", err)
            }
        }
        return nil
    }
    

    【讨论】:

    • 非常感谢您的详细回答!
    猜你喜欢
    • 2022-09-29
    • 1970-01-01
    • 2020-07-24
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-10-06
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多