首页
/ LobeHub Upstash Workflow 实战指南:基于 QStash 的三层异步工作流架构与三大核心模式

LobeHub Upstash Workflow 实战指南:基于 QStash 的三层异步工作流架构与三大核心模式

2026-09-06 13:56:46作者:翟江哲Frasier

LobeHub 使用 Upstash Workflow + QStash 在后台执行大量异步批量任务(评测运行、记忆提取、onboarding 处理等),并沉淀出一套标准化的"三层架构 + 三大模式"实现规范。读完本文,你将掌握 dry-run 预演、fan-out 分片扇出、单任务执行这三种模式的设计动机与完整代码模板,并能按 LobeHub 仓库中 agent-eval-run 等真实工作流的写法,独立开发一个具备幂等性、限流控制和可观测性的新工作流。

一、为什么是"三大模式":平台约束决定了设计

LobeHub 的每个 Upstash Workflow 都组合了三个核心模式。它们之所以存在,是因为平台在三个维度上约束了你:速率限制(rate limit)让盲目 fan-out 很危险,步骤上限(step limit)限制了单个工作流的规模,而幂等性(idempotency)要求重试不能造成重复处理。

  1. Dry-Run Mode(预演模式)——不触发实际执行,先拿到统计信息;
  2. Fan-Out Pattern(扇出模式)——把大批量拆成小分片并行处理;
  3. Single Task Execution(单任务执行)——每次工作流执行只处理一条数据。

完整的实现规范定义在仓库的 SKILL.md 中,配套的四个参考文档分别为 实现模板最佳实践真实案例云项目部署

二、三层架构总览

所有工作流遵循同一套 3 层架构:

Layer 1: Entry Point (process-*)
  ├─ 校验前置条件
  ├─ 计算待处理数据总量
  ├─ 过滤已处理数据
  ├─ 支持 dry-run 模式(仅返回统计)
  └─ 有工作要处理时触发 Layer 2

Layer 2: Pagination (paginate-*)
  ├─ 处理基于游标(cursor)的分页
  ├─ 对大批量实现 fan-out
  ├─ 递归处理所有页
  └─ 为每条数据触发 Layer 3

Layer 3: Single Task Execution (execute-* / generate-*)
  └─ 执行单条数据的实际业务逻辑

仓库中的实际落地位置

