ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

LobeHub Upstash Workflow 三层模式实战:以 welcome-placeholder 与 agent-welcome 两个真实工作流为例

LobeHub Upstash Workflow 三层模式实战:以 welcome-placeholder 与 agent-welcome 两个真实工作流为例 LobeHub Upstash Workflow 三层模式实战以 welcome-placeholder 与 agent-welcome 两个真实工作流为例【免费下载链接】lobehub LobeHub is your Chief Agent Operator, organizing your agents into 7×24 operations by hiring, scheduling, and reporting on your entire AI team.项目地址: https://gitcode.com/GitHub_Trending/lo/lobehub本文基于 LobeHub 仓库中的技能参考文档 examples.md完整拆解两个已经落地代码库的 Upstash Workflow QStash 异步工作流——welcome-placeholderAI 欢迎占位符生成与agent-welcomeAgent 欢迎语与开放问题生成。读完本文你将掌握 LobeHub 统一的「入口 → 分页 → 单任务」三层工作流架构理解 dry-run 预演、fan-out 扇出、单任务幂等执行三大核心模式如何在真实业务中落地并学会如何以「实体替换」的方式低成本新增自己的工作流。一、模式背景为什么 LobeHub 的工作流都是同一个形状在 LobeHub 中基于 Upstash Workflow 的异步批量任务例如为全量用户生成内容、为全量 Agent 预热数据必须同时应对三个平台层面的约束见 SKILL.md速率限制让无脑全量扇出blind fan-out非常危险步骤上限限制单个 workflow 的执行规模幂等性要求决定了重试时不能重复处理同一个对象。因此每个工作流都由三个核心模式组合而成模式作用出现层级Dry-Run 模式只统计不执行预览「将会处理多少条」Layer 1 入口层Fan-Out 扇出模式把超大批次切成小分片并行递归处理Layer 2 分页层单任务执行每次 workflow 执行恰好处理一个对象Layer 3 执行层所有工作流遵循同一套三层架构Layer 1: 入口层 (process-*) ├─ 校验前置条件 ├─ 计算待处理对象总数 ├─ 过滤已有结果的对象 ├─ 支持 dry-run 模式仅统计 └─ 有待处理工作时触发 Layer 2 Layer 2: 分页层 (paginate-*) ├─ 游标cursor分页 ├─ 大批次 fan-out ├─ 递归处理所有页 └─ 为每个对象触发 Layer 3 Layer 3: 单任务执行层 (execute-* / generate-*) └─ 针对单个对象执行真正的业务逻辑文件布局约定上API 路由放在src/app/(backend)/api/workflows/{workflow-name}/下三个 route 文件分别对应三层而工作流类放在apps/server/src/workflows/{workflowName}/index.ts。当前仓库的 apps/server/src/workflows 目录中可以看到多个按该约定组织的实现如onboardingUnderstanding、topicAutoSummary等印证了这是代码库的长期约定。需要说明的是下文引用的welcome-placeholder、agent-welcome两个示例的路径/api/workflows/welcome-placeholder/...等来自 examples.md 文档的历史快照从当前仓库源码结构看代码库经过重构后部分具体文件已迁移但「三层切分 三大模式」这一模式本身在文档与技能体系中仍被作为标准范式保留。二、示例一welcome-placeholder用户欢迎占位符使用场景为平台用户生成 AI 驱动的欢迎占位符welcome placeholders。三层结构Layer 1process-users— 入口层检查哪些用户符合条件Layer 2paginate-users— 分页遍历所有活跃用户Layer 3generate-user— 为一个用户生成占位符。关键特性这一节是 examples.md 的核心内容完整列出过滤掉 Redis 中已缓存占位符的用户——这是幂等性的第一道闸门paidOnly标志位将处理范围限定为订阅付费用户dryRun模式只输出统计信息不产生任何副作用大批量用户时的 fan-out分片大小CHUNK_SIZE20。Layer 3 的代码形态export const { POST } serveGenerateUserPlaceholderPayload(async (context) { const { userId } context.requestPayload ?? {}; const workflow new WelcomePlaceholderWorkflow(db, userId); const placeholders await context.run(generate, () workflow.generate()); return { success: true, userId, placeholdersCount: placeholders.length }; });注意这里的几个模式要点路由处理器通过upstash/workflow/nextjs的servePayload包裹payload 类型化GenerateUserPlaceholderPayload真正的业务被封装进WelcomePlaceholderWorkflow类构造函数接收db与单个userId路由层只做「取 payload → 调类方法 → 返回统一响应」的薄壳context.run(generate, ...)为这一步赋予唯一步骤名——步骤名是 workflow 持久化与断点恢复的依据SKILL.md 的清单中明确要求「Uniquecontext.run()step names」。文档中给出的对应文件历史路径供溯源/api/workflows/welcome-placeholder/process-users/route.ts /api/workflows/welcome-placeholder/paginate-users/route.ts /api/workflows/welcome-placeholder/generate-user/route.ts /server/workflows/welcomePlaceholder/index.ts三、示例二agent-welcomeAgent 欢迎数据使用场景为 AI Agent 生成欢迎消息与开放问题welcome messages and open questions。三层结构Layer 1process-agents— 入口层检查哪些 Agent 符合条件Layer 2paginate-agents— 分页遍历活跃 AgentLayer 3generate-agent— 为一个Agent 生成欢迎数据。关键特性与示例一完全对应的四件套——过滤 Redis 中已有缓存数据的 AgentpaidOnly标志位只处理订阅用户的 AgentdryRun统计模式;大批量 Agent 时 fan-outCHUNK_SIZE20。Layer 3 的代码形态export const { POST } serveGenerateAgentWelcomePayload(async (context) { const { agentId } context.requestPayload ?? {}; const workflow new AgentWelcomeWorkflow(db, agentId); const data await context.run(generate, () workflow.generate()); return { success: true, agentId, data }; });文档中给出的对应文件历史路径供溯源/api/workflows/agent-welcome/process-agents/route.ts /api/workflows/agent-welcome/paginate-agents/route.ts /api/workflows/agent-welcome/generate-agent/route.ts /server/workflows/agentWelcome/index.ts四、两者的相同与差异「实体替换」才是这套模式的真正卖点examples.md 用一句话点破了两个示例的本质关系它们是同一个模式只在三点上不同实体类型users vs agents业务逻辑placeholder 生成 vs welcome 数据生成数据源不同的数据库查询。除此之外的一切——三层切分、dry-run 处理、fan-out、filter-existing 过滤、flowControl调参——完全一致。文档原话的结论是一旦内化了这个模式新增一个工作流主要就是实体替换entity-substitution。这正是「示例文档」的价值所在读这两个端到端实现等于同时读到了 LobeHub 工作流范式的模板实例。五、纵深三大模式在完整模板中的落点为了把示例中「只展示了 Layer 3」的部分补全references/implementation.md 给出了与上述两个示例逐行同构的完整代码模板其中与示例一、示例二直接对应的实现细节如下。5.1 工作流类静态 trigger 方法 幂等过滤器export class {WorkflowName}Workflow { private static client: Client; static triggerProcessItems(payload: ProcessItemsPayload) { const url getWorkflowUrl(WORKFLOW_PATHS.processItems); return this.getClient().trigger({ body: payload, url }); } static triggerPaginateItems(payload: PaginateItemsPayload) { const url getWorkflowUrl(WORKFLOW_PATHS.paginateItems); return this.getClient().trigger({ body: payload, url }); } static triggerExecuteItem(payload: ExecuteItemPayload) { const url getWorkflowUrl(WORKFLOW_PATHS.executeItem); return this.getClient().trigger({ body: payload, url }); } /** * Filter items that need processing (e.g. check Redis cache, database state). * Return only the ones that actually need work — keeps the pipeline idempotent. */ static async filterItemsNeedingProcessing(itemIds: string[]): Promisestring[] { // ... } }对应到两个示例WelcomePlaceholderWorkflow/AgentWelcomeWorkflow的实例方法generate()就是 Layer 3 调用的业务入口而静态trigger*方法承担了「workflow 之间通过 QStash HTTP 触发」的桥接——每个静态方法都对应一个 HTTP 路由 URL。5.2 Layer 1dry-run 短路与并行度锁定入口层模板中dry-run 在一切副作用之前短路返回统计信息if (dryRun) { console.log([{workflow}:process] Dry run mode, returning statistics only); return { ...result, dryRun: true, message: [DryRun] Would process ${itemsNeedingProcessing.length} items, }; }这解释了示例中「dryRunmode for statistics」的实际机制先get-all-items查询全量符合条件的实体 id再经filter-existing剔除已处理项然后先判断 dryRun 再决定是否触发分页。入口层的 flowControl 被刻意压到parallelism: 1, ratePerSecond: 1从结构上看是为了保证「同一时刻只有一个入口实例在派发」避免重复触发。5.3 Layer 2fan-out 的递归自触发示例中标注的CHUNK_SIZE20在模板中对应这样的分支逻辑const PAGE_SIZE 50; // 每页拉取的对象数 const CHUNK_SIZE 20; // 每个 fan-out 分片的对象数 if (itemIds.length CHUNK_SIZE) { const chunks chunk(itemIds, CHUNK_SIZE); await Promise.all( chunks.map((ids, idx) context.run({workflow}:fanout:${idx 1}/${chunks.length}, () WorkflowClass.triggerPaginateItems({ itemIds: ids }), ), ), ); } else { // 直接为每个对象触发 Layer 3 await Promise.all( itemIds.map((itemId) context.run({workflow}:execute:${itemId}, () WorkflowClass.triggerExecuteItem({ itemId }), ), ), ); }两个值得注意的机制fan-out 是递归自触发分片不是在本 workflow 内继续循环而是带着itemIds重新trigger一次 paginate 路由天然绕开了单个 workflow 的步骤数上限双模式 payloadPaginateItemsPayload同时支持cursor常规分页与itemIdsfan-out 分片直入两种入参收到itemIds时直接走执行分支不再翻页。分页层的 flowControl 默认parallelism: 20, ratePerSecond: 5与执行层的parallelism: 10, ratePerSecond: 5相配合形成「入口串行 → 分页半并行 → 执行受控并行」的速率漏斗对应 SKILL.md 中「速率限制让盲目 fan-out 危险」的约束。5.4 环境变量与触发前提三个trigger*方法在运行时依赖两个必需环境变量缺失时直接抛错# Required for all workflows APP_URLhttps://your-app.com # Base URL for workflow endpoints QSTASH_TOKENqstash_xxx # QStash authentication token # Optional (for custom QStash URL) QSTASH_URLhttps://custom-qstash.com即APP_URL用于拼接/api/workflows/...的完整触发地址QSTASH_TOKEN用于构造 UpstashClient的鉴权QSTASH_URL可选用于指向自建/自定义的 QStash 端点。六、复用到你自己的实体从示例到新工作流的最小路径把示例一/二抽象后在 LobeHub 代码库中新增一个「对全量实体做 X」的工作流路径基本是填空确定实体与业务逻辑明确处理对象users / agents / items…与单对象业务函数写工作流类在apps/server/src/workflows/{workflowName}/index.ts定义WORKFLOW_PATHS、三种 payload 接口、三个静态 trigger 方法与filterItemsNeedingProcessingRedis 缓存或数据库状态判断写三层路由process-*查询全量候选 →filter-existing→ dry-run 短路或触发分页flowControlparallelism: 1paginate-*PAGE_SIZE50游标分页CHUNK_SIZE20超阈 fan-outflowControlparallelism: 20, ratePerSecond: 5generate-*/execute-*单对象执行context.run步骤名唯一flowControlparallelism: 10, ratePerSecond: 5日志与响应统一[workflow:layer]前缀日志、统一{ success, ... }响应形状;验证顺序先dryRun冒烟再用小批量试跑最后全量。SKILL.md 中还维护了一份覆盖「规划 / 实现 / 质量与部署」三段的新工作流清单可作为发布前的逐项核对flowControl 调参、错误处理、日志与测试的进一步细节见 best-practices.mdlobehub-cloud 环境下的 re-export 与云上运维见 cloud.md。七、小结examples.md 用welcome-placeholder与agent-welcome两个真实工作流证明LobeHub 的 Upstash 异步工作流是一个模式、多次实例化——三层切分、dry-run、fan-outCHUNK_SIZE20、filter-existing、flowControl 调参全部同构差异仅在于实体、业务逻辑与数据源Layer 3 的servePayloadcontext.run 工作流类薄壳调用是可直接套用的代码形态从示例到落地implementation.md 提供逐行同构的完整模板配合APP_URL/QSTASH_TOKEN环境变量即可在本地或云端触发掌握「实体替换」这一视角后新增一个面向任意实体的批量 AI 工作流主要就是替换查询、过滤与生成函数三处代码。【免费下载链接】lobehub LobeHub is your Chief Agent Operator, organizing your agents into 7×24 operations by hiring, scheduling, and reporting on your entire AI team.项目地址: https://gitcode.com/GitHub_Trending/lo/lobehub创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表