【问题标题】:Cadence throwing WorkflowRejectedExecutionError when executing workflow with many child workflows/activities执行具有许多子工作流/活动的工作流时,Cadence 抛出 WorkflowRejectedExecutionError
【发布时间】:2020-08-21 12:29:44
【问题描述】:

我正在评估使用 Cadence 执行长时间运行的批量操作。我有以下(Kotlin)代码:

class UpdateNameBulkWorkflowImpl : UpdateNameBulkWorkflow {

    private val changeNamePromises = mutableListOf<Promise<ChangeNameResult>>()

    override fun updateNames(newName: String, entityIds: Collection<String>) {
        entityIds.forEach { entityId ->
            val childWorkflow = Workflow.newChildWorkflowStub(
                    UpdateNameBulkWorkflow.UpdateNameSingleWorkflow::class.java
            )
            val promise = Async.function(childWorkflow::setName, newName, entityId)

            changeNamePromises.add(promise)
        }

        val allDone = Promise.allOf(changeNamePromises)
        allDone.get()
    }

    class UpdateNameSingleWorkflowImpl : UpdateNameBulkWorkflow.UpdateNameSingleWorkflow {
        override fun setName(newName: String, entityId: String): SetNameResult {
            return Async.function(activities::setName, newName, entityId).get()
        }
    }
}

这适用于较少数量的实体,但我很快遇到了以下异常:

java.lang.RuntimeException: Failure processing decision task. WorkflowID=b5327d20-6ea6-4aba-b863-2165cb21e038, RunID=c85e2278-e483-4c81-8def-f0cc0bd309fd
    at com.uber.cadence.internal.worker.WorkflowWorker$TaskHandlerImpl.wrapFailure(WorkflowWorker.java:283) ~[cadence-client-2.7.4.jar:na]
    at com.uber.cadence.internal.worker.WorkflowWorker$TaskHandlerImpl.wrapFailure(WorkflowWorker.java:229) ~[cadence-client-2.7.4.jar:na]
    at com.uber.cadence.internal.worker.PollTaskExecutor.lambda$process$0(PollTaskExecutor.java:76) ~[cadence-client-2.7.4.jar:na]
    at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128) ~[na:na]
    at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628) ~[na:na]
    at java.base/java.lang.Thread.run(Thread.java:834) ~[na:na]
Caused by: com.uber.cadence.internal.sync.WorkflowRejectedExecutionError: java.util.concurrent.RejectedExecutionException: Task java.util.concurrent.FutureTask@7f17a605[Not completed, task = java.util.concurrent.Executors$RunnableAdapter@7fa9f240[Wrapped task = com.uber.cadence.internal.sync.WorkflowThreadImpl$RunnableWrapper@1a27000b]] rejected from java.util.concurrent.ThreadPoolExecutor@22188bd0[Running, pool size = 600, active threads = 600, queued tasks = 0, completed tasks = 2400]
    at com.uber.cadence.internal.sync.WorkflowThreadImpl.start(WorkflowThreadImpl.java:281) ~[cadence-client-2.7.4.jar:na]
    at com.uber.cadence.internal.sync.AsyncInternal.execute(AsyncInternal.java:300) ~[cadence-client-2.7.4.jar:na]
    at com.uber.cadence.internal.sync.AsyncInternal.function(AsyncInternal.java:111) ~[cadence-client-2.7.4.jar:na]
...

看来我很快就会耗尽线程池,Cadence 无法安排新任务。

我通过将updateNames 的定义更改为:

    override fun updateNames(newName: String, entityIds: Collection<String>) {

        entityIds.chunked(200).forEach { sublist ->
            val promises = sublist.map { entityId ->
                val childWorkflow = Workflow.newChildWorkflowStub(
                        UpdateNameBulkWorkflow.UpdateNameSingleWorkflow::class.java
                )
                Async.function(childWorkflow::setName, newName, entityId)
            }

            val allDone = Promise.allOf(promises)
            allDone.get()
        }
    }

这基本上处理了 200 个块中的项目,并等待每个块完成,然后再移动到下一个块。我担心这将如何执行(一个块中的单个错误将在重试时停止处理以下块中的所有记录)。我还关心 Cadence 在发生崩溃时恢复此功能的进度的能力。

我的问题是:是否有一种惯用的 Cadence 方式来执行此操作,而不会立即导致资源耗尽?是我使用了错误的技术还是这只是一种幼稚的方法?

【问题讨论】:

    标签: java kotlin cadence-workflow


    【解决方案1】:

    Cadence 工作流程对单个工作流程运行的大小限制相对较小。它随着并行工作流运行的数量而扩展。因此,在单个工作流中执行大量任务是一种反模式。

    惯用的解决方案是:

    • 运行一个有限大小的 chank,然后调用 continue as new。这样,单个运行大小是有界的。
    • 使用分层工作流。具有 1k 个子项的单个父项,每个子项执行 1k 个活动,允许执行 100 万个活动,同时保持每个工作流历史记录大小有界。

    【讨论】:

    • 您好,感谢您的回复马克西姆!当您说“在单个工作流中执行大量任务是一种反模式”时,这种情况下的“任务”是否意味着子工作流和活动?在第一个示例中,我给出的父级不直接启动任何活动,但它确实等待许多子工作流完成,因此在我的测试中,我有一个包含 5000 个子工作流的父工作流,每个子工作流都有 1 个活动。这就是导致资源耗尽的原因。
    • 我相信它应该适用于 5k 个孩子。你会针对github.com/temporalio/java-sdk/issues 提出问题吗?无论如何,我会推荐“继续作为新的”方法,因为它可以扩展到任意数量的孩子。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多