从源码结构看,当前仓库的落地方式与规范文档略有演进:规范文档描述的是 Next.js App Router 的逐目录路由形态(src/app/(backend)/api/workflows/{name}/process-*/route.ts),而当前主服务已迁移到 Hono 服务端的集中式路由。对应关系如下:

  • 统一入口:工作流捕获路由/api/workflows/%5B%5B...route%5D%5D/route.ts) 把所有 POST /api/workflows/* 请求委托给 Hono 应用;
  • 路由注册workflows Hono 根路由new Hono().basePath('/api/workflows') 挂载了 agent-eval-runmemory-user-memoryagent-signaltaskgoaltopic-auto-summaryverify 等九个业务工作流;
  • Workflow 类:位于 apps/server/src/workflows/ 下,如 agentEvalRun
  • 各层 Handler:位于 apps/server/src/router-hono/workflows/{name}/workflows/ 下,如 runBenchmarkpaginateTestCases

规范中给出的标准文件结构如下(新工作流在开源仓中按此组织):

src/
├── app/(backend)/api/workflows/
│   └── {workflow-name}/
│       ├── process-{entities}/route.ts      # Layer 1
│       ├── paginate-{entities}/route.ts     # Layer 2
│       └── execute-{entity}/route.ts        # Layer 3
└── server/workflows/
    └── {workflowName}/
        └── index.ts                          # Workflow 类

规范文档推荐的延伸阅读导航(均已转换为本仓库路径):

你想做的事 参考文档
从零写 Workflow 类 + 三层路由 references/implementation.md
调优 flowControl、错误处理、日志与测试 references/best-practices.md
看两个端到端的真实工作流 references/examples.md
在 lobehub-cloud 上部署(re-export、云专属操作) references/cloud.md

三、三大模式详解(60 秒版)

1. Dry-Run Mode:在任何副作用之前短路

Layer 1 在产生任何副作用之前短路返回,让调用方能"预演"将要发生什么:

if (dryRun) {
  return {
    ...result,
    dryRun: true,
    message: `[DryRun] Would process ${itemsNeedingProcessing.length} items`,
  };
}

典型用途:在真正提交批量任务之前,先确认将要处理多少条数据。仓库中的 runBenchmark 就是这个写法的直接实例:

// If dryRun mode, return statistics only
if (dryRun) {
  console.info('[run-benchmark] Dry run: %d test cases would execute', testCaseIds.length);
  return {
    ...result,
    dryRun: true,
    message: `[DryRun] Would execute ${testCaseIds.length} test cases`,
  };
}

2. Fan-Out Pattern:拆片扇出,绕开步骤上限

Layer 2 把超大批量切分成块(chunk),并用每个分片递归重新触发自身。这样做是为了在一页数据条数过多时避免触碰工作流的步骤上限:

const CHUNK_SIZE = 20;

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 }),
      ),
    ),
  );
}

默认常量:PAGE_SIZE = 50(每页条数)、CHUNK_SIZE = 20(每个扇出分片的条数)。这两个常量在当前仓库的 paginateTestCases 中一字不差地保留:

const CHUNK_SIZE = 20; // Max items to process directly
const PAGE_SIZE = 50; // Items per page

当 fan-out 分片直接带着 itemIds 回到 Layer 2 时,会跳过分页逻辑直接逐条触发 Layer 3(见 paginateTestCases L42-L57)。

3. Single Task Execution:一次执行只处理一条

Layer 3 每次调用恒定处理一条数据。并行度由 Layer 2 扇出到大量 Layer 3 调用来实现,并由 flowControl 控制节奏:

export const { POST } = serve<ExecutePayload>(
  async (context) => {
    const { itemId } = context.requestPayload ?? {};
    if (!itemId) return { success: false, error: 'Missing itemId' };

    const item = await context.run('workflow:get-item', () => getItem(itemId));
    const result = await context.run('workflow:execute', () => processItem(item));
    await context.run('workflow:save', () => saveResult(itemId, result));

    return { success: true, itemId, result };
  },
  {
    flowControl: { key: 'workflow.execute', parallelism: 10, ratePerSecond: 5 },
  },
);

四、Workflow 类:触发器与过滤器的统一封装

规范模板(详见 implementation.md)要求每个工作流配一个 Workflow 类,包含三部分:WORKFLOW_PATHS 常量、各层 payload 类型、静态 trigger 方法,外加一个用于保持管道幂等的静态过滤方法:

import { Client } from '@upstash/workflow';
import debug from 'debug';

const log = debug('lobe-server:workflows:{workflow-name}');

// Workflow paths
const WORKFLOW_PATHS = {
  processItems: '/api/workflows/{workflow-name}/process-items',
  paginateItems: '/api/workflows/{workflow-name}/paginate-items',
  executeItem: '/api/workflows/{workflow-name}/execute-item',
} as const;

// Payload types
export interface ProcessItemsPayload {
  dryRun?: boolean;
  force?: boolean;
}

export interface PaginateItemsPayload {
  cursor?: string;
  itemIds?: string[]; // For fanout chunks
}

export interface ExecuteItemPayload {
  itemId: string;
}

const getWorkflowUrl = (path: string): string => {
  const baseUrl = process.env.APP_URL;
  if (!baseUrl) throw new Error('APP_URL is required to trigger workflows');
  return new URL(path, baseUrl).toString();
};

const getWorkflowClient = (): Client => {
  const token = process.env.QSTASH_TOKEN;
  if (!token) throw new Error('QSTASH_TOKEN is required to trigger workflows');

  const config: ConstructorParameters<typeof Client>[0] = { token };
  if (process.env.QSTASH_URL) {
    (config as Record<string, unknown>).url = process.env.QSTASH_URL;
  }
  return new Client(config);
};

export class {WorkflowName}Workflow {
  private static client: Client;

  private static getClient(): Client {
    if (!this.client) this.client = getWorkflowClient();
    return this.client;
  }

  static triggerProcessItems(payload: ProcessItemsPayload) {
    const url = getWorkflowUrl(WORKFLOW_PATHS.processItems);
    log('Triggering process-items workflow');
    return this.getClient().trigger({ body: payload, url });
  }

  static triggerPaginateItems(payload: PaginateItemsPayload) {
    const url = getWorkflowUrl(WORKFLOW_PATHS.paginateItems);
    log('Triggering paginate-items workflow');
    return this.getClient().trigger({ body: payload, url });
  }

  static triggerExecuteItem(payload: ExecuteItemPayload) {
    const url = getWorkflowUrl(WORKFLOW_PATHS.executeItem);
    log('Triggering execute-item workflow: %s', payload.itemId);
    return this.getClient().trigger({ body: payload, url });
  }

  /**
   * 过滤出仍需要处理的数据(如检查 Redis 缓存、数据库状态)。
   * 只返回真正需要处理的那些 —— 让整条管道保持幂等。
   */
  static async filterItemsNeedingProcessing(itemIds: string[]): Promise<string[]> {
    if (itemIds.length === 0) return [];
    // 检查已有状态,返回需要处理的数据
    return itemIds;
  }
}

