无法向已调度的Temporal工作流发送信号:WorkflowNotFoundError问题求助
无法向已调度的Temporal工作流发送信号:WorkflowNotFoundError问题求助
我帮你梳理下你遇到的问题和可能的解决办法,看起来你在Temporal调度工作流接收信号这块踩了几个典型的坑:
一、先搞懂WorkflowNotFoundError的核心原因
你遇到的这个错误,本质是你要发信号的工作流实例要么不存在,要么已经不在运行状态,结合你的代码和配置,大概率是这几个情况:
- 工作流实例已经运行结束(处理完所有jobs后退出),Temporal不会给已完成的实例转发信号;
- 固定
workflowId导致调度无法启动新实例,而旧实例可能已经终止; - 工作流代码里的逻辑错误,导致实例提前退出或者启动失败。
二、逐个修复你的代码问题
1. 先修正工作流里的变量错误(最紧急)
你的工作流代码里有个明显的变量名笔误,直接导致activeJobs初始化失败,工作流逻辑完全走偏:
// 原错误代码 for (const req of requests) { const jobs = await createJob(req); // 这里的scrapeJob是未定义的变量!应该是你刚创建的jobs activeJobs[Object.keys(scrapeJob)[0]] = Object.values(scrapeJob)[0]; } // 修正后 for (const req of requests) { const jobs = await createJob(req); const jobId = Object.keys(jobs)[0]; activeJobs[jobId] = false; // 初始标记为未完成 }
这个错误会导致activeJobs没有被正确填充,工作流可能提前退出或者陷入无效循环,直接影响实例的运行状态。
2. 调整调度的workflowId策略
你给调度的工作流设置了固定的workflowId,但Temporal默认不允许同一个workflowId启动多个实例,哪怕你开了ALLOW_ALL的重叠策略也没用——这个策略管的是调度触发时机,不管workflowId重复的问题。
方案A:每个调度实例用唯一workflowId
如果你的需求是每个调度周期处理一批独立的jobs,那给每个实例生成唯一id,同时在创建job时把这个id传给job服务,让webhook回调时带上对应的workflowId:
const schedule = await client.schedule.create({ action: { type: 'startWorkflow', workflowType: jobManagerWorkflow, args: [], taskQueue: 'job-queue', // 用动态唯一id,结合调度id和时间戳 workflowId: `job-manager-{{scheduleId}}-${Date.now()}`, // 可选:允许id复用(仅当实例失败时) workflowIdReusePolicy: WorkflowIdReusePolicy.ALLOW_DUPLICATE_FAILED_ONLY }, // ...其他调度配置 });
然后在工作流里,把当前实例的workflowId传给createJob:
// 工作流里获取自己的workflowId const currentWorkflowId = info().workflowId; for (const req of requests) { // 把workflowId传给job服务,让它回调时带上 const jobs = await createJob(req, currentWorkflowId); const jobId = Object.keys(jobs)[0]; activeJobs[jobId] = false; }
webhook里就可以根据回调传过来的workflowId找到对应的实例。
方案B:用长期运行的工作流接收调度信号
如果你的需求是用同一个工作流持续处理所有调度的jobs,那不要用调度启动工作流,而是先启动一个长期运行的实例,再让调度给它发信号触发job创建:
// 先启动长期运行的工作流 await client.workflow.start(longRunningJobManagerWorkflow, { workflowId: 'long-running-job-manager', taskQueue: 'job-queue' }); // 调度配置改为发信号 const schedule = await client.schedule.create({ action: { type: 'signalWorkflow', workflowId: 'long-running-job-manager', signal: 'createNewJobs', // 工作流里要定义这个信号 args: [requests], // 传递要创建的job请求 taskQueue: 'job-queue', }, // ...其他调度配置 });
对应的长期运行工作流代码:
const signalCreateNewJobs = defineSignal<[Request[]]>('createNewJobs'); const signalJobCompleted = defineSignal<[string]>('signalJobCompleted'); export async function longRunningJobManagerWorkflow(): Promise<void> { const activeJobs: Record<string, boolean> = {}; // 处理调度发来的创建job信号 setHandler(signalCreateNewJobs, async (requests) => { const currentWorkflowId = info().workflowId; for (const req of requests) { const jobs = await createJob(req, currentWorkflowId); const jobId = Object.keys(jobs)[0]; activeJobs[jobId] = false; } }); // 处理job完成信号 setHandler(signalJobCompleted, (jobId: string) => { if (activeJobs[jobId] === false) { activeJobs[jobId] = true; } }); // 长期运行,休眠等待信号或job完成 while (true) { // 处理已完成的job for (const [jobId, isComplete] of Object.entries(activeJobs)) { if (isComplete) { await processJobA(jobId); delete activeJobs[jobId]; } } // 用condition替代轮询,节省资源 await condition(() => Object.values(activeJobs).some(status => status), '1 minute'); } }
3. 替换轮询逻辑为Temporal推荐的condition
你原代码里用while循环轮询是Temporal的反模式,会频繁触发工作流任务,浪费资源。改用condition让工作流休眠,直到满足条件再唤醒:
// 原错误轮询 while (Object.values(activeJobs).some((status) => !status)) { // ...处理逻辑 } // 修正后 await condition(() => Object.values(activeJobs).every(status => status)); // 所有job完成后统一处理 for (const jobId of Object.keys(activeJobs)) { await processJobA(jobId); }
三、验证修复效果
- 打开Temporal Web UI,查看
unique-workflow-id对应的实例状态:- 如果是
Completed:说明工作流已经运行结束,无法接收信号; - 如果是
Failed:查看工作流日志找启动失败原因; - 如果是
Running:在webhook里先调用workflow.describe()确认实例状态,再发信号。
- 如果是
- 测试调度触发逻辑,确保每个实例能正常启动,且webhook能正确找到对应的运行中实例。
备注:内容来源于stack exchange,提问作者Freddy.
相关产品推荐
相关产品推荐

