☰
Medusa 工作流引擎内存实现(workflow-engine-inmemory)深度解析:从 v2.0 到 v2.20 的演进与源码实现
2026/10/8 22:27:54 网站建设 项目流程

Medusa 工作流引擎内存实现(workflow-engine-inmemory)深度解析:从 v2.0 到 v2.20 的演进与源码实现

【免费下载链接】medusaThe world's most flexible commerce platform for agents and developers项目地址: https://gitcode.com/GitHub_Trending/me/medusa

@medusajs/workflow-engine-inmemory是 Medusa 2.x 中负责执行业务工作流(Workflow)的模块化编排引擎的内存版实现。本文以其 CHANGELOG.md 记录的版本演变为骨架,结合模块源码与集成测试,梳理该引擎从诞生到 v2.20 的关键能力(分布式事务存储、幂等与竞态防护、重试/超时/定时调度、订阅通知、过期清理),并给出可直接对照源码的实践指引。读完本文,你将掌握内存版工作流引擎的内部架构、核心服务方法语义以及其与 Redis 版的取舍边界。

模块定位:Medusa 工作流编排器的内存实现

在 Medusa 2.x 的模块化架构中,工作流(Workflow)由@medusajs/framework/orchestration的分布式事务(Distributed Transaction)机制驱动,而执行状态的存取、重试/超时定时器与定时任务的调度,则由工作流引擎模块负责。仓库中提供两套实现:

  • packages/modules/workflow-engine-inmemory(本文主题):状态以进程内内存为主、同时可落库为快照,适合单实例与本地开发。
  • packages/modules/workflow-engine-redis:基于 Redis 的分布式实现,支持多实例横向扩展。