对照真实实现 AgentEvalRunWorkflow,可以看到两点演进:

  • 触发客户端不再在类内部惰性构造,而是统一复用全局的 workflowClient(来自 qstash 客户端封装,详见第七节);
  • filterTestCasesNeedingExecution 给出幂等过滤的完整示范:查询该 run 下所有 RunTopic,只保留 status === 'pending' 的用例 ID,其余一律跳过——这就是"重试不会重复处理"的落点。

五、三层路由完整模板

Layer 1:入口点(process-*)

职责:校验前置条件、计算统计信息、支持 dry-run 模式。

import { serve } from '@upstash/workflow/nextjs';
import { getServerDB } from '@/database/server';
import { WorkflowClass, type ProcessPayload } from '@/server/workflows/{workflowName}';

export const { POST } = serve<ProcessPayload>(
  async (context) => {
    const { dryRun, force } = context.requestPayload ?? {};

    console.log('[{workflow}:process] Starting with payload:', { dryRun, force });

    const allItemIds = await context.run('{workflow}:get-all-items', async () => {
      const db = await getServerDB();
      // 查询数据库中符合条件的项目
      return items.map((item) => item.id);
    });

    console.log('[{workflow}:process] Total eligible items:', allItemIds.length);

    if (allItemIds.length === 0) {
      return { success: true, totalEligible: 0, message: 'No eligible items found' };
    }

    const itemsNeedingProcessing = await context.run('{workflow}:filter-existing', () =>
      WorkflowClass.filterItemsNeedingProcessing(allItemIds),
    );

    const result = {
      success: true,
      totalEligible: allItemIds.length,
      toProcess: itemsNeedingProcessing.length,
      alreadyProcessed: allItemIds.length - itemsNeedingProcessing.length,
    };

    // 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`,
      };
    }

    if (itemsNeedingProcessing.length === 0) {
      return { ...result, message: 'All items already processed' };
    }

    await context.run('{workflow}:trigger-paginate', () => WorkflowClass.triggerPaginateItems({}));

    return {
      ...result,
      message: `Triggered pagination for ${itemsNeedingProcessing.length} items`,
    };
  },
  {
    flowControl: {
      key: '{workflow}.process',
      parallelism: 1, // 单实例 —— 避免重复处理
      ratePerSecond: 1,
    },
  },
);

当前仓库中 runBenchmarkHandler 的执行顺序与模板完全一致:校验 runId/userId → 检查 run 状态(running 且未 force 时拒绝)→ 拉取全部测试用例 → filter-existing 过滤 → dry-run 短路 → 无待执行项提前返回 → 更新 run 状态为 running → 触发 paginate 层。

Layer 2:分页(paginate-*)

职责:处理基于游标的分页,对大批量实现 fan-out。

import { serve } from '@upstash/workflow/nextjs';
import { chunk } from 'es-toolkit/compat';
import { getServerDB } from '@/database/server';
import { WorkflowClass, type PaginatePayload } from '@/server/workflows/{workflowName}';

const PAGE_SIZE = 50;
const CHUNK_SIZE = 20;

