ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

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

2026/9/10 15:43:09 拓冰建站 浏览量
Medusa 工作流引擎内存实现(workflow-engine-inmemory)深度解析:从 v2.0 到 v2.20 的演进与源码实现 Medusa 工作流引擎内存实现workflow-engine-inmemory深度解析从 v2.0 到 v2.20 的演进与源码实现【免费下载链接】medusaThe worlds most flexible commerce platform for agents and developers项目地址: https://gitcode.com/GitHub_Trending/me/medusamedusajs/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作为 peerDependency2.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.2PR #6128引入 in-memory 与 Redis 两套工作流引擎确立模块骨架。0.0.3PR #6330正式定义「Workflow engine API」即后续稳定下来的run、getRunningTransaction、retryStep、setStepSuccess、setStepFailure、subscribe等编排器方法详见下文源码解析。0.0.4PR #6869修复引擎订阅者subscribers的响应与错误处理为 subscribe.ts 等订阅测试打下基础。2.0.0Medusa 2.0 大版本对齐2.0.0PR #7341将整个模块随 Medusa 2.0 同步发布依赖从medusajs/types、medusajs/workflows-sdk、medusajs/modules-sdk、medusajs/utils多个包收敛为单一的medusajs/framework此后每个版本号的变更几乎都伴随 framework 的同步升级。2.1.1迁移到 DMLPR #10477「Migrate to DML」将数据模型从旧式实体定义迁移到 Medusa 2.x 的声明式模型语言DML。对应源码即 src/models/workflow-execution.tsexport 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联合主键唯一标识一次执行runexecution存整个事务流的检查点TransactionCheckpointstate记录TransactionState枚举如NOT_STARTED、INVOKING、DONE、FAILED、REVERTED、WAITING_TO_COMPENSATE等retention_time用于过期清理。仓库src/migrations/下从Migration20231228143900.ts到Migration20250908080305.ts共 7 个迁移文件反映模型随版本持续演化的痕迹。2.4.0MikroORM 6 与软删除约束PR #10292 将底层 ORM 升级到 MikroORM 6PR #11048 修复「Unique constraint should account for soft deleted records」——即唯一索引需考虑软删除记录这正是上述模型中大量where: deleted_at IS NULL部分索引partial index存在的直接原因。2.5.0 与 2.8.x生命周期与幂等重跑2.5.0PR #11200引入「remove expired workflow executions」即过期执行清理能力的雏形。2.8.0PR #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.0PR #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.0PR #13151fix(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.0PR #14037fix(): Identify step that force save checkpoint——识别强制保存检查点的步骤。对应存储层saveToDb()中的shouldStoreCurrentSteps逻辑当找到当前步骤且同深度的步骤定义中store true时即使事务未结束也会把检查点写入数据库保证关键步骤的状态不会因进程崩溃而丢失。2.17.x 及之后可观测性与测试加速2.17.0PR #15805使用数据库快照/模板加速测试运行。2.17.1PR #15815pass scheduled_for to job handlers——把本次调度的执行时间scheduledFor作为input.scheduledForISO 字符串传给定时工作流的 handler对应jobHandler()中input: { scheduledFor: scheduledFor.toISOString() }的实现。2.17.2PR #15683为包补充bugs元数据。2.20.x仅随 framework 同步发布无行为变更。核心架构与源码实现服务分层模块由三个核心组件协作WorkflowsModuleService对外暴露的模块服务继承ModulesSdkUtils.MedusaService提供run、cancel、retryStep、setStepSuccess、setStepFailure、subscribe、unsubscribe、getRunningTransaction以及基于WorkflowExecution模型的listWorkflowExecutions/listAndCountWorkflowExecutions查询方法支持按q对workflow_id、transaction_id、state、runId做模糊检索。WorkflowOrchestratorService真正的编排逻辑。构造时完成三件关键接线inMemoryDistributedTransactionStorage.setWorkflowOrchestratorService(this) DistributedTransaction.setStorage(inMemoryDistributedTransactionStorage) WorkflowScheduler.setStorage(inMemoryDistributedTransactionStorage)即把「事务存储」与「调度器存储」都指向内存存储实现使框架层的分布式事务与调度器统一走本模块。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维护一个模块级静态MapworkflowId, MaptransactionId, handlers[]传transactionId时只订阅该事务的事件否则订阅该工作流的全部事件键any。notify使用setImmediate异步分发订阅者避免阻塞工作流执行当事件为onFinish时自动删除该事务的订阅者防止内存泄漏。buildWorkflowEvents把框架事务事件的回调统一转发为订阅者通知覆盖onBegin、onResume、onTimeout、onCompensateBegin、onFinish、onStepBegin、onStepSuccess、onStepFailure、onStepAwaiting、onCompensateStepSuccess、onCompensateStepFailure等全套生命周期同时保留调用方传入的自定义事件处理器。内存存储检查点、定时器与竞态防护存储层以Recordstring, 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 worlds most flexible commerce platform for agents and developers项目地址: https://gitcode.com/GitHub_Trending/me/medusa创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考