CHANGELOG 的开篇记录(0.0.2 版本,PR #6128「Modules: Workflows Engine in-memory and Redis」)表明两个引擎在最初是成对设计与发布的,共享同一套编排器接口语义。

包结构与依赖

从 package.json 可以看到该包的核心信息:

  • 包名@medusajs/workflow-engine-inmemory,当前版本 2.20.1,描述为 "Medusa Workflow Orchestrator module";
  • 运行时依赖仅两个:cron-parser(解析定时任务的 cron 表达式)与ulid(生成事务 ID 等标识);
  • @medusajs/framework作为 peerDependency(2.20.1),符合 Medusa 2.x「单一 framework 包集中导出所有 SDK」的依赖策略——这也对应 CHANGELOG 2.11.0 中「Move peer deps into a single package and re export from framework」的清理工作;
  • Node 要求>=20。

模块的入口 src/index.ts 只有寥寥几行,它通过Module(Modules.WORKFLOW_ENGINE, ...)向框架注册服务与加载器:

import { Module, Modules } from "@medusajs/framework/utils" import { WorkflowsModuleService } from "@services" import { loadUtils } from "./loaders" export default Module(Modules.WORKFLOW_ENGINE, { service: WorkflowsModuleService, loaders: [loadUtils], })

从 CHANGELOG 看引擎演进主线

虽然 CHANGELOG 以依赖更新记录为主,但它忠实地记录了引擎能力的关键节点,可与源码一一对应。

0.0.x:诞生与 API 成型

  • 0.0.2(PR #6128):引入 in-memory 与 Redis 两套工作流引擎,确立模块骨架。
  • 0.0.3(PR #6330):正式定义「Workflow engine API」,即后续稳定下来的run、getRunningTransaction、retryStep、setStepSuccess、setStepFailure、subscribe等编排器方法(详见下文源码解析)。
  • 0.0.4(PR #6869):修复引擎订阅者(subscribers)的响应与错误处理,为 subscribe.ts 等订阅测试打下基础。

2.0.0:Medusa 2.0 大版本对齐

2.0.0(PR #7341)将整个模块随 Medusa 2.0 同步发布,依赖从@medusajs/types、@medusajs/workflows-sdk、@medusajs/modules-sdk、@medusajs/utils多个包收敛为单一的@medusajs/framework,此后每个版本号的变更几乎都伴随 framework 的同步升级。

2.1.1:迁移到 DML

PR #10477「Migrate to DML」将数据模型从旧式实体定义迁移到 Medusa 2.x 的声明式模型语言(DML)。对应源码即 src/models/workflow-execution.ts:

export const WorkflowExecution = model .define("workflow_execution", { id: model.id({ prefix: "wf_exec" }), workflow_id: model.text().primaryKey(), transaction_id: model.text().primaryKey(), run_id: model.text().primaryKey(), execution: model.json().nullable(), context: model.json().nullable(), state: model.enum(TransactionState), retention_time: model.number().nullable(), }) .indexes([ ... ])

模型以workflow_id + transaction_id + run_id联合主键唯一标识一次执行(run),execution存整个事务流的检查点(TransactionCheckpoint),state记录TransactionState枚举(如NOT_STARTED、INVOKING、DONE、FAILED、REVERTED、WAITING_TO_COMPENSATE等),retention_time用于过期清理。仓库src/migrations/下从Migration20231228143900.ts到Migration20250908080305.ts共 7 个迁移文件,反映模型随版本持续演化的痕迹。

2.4.0:MikroORM 6 与软删除约束

PR #10292 将底层 ORM 升级到 MikroORM 6;PR #11048 修复「Unique constraint should account for soft deleted records」——即唯一索引需考虑软删除记录,这正是上述模型中大量where: "deleted_at IS NULL"部分索引(partial index)存在的直接原因。

2.5.0 与 2.8.x:生命周期与幂等重跑

  • 2.5.0(PR #11200):引入「remove expired workflow executions」,即过期执行清理能力的雏形。
  • 2.8.0(PR #12362):允许「re run non idempotent but stored workflow with the same transaction id if considered done」——当一个非幂等但已存储的工作流被视为已完成(done)时,允许用同一事务 ID 重新运行。源码中get()方法对idempotent选项与非幂等场景下DONE/FAILED/REVERTED终态的过滤逻辑(workflow-orchestrator-storage.ts)即为此服务。
  • 2.7.0(PR #11873):防止共享 context 引用,并暴露cancel方法(PR #11844「expose cancel method」)。对应源码 workflow-orchestrator.ts 中的cancel():它通过getRunningTransaction找到运行中的事务,调用exportedWorkflow.cancel触发补偿流程,并返回包含transactionId、hasFinished、hasFailed、exists等字段的acknowledgement。

2.9.0 ~ 2.12.x:调度、竞态与检查点稳定化

  • 2.9.0(PR #13151):fix(workflow-engine-inmemory): fix cron job schedule——修复定时任务调度。对应InMemoryDistributedTransactionStorage.schedule():先用cron-parser解析表达式并计算到下次执行的延迟,setTimeout触发jobHandler,且调用timer.unref()避免定时器阻止进程退出。
  • 2.10.x:引入autoRetry步骤配置支持(PR #13391);改进工作流引擎的竞态条件(PR #13345);仅对异步工作流执行竞态检查(PR #13396);改进定时器与通知(PR #13434)。
  • 2.11.0:一批并发相关修复——「workflows concurrency」「workflow async concurrency」(PR #13645、#13769)、「Always create cleaner job」(即周期性清理任务的定时器恒定创建)、「Workflow save to db + index integration instability」以及工作流引擎迁移问题修复。
  • 2.12.0(PR #14037):fix(): Identify step that force save checkpoint——识别强制保存检查点的步骤。对应存储层saveToDb()中的shouldStoreCurrentSteps逻辑:当找到当前步骤且同深度的步骤定义中store === true时,即使事务未结束也会把检查点写入数据库,保证关键步骤的状态不会因进程崩溃而丢失。

2.17.x 及之后:可观测性与测试加速

  • 2.17.0(PR #15805):使用数据库快照/模板加速测试运行。
  • 2.17.1(PR #15815):pass scheduled_for to job handlers——把本次调度的执行时间scheduledFor作为input.scheduledFor(ISO 字符串)传给定时工作流的 handler,对应jobHandler()中input: { scheduledFor: scheduledFor.toISOString() }的实现。
  • 2.17.2(PR #15683):为包补充bugs元数据。
  • 2.20.x:仅随 framework 同步发布,无行为变更。

核心架构与源码实现

服务分层

模块由三个核心组件协作:

  1. WorkflowsModuleService:对外暴露的模块服务,继承ModulesSdkUtils.MedusaService,提供run、cancel、retryStep、setStepSuccess、setStepFailure、subscribe、unsubscribe、getRunningTransaction,以及基于WorkflowExecution模型的listWorkflowExecutions/listAndCountWorkflowExecutions查询方法(支持按q对workflow_id、transaction_id、state、runId做模糊检索)。
  2. WorkflowOrchestratorService:真正的编排逻辑。构造时完成三件关键接线:
    inMemoryDistributedTransactionStorage.setWorkflowOrchestratorService(this) DistributedTransaction.setStorage(inMemoryDistributedTransactionStorage) WorkflowScheduler.setStorage(inMemoryDistributedTransactionStorage)

    即把「事务存储」与「调度器存储」都指向内存存储实现,使框架层的分布式事务与调度器统一走本模块。

  3. InMemoryDistributedTransactionStorage:同时实现IDistributedTransactionStorage与IDistributedSchedulerStorage两个接口的存储层。

编排器方法语义

run()是核心入口:接受工作流 ID(字符串)或工作流对象,自动生成transactionId(未指定时为"auto-" + ulid()),通过MedusaWorkflow.getWorkflow(workflowId)找到已注册工作流并执行,最后返回包含acknowledgement(含transactionId、workflowId、hasFinished、hasFailed等)与执行结果的对象;若throwOnError(默认 true)且存在错误则抛出首个错误。

setStepSuccess/setStepFailure/retryStep均以idempotencyKey定位步骤。buildIdempotencyKeyAndParts支持两种形态:结构化对象{ workflowId, transactionId, stepId, action }或形如workflowId:transactionId:stepId:action的拼接字符串。setStepFailure还支持forcePermanentFailure强制标记永久失败。

subscribe/unsubscribe维护一个模块级静态Map<workflowId, Map<transactionId, handlers[]>>:传transactionId时只订阅该事务的事件,否则订阅该工作流的全部事件(键"any")。notify使用setImmediate异步分发订阅者,避免阻塞工作流执行;当事件为onFinish时自动删除该事务的订阅者,防止内存泄漏。

buildWorkflowEvents把框架事务事件的回调统一转发为订阅者通知,覆盖onBegin、onResume、onTimeout、onCompensateBegin、onFinish、onStepBegin、onStepSuccess、onStepFailure、onStepAwaiting、onCompensateStepSuccess、onCompensateStepFailure等全套生命周期,同时保留调用方传入的自定义事件处理器。

内存存储:检查点、定时器与竞态防护

存储层以Record<string, TransactionCheckpoint>作为进程内主存储,同时按策略写库:

  • 写库策略(saveToDb):仅当事务未开始、已结束、等待补偿、存在强制保存步骤或存在异步版本(flow._v)时才落库,其余中间态仅驻留内存,以换取性能。
  • 终态优化:事务达到DONE/FAILED/REVERTED时,若无retentionTime且无父步骤幂等键(非子工作流),直接删除 DB 记录;否则保留并写入retention_time。
  • 竞态防护(#preventRaceConditionExecutionIfNecessary):通过检查已存储检查点判断并发执行状态——若某步骤已被其他执行完成则抛SkipStepAlreadyFinishedError;若非初始检查点但已找不到检查点(说明另一执行已结束并清理),抛SkipExecutionError;若最新执行已被取消而当前未取消,抛SkipCancelledExecutionError。run()内部对SkipExecutionError族异常做静默吞掉处理,实现「重复执行时后到者让位」的语义。
  • 重试与超时:scheduleRetry、scheduleStepTimeout、scheduleTransactionTimeout均以Map维护 Node 定时器(键为workflowId:transactionId[:stepId]),超时后通过executeTransaction重新拉起编排器run以推进流程;clearRetry/clearStepTimeout/clearTransactionTimeout负责在步骤完成时取消定时器。
  • 生命周期:onApplicationStart启动一个 30 分钟(THIRTY_MINUTES_IN_MS)周期的清理定时器调用clearExpiredExecutions,按retention_time与updated_at删除已进入终态的超期记录;onApplicationShutdown统一清理重试、超时、调度与挂起定时器。

定时任务调度

schedule()接受 cron 表达式(如"0 0 * * *")或固定interval(毫秒数),先remove旧任务保证配置始终最新,再用cron-parser解析并计算首次延迟。jobHandler在每次触发时:检查numberOfExecutions上限 → 调用编排器run(jobId, { input: { scheduledFor } })→ 安排下一次执行;若工作流不存在(NOT_FOUND),记录警告并移除该任务。cron-parser正是 package.json 中声明的运行时依赖。

测试验证:仓库内的完整证据链

模块在 integration-tests 下提供了体系化的测试夹具与用例,是理解引擎行为的绝佳入口:

  • 工作流夹具(__fixtures__/):覆盖同步/异步(workflow_sync.ts、workflow_async.ts、workflow_parallel_async.ts)、幂等/非幂等(workflow_idempotent.ts、workflow_not_idempotent_with_retention.ts)、条件步骤、事件分组(workflow_event_group_id.ts)、定时任务(workflow_scheduled.ts)、步骤超时(workflow_step_timeout.ts)、事务超时(workflow_transaction_timeout.ts)、自动重试(workflow_1_auto_retries.ts及其关闭版)、手动重试(workflow_1_manual_retry_step.ts)、重试间隔(workflow_retry_interval.ts、workflow_sync_retry_interval.ts)等场景。
  • 测试用例(__tests__/):index.spec.ts覆盖主流程,race.spec.ts验证并发竞态防护,subscribe.spec.ts验证订阅通知机制,retry-interval.spec.ts验证重试间隔语义。这些用例与 CHANGELOG 中 2.10.x(竞态改进、autoRetry)、2.11.0(并发修复)等条目一一呼应。

适用场景与边界

综合 README.md(仅一句 "Workflow Orchestrator")与源码结构,可以给出如下谨慎结论:

  • 适用:单实例部署、本地开发、测试环境,或对持久化要求不高的异步工作流场景。检查点数据仍会按策略写入数据库(如workflow_execution表),从而支持进程重启后的执行恢复。
  • 边界:内存存储与进程内定时器意味着工作流状态、重试/超时/定时任务不跨实例共享;多实例横向扩展时应选用workflow-engine-redis。CHANGELOG 中反复出现的「竞态条件」「并发」修复条目,正说明这套内存实现必须在并发执行语义上保持与 Redis 版一致的行为契约。

参考路径速查

  • 版本变更记录:CHANGELOG.md
  • 模块注册入口:src/index.ts
  • 对外模块服务:src/services/workflows-module.ts
  • 编排器服务:src/services/workflow-orchestrator.ts
  • 内存存储与调度实现:src/utils/workflow-orchestrator-storage.ts
  • 数据模型:src/models/workflow-execution.ts
  • 集成测试:integration-tests/tests与 integration-tests/fixtures

【免费下载链接】medusaThe world's most flexible commerce platform for agents and developers项目地址: https://gitcode.com/GitHub_Trending/me/medusa

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询