export const { POST } = serve<PaginatePayload>(
  async (context) => {
    const { cursor, itemIds: payloadItemIds } = context.requestPayload ?? {};

    console.log('[{workflow}:paginate] Starting:', {
      cursor,
      itemIdsCount: payloadItemIds?.length ?? 0,
    });

    // 若传入了特定 itemIds(来自 fanout 分片),直接处理
    if (payloadItemIds && payloadItemIds.length > 0) {
      await Promise.all(
        payloadItemIds.map((itemId) =>
          context.run(`{workflow}:execute:${itemId}`, () =>
            WorkflowClass.triggerExecuteItem({ itemId }),
          ),
        ),
      );
      return { success: true, processedItems: payloadItemIds.length };
    }

    // 逐页遍历所有数据
    const itemBatch = await context.run('{workflow}:get-batch', async () => {
      const db = await getServerDB();
      const items = await db.query(...);
      if (!items.length) return { ids: [] };
      const last = items.at(-1);
      return {
        ids: items.map((item) => item.id),
        cursor: last ? last.id : undefined,
      };
    });

    const batchItemIds = itemBatch.ids;
    const nextCursor = 'cursor' in itemBatch ? itemBatch.cursor : undefined;

    if (batchItemIds.length === 0) {
      return { success: true, message: 'Pagination complete' };
    }

    const itemIds = await context.run('{workflow}:filter-existing', () =>
      WorkflowClass.filterItemsNeedingProcessing(batchItemIds),
    );

    if (itemIds.length > 0) {
      if (itemIds.length > CHUNK_SIZE) {
        // Fan out —— 用每个分片递归重入分页层
        const chunks = chunk(itemIds, CHUNK_SIZE);
        console.log('[{workflow}:paginate] Fanout mode:', {
          chunks: chunks.length,
          chunkSize: CHUNK_SIZE,
        });

        await Promise.all(
          chunks.map((ids, idx) =>
            context.run(`{workflow}:fanout:${idx + 1}/${chunks.length}`, () =>
              WorkflowClass.triggerPaginateItems({ itemIds: ids }),
            ),
          ),
        );
      } else {
        // 直接处理这一页
        await Promise.all(
          itemIds.map((itemId) =>
            context.run(`{workflow}:execute:${itemId}`, () =>
              WorkflowClass.triggerExecuteItem({ itemId }),
            ),
          ),
        );
      }
    }

    // 尾调用进入下一页
    if (nextCursor) {
      await context.run('{workflow}:next-page', () =>
        WorkflowClass.triggerPaginateItems({ cursor: nextCursor }),
      );
    }

    return {
      success: true,
      processedItems: itemIds.length,
      skippedItems: batchItemIds.length - itemIds.length,
      nextCursor: nextCursor ?? null,
    };
  },
  {
    flowControl: {
      key: '{workflow}.paginate',
      parallelism: 20,
      ratePerSecond: 5,
    },
  },
);

实际的 paginateTestCasesHandler 还示范了一个规范未展开但生产中必要的细节:分页前先检查 run 是否被中止runStatus === 'aborted' 时返回 { cancelled: true }),避免用户取消后分页层仍在盲目触发下游任务。

Layer 3:执行(execute-* / generate-*)

职责:对恰好一条数据执行实际业务逻辑。

import { serve } from '@upstash/workflow/nextjs';
import { getServerDB } from '@/database/server';
import { WorkflowClass, type ExecutePayload } from '@/server/workflows/{workflowName}';

export const { POST } = serve<ExecutePayload>(
  async (context) => {
    const { itemId } = context.requestPayload ?? {};

    if (!itemId) {
      return { success: false, error: 'Missing itemId' };
    }

    const db = await getServerDB();

    const item = await context.run('{workflow}:get-item', async () => {
      // 查询数据库获取该条目
      return item;
    });

    if (!item) {
      return { success: false, error: 'Item not found' };
    }

    const result = await context.run('{workflow}:process-item', async () => {
      const workflow = new WorkflowClass(db, itemId);
      return workflow.generate(); // 或 process()、execute() 等
    });

    await context.run('{workflow}:save-result', async () => {
      const workflow = new WorkflowClass(db, itemId);
      return workflow.saveToRedis(result); // 或 saveToDatabase() 等
    });

    return { success: true, itemId, result };
  },
  {
    flowControl: {
      key: '{workflow}.execute',
      parallelism: 10,
      ratePerSecond: 5,
    },
  },
);

