Temporal Workflow调用Activity时的线程阻塞与扩展相关问题
Temporal Java SDK 常见问题解答
代码示例
Map<String, String> activityResult = customActivity.fetchCustomMap(); Workflow.sleep(Duration.ofSeconds(5)); Map<String, String> activityResult1 = customActivity.fetchCustomMap(); System.out.println(Thread.currentThread() + "activity result " + activityResult); System.out.println(Thread.currentThread() + "activity result1 " + activityResult1); Workflow.sleep(Duration.ofSeconds(2)); System.out.println(Thread.currentThread() + "workflow execution ends"); Workflow.sleep(Duration.ofSeconds(2));
疑问与解答
补充确认:Activity调用时发起线程是否全程阻塞?
不会。Temporal SDK拦截Activity调用后,会把操作转化为事件提交给Temporal Server,发起调用的线程T1会被立即释放,不会一直等待Activity执行完成。当Server收到Activity执行完成的事件后,会触发Workflow继续执行,此时SDK会重新调度Workflow代码,从之前暂停的位置(Activity调用的下一行)继续执行,处理的线程不一定是原来的T1。
问题1:同步模式下如何扩展?Workflow是否固定在某个Worker实例上?
- 扩展方式:同步调用Activity的模式不影响扩展,只需增加Worker实例数量就能提升整体处理能力,Temporal Server会自动将Activity任务分发给注册了对应Activity类型的Worker。
- Workflow实例不会固定在某个Worker上。当Workflow需要继续执行(比如Activity完成、sleep结束)时,Server会将Workflow任务分发给任意一个注册了对应Workflow类型的空闲Worker。即使之前处理该Workflow的Worker下线,其他Worker也能接手继续执行。
问题2:如何将Activity异步分发到不同Worker?
你当前的代码已经是异步分发模式——Temporal Server本身就会把Activity任务路由到任意可用的、注册了该Activity类型的Worker,无需额外修改代码。如果需要更细粒度控制(比如指定特定Worker组执行),可以在调用Activity时通过ActivityOptions设置taskQueue参数,将不同Activity分配到不同任务队列,再让对应的Worker监听该队列即可。
示例代码片段:
ActivityOptions options = ActivityOptions.newBuilder() .setTaskQueue("specific-task-queue") .build(); CustomActivity customActivity = Workflow.newActivityStub(CustomActivity.class, options);
启动Worker时指定任务队列:
WorkerFactory factory = WorkerFactory.newInstance(client); Worker worker = factory.newWorker("specific-task-queue"); worker.registerActivityImplementation(CustomActivity.class, new CustomActivityImpl());
问题3:重放时缓存机制如何运作?不同Worker是否需知晓已执行的所有Activity?
- 重放机制:Temporal的Workflow重放是通过重新执行Workflow代码、匹配历史事件来恢复状态的。SDK会缓存已执行的Activity调用结果,重放时遇到之前已完成的Activity调用,不会再次发起实际执行,而是直接使用缓存结果,保证Workflow状态与之前一致。
- 不同Worker不需要知晓已执行的所有Activity。Worker只需要处理当前分配给自己的任务(Workflow或Activity),已执行的Activity结果会记录在Temporal Server的Workflow历史中,当其他Worker接手该Workflow进行重放或继续执行时,会从Server获取历史事件,SDK自动处理缓存和状态恢复,无需Worker间同步信息。
内容的提问来源于stack exchange,提问作者user1189332
相关产品推荐
相关产品推荐