一个实现细节值得注意:当前仓库用 runStep 辅助函数 统一包装 context.run() 调用(runStep(context, 'step-name', fn)),所有层的步骤名依然保持"工作流前缀 + 语义名"的约定,与规范中的 context.run() 用法等价。

六、最佳实践与常见坑

以下内容来自 best-practices.md,是脚手架落地后必须逐条对照的清单。

1. 错误处理:显式失败,而不是异常炸掉工作流

export const { POST } = serve<Payload>(
  async (context) => {
    const { itemId } = context.requestPayload ?? {};

    if (!itemId) {
      return { success: false, error: 'Missing itemId in payload' };
    }

    try {
      const result = await context.run('step-name', () => doWork(itemId));
      return { success: true, itemId, result };
    } catch (error) {
      console.error('[workflow:error]', error);
      return {
        success: false,
        error: error instanceof Error ? error.message : 'Unknown error',
      };
    }
  },
  { flowControl: { ... } },
);

2. 日志:统一前缀方便跨 QStash 仪表盘与 grep 排查

console.log('[{workflow}:{layer}] Starting with payload:', payload);
console.log('[{workflow}:{layer}] Processing items:', { count: items.length });
console.log('[{workflow}:{layer}] Completed:', result);
console.error('[{workflow}:{layer}:error]', error);

仓库实现中则使用 debug 包(如 debug('lobe-server:workflows:run-benchmark'))作为结构化替代,前缀约定不变。

3. 返回值:形状要匹配层的职责

入口层返回统计信息,执行层返回单条结果:

// 成功
return { success: true, itemId, result, message: 'Optional success message' };

// 错误
return { success: false, error: 'Error description', itemId };

// 统计(入口点)
return {
  success: true,
  totalEligible: 100,
  toProcess: 80,
  alreadyProcessed: 20,
  dryRun: true, // 如适用
  message: 'Summary message',
};

4. flowControl:按层调并发

// Layer 1: 入口 —— 单实例,避免重复处理
flowControl: { key: '{workflow}.process',  parallelism: 1,  ratePerSecond: 1 }

// Layer 2: 分页 —— 中等并发
flowControl: { key: '{workflow}.paginate', parallelism: 20, ratePerSecond: 5 }

// Layer 3: 执行 —— 更高并发以并行处理条目
flowControl: { key: '{workflow}.execute',  parallelism: 10, ratePerSecond: 5 }

为什么是这些默认值

  • Layer 1 恒用 parallelism: 1,保证并发触发不会让两个实例启动同一批任务;
  • Layer 2 可以较宽地扇出(10–20),因为分页本身开销小;
  • Layer 3 默认封顶 5–10,依据外部 API 的速率限制上下调。

5. context.run() 的黄金法则:步骤名必须唯一

Upstash 按步骤名(step name)做去重(dedupe)——这是"重试不重复执行"的机制本身,但也意味着重复步骤名会静默丢数据

// ✅ 唯一步骤名
await Promise.all(
  items.map((item) => context.run(`{workflow}:execute:${item.id}`, () => processItem(item))),
);

// ❌ 同名步骤 —— Upstash 按步骤名去重,你会丢失数据
await Promise.all(items.map((item) => context.run(`{workflow}:execute`, () => processItem(item))));

其他要求:步骤名带前缀且描述性强({workflow}:step-name);每个步骤保持幂等(可安全重试);不要嵌套 context.run(),保持扁平。

6. Payload 校验:快速失败

在函数顶部校验,让失败显式化,而不是让 undefined 级联成后期难以定位的故障:

const { itemId, configId } = context.requestPayload ?? {};

if (!itemId)   return { success: false, error: 'Missing itemId in payload' };
if (!configId) return { success: false, error: 'Missing configId in payload' };

7. 数据库连接:每个工作流只取一次

getServerDB() 是异步的,在每个步骤里重复获取会白白增加延迟:

const db = await getServerDB();

const item = await context.run('get-item', () => itemModel.findById(db, itemId));
const result = await context.run('save-result', () => resultModel.create(db, result));

8. 测试:dry-run 路径与完整路径都要覆盖

describe('WorkflowName', () => {
  it('should process items successfully', async () => {
    const items = await createTestItems();
    await WorkflowClass.triggerProcessItems({ dryRun: false });
    await waitForCompletion();
    const results = await getResults();
    expect(results).toHaveLength(items.length);
  });

  it('should support dryRun mode', async () => {
    const result = await WorkflowClass.triggerProcessItems({ dryRun: true });
    expect(result).toMatchObject({
      success: true,
      dryRun: true,
      totalEligible: expect.any(Number),
      toProcess: expect.any(Number),
    });
  });
});

常见坑速查

后果 正确做法
复用 context.run() 步骤名 Upstash 按名去重,批量处理只剩一条 步骤名带上条目 ID:`process:${item.id}`
跳过 payload 校验 undefined 级联成后期难懂的失败 顶部 fail fast,返回 { success: false, error }
跳过 filter 步骤 已处理的数据被重复处理 filterExisting(allItems) 再触发执行
日志前缀不一致 无法按 workflow+layer grep 定位 统一 [workflow:layer] 前缀

七、环境配置与 QStash 客户端

环境变量

所有工作流依赖以下环境变量(SKILL.md 的"Environment Variables"一节):

# 所有工作流必需
APP_URL=https://your-app.com # 工作流端点的基础 URL
QSTASH_TOKEN=qstash_xxx      # QStash 认证令牌

# 可选(自定义 QStash URL)
QSTASH_URL=https://custom-qstash.com

APP_URLQSTASH_TOKEN 的强校验体现在 Workflow 类的 getWorkflowUrl / 客户端构造中:缺少任一项会直接抛出 'APP_URL is required to trigger workflows' / 'QSTASH_TOKEN is required to trigger workflows' 错误(见 agentEvalRun/index.ts L131-L135)。

仓库中统一封装的客户端

当前仓库不鼓励每个 Workflow 类各自 new Client,而是复用 src/libs/qstash 中的单例:

  • qstashClientOtelQstashClient,作为 Upstash Workflow serve() 的选项传入,携带 Vercel Deployment Protection 的 bypass 头(VERCEL_AUTOMATION_BYPASS_SECRET),并对每次 publishJSON 记录 OTEL 指标;
  • workflowClientOtelWorkflowClient,用于 workflowClient.trigger() 触发工作流,成功/失败都会经 recordUpstashWorkflowEvent 上报(含 workflowRunIdstatus、重试参数等);
  • verifyQStashSignature:用 QSTASH_CURRENT_SIGNING_KEY / QSTASH_NEXT_SIGNING_KEY 校验 QStash 回调签名(未配置签名密钥时跳过校验),配套中间件见 qstashAuth

Hono 路由的注册方式(serve 来自 @upstash/workflow/hono,并叠加 OTEL 包装)可参考 agent-eval-run 路由

app.post(
  '/run-benchmark',
  serve(
    withOtelMetricsForUpstashWorkflows(runBenchmarkHandler, {
      url: '/api/workflows/agent-eval-run/run-benchmark',
    }),
    { ...runBenchmarkWorkflowOptions, qstashClient },
  ),
);

八、仓库内的端到端实例

examples.md 给出了两个"逐字遵循本模式"的工作流(位于 lobehub-cloud 私有仓,作为模式参照):

实例 1:welcome-placeholder(为用户生成 AI 欢迎占位内容)

  • Layer 1 process-users:入口,检查符合条件的用户;
  • Layer 2 paginate-users:分页遍历活跃用户;
  • Layer 3 generate-user:为单个用户生成占位内容。

关键特性:过滤掉 Redis 中已有缓存占位内容的用户;paidOnly 标记限定订阅用户;dryRun 统计模式;大批量用户走 fan-out(CHUNK_SIZE=20)。Layer 3 形态:

export const { POST } = serve<GenerateUserPlaceholderPayload>(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 };
});

实例 2:agent-welcome(为 AI Agent 生成欢迎语与开场问题)

三层结构同构(process-agents / paginate-agents / generate-agent),同样具备 Redis 缓存过滤、paidOnlydryRun 与 fan-out。

两个工作流的差异仅在于:实体类型(users vs agents)、业务逻辑(占位内容生成 vs 欢迎语生成)、数据源(不同的数据库查询)。其余一切——三层拆分、dry-run 处理、fan-out、filter-existing、flowControl 调参——完全相同。这正是该模式的意义:一旦内化该模式,新增工作流基本就是"实体替换"

在开源仓内,agent-eval-run(评测基准运行)与 memory-user-memory(用户记忆提取)是可直接阅读的端到端实现,目录分别位于 agent-eval-runmemory-user-memory

九、云项目(lobehub-cloud)部署模式

如果工作流要在 lobehub-cloud 私有仓中部署,遵循 cloud.md 的约定:

两种归属模式

  1. Cloud-Only 工作流:功能专属云用户(AI 生成、付费特性),直接实现在 lobehub-cloud/src/app/(backend)/api/workflows/ 下,无需 re-export;
  2. Re-export 工作流:核心实现在开源 lobehub,云端只需转发。原因是云部署必须 serve 这些端点,而 lobehub 子模块代码不能被云路由直接访问;同时为未来的云专属覆写留了口子。

Re-export 的正确写法

// lobehub-cloud/src/app/(backend)/api/workflows/feature/layer/route.ts
export { POST } from 'lobehub/src/app/(backend)/api/workflows/feature/layer/route';

关键点:必须使用 lobehub/src/... 路径,不能@/...,否则会触发 Circular definition of import alias 'POST' 循环导入错误。

TypeScript 路径映射

云项目通过 tsconfig paths 实现"云代码优先、开源回退"的解析顺序:

// lobehub-cloud/tsconfig.json
{
  "compilerOptions": {
    "paths": {
      "@/*": ["./src/*", "./lobehub/src/*"]
    }
  }
}

解析顺序:先查 ./src/*(云代码),再回退 ./lobehub/src/*(开源),允许云端覆写特定模块而其余部分沿用 lobehub 默认。

归属决策与迁移

  • 放到 lobehub(开源):功能对所有用户有用、无专有业务逻辑、可以开源;
  • 放到 cloud(私有):付费/高级特性、依赖云专属服务、包含专有算法;
  • 迁移工作流从云到开源的五步:拷贝路由 → 去除云专属依赖(换成通用接口)→ 云端建 re-export → Workflow 类移入 apps/server/src/workflows/ → 更新云端 import 为 lobehub/apps/server/src/workflows/feature

十、新工作流开发 Checklist

以下清单完整继承自 SKILL.md 的"Checklist for New Workflows"一节,建议在新建工作流时逐项打勾。

规划

  • [ ] 明确要处理的实体(users、agents、items、…)
  • [ ] 定义单条数据的业务逻辑
  • [ ] 确定过滤逻辑(Redis 缓存、数据库状态、…)

实现

  • [ ] 用 TypeScript interface 定义各层 payload 类型
  • [ ] 创建带静态 trigger 方法的 Workflow 类
  • [ ] Layer 1:入口路由,支持 dry-run
  • [ ] Layer 1:过滤逻辑,避免重复工作
  • [ ] Layer 2:带 fan-out 的分页
  • [ ] Layer 3单任务执行(每次运行一条)
  • [ ] 为每一层配置合适的 flowControl
  • [ ] 使用带工作流前缀的一致化日志
  • [ ] 校验所有必需的 payload 参数
  • [ ] 保证 context.run() 步骤名唯一

质量与部署

  • [ ] 返回一致的响应形状
  • [ ] 配置云部署(如在 lobehub-cloud,参见 references/cloud.md
  • [ ] 编写集成测试(dryRun 路径 + 完整执行路径)
  • [ ] 先用 dry-run 做冒烟测试
  • [ ] 全量上线前用小批量试跑验证

结语

LobeHub 的 Upstash Workflow 规范本质上是对"平台三大约束"(限流、步骤上限、幂等)的一组固定应答:入口层单实例 + dry-run 预演守住幂等与可观测性,分页层用游标 + fan-out 把任意规模的批量拆到步骤上限以内,执行层恒定单条 + flowControl 把并发压进外部 API 的速率预算。三层各自独立、可单独重试与限流,组合起来就是一条可水平扩展的后台处理管道。仓库中的 agent-eval-runmemory-user-memory 等实现都是这套模式的直接产物,新工作流只需按本文的模板与 Checklist 完成"实体替换"即可上线。

登录后查看全文
热门项目推荐
相关项目推荐