inngest-durable-functions

inngest-durable-functions

用于构建必须能在进程崩溃后存活、失败时自动重试、按计划运行、响应事件或在基础设施故障时保持状态的函数——例如,丢弃事件的 Webhook 处理器、不稳定的定时任务、中途失败的后台作业,或需要从中断处恢复的工作流。涵盖 Inngest 函数配置、触发器(事件、cron、invoke)、步骤执行与记忆化、幂等性、取消、错误处理、重试、日志记录和可观测性。

26Star
5Fork
更新于 2026/7/2
SKILL.md
readonly只读
name
inngest-durable-functions
description

用于构建必须能在进程崩溃后存活、失败时自动重试、按计划运行、响应事件或在基础设施故障时保持状态的函数——例如,丢弃事件的 Webhook 处理器、不稳定的定时任务、中途失败的后台作业,或需要从中断处恢复的工作流。涵盖 Inngest 函数配置、触发器(事件、cron、invoke)、步骤执行与记忆化、幂等性、取消、错误处理、重试、日志记录和可观测性。

Inngest 持久化函数

掌握 Inngest 的持久化执行模型,构建容错、长时间运行的工作流。本技能涵盖从触发器到错误处理的完整生命周期。

这些技能专注于 TypeScript。 对于 Python 或 Go,请参考 Inngest 文档 获取语言特定指导。核心概念适用于所有语言。

你需要了解的核心概念

持久化执行模型

  • 每个步骤 应封装副作用和非确定性代码
  • 记忆化 防止已完成的步骤重新执行
  • 状态持久化 在基础设施故障时存活
  • 自动重试 可配置重试次数

步骤执行流程

// ❌ 错误:步骤外的非确定性逻辑
async ({ event, step }) => {
  const timestamp = Date.now(); // 这会运行多次!

  const result = await step.run("process-data", () => {
    return processData(event.data);
  });
};

// ✅ 正确:所有非确定性逻辑放在步骤中
async ({ event, step }) => {
  const result = await step.run("process-with-timestamp", () => {
    const timestamp = Date.now(); // 只运行一次
    return processData(event.data, timestamp);
  });
};

函数限制

每个 Inngest 函数都有以下硬限制:

  • 每个函数运行最多 1,000 个步骤
  • 每个步骤返回的数据最多 4MB
  • 函数运行状态(包括事件数据、步骤输出和函数输出)合计最多 32MB
  • 每个步骤 = 单独的 HTTP 请求(约 50-100ms 开销)

如果达到这些限制,请将函数拆分为更小的函数,通过 step.invoke()step.sendEvent() 连接。

何时使用步骤

始终包裹在 step.run() 中:

  • API 调用和网络请求
  • 数据库读写
  • 文件 I/O 操作
  • 任何非确定性操作
  • 任何希望在失败时独立重试的操作

永远不要包裹在 step.run() 中:

  • 纯计算和数据转换
  • 简单的验证逻辑
  • 无副作用的确定性操作
  • 日志记录(在步骤外使用)

函数创建

基本函数结构

const processOrder = inngest.createFunction(
  {
    id: "process-order", // 唯一,切勿更改
    triggers: [{ event: "order/created" }],
    retries: 4, // 默认:每个步骤 4 次重试
    concurrency: 10 // 最大并发执行数
  },
  async ({ event, step }) => {
    // 你的持久化工作流
  }
);

步骤 ID 和记忆化

// 步骤 ID 可以重用 - Inngest 自动处理计数器
const data = await step.run("fetch-data", () => fetchUserData());
const more = await step.run("fetch-data", () => fetchOrderData()); // 不同的执行

// 使用描述性 ID 以提高清晰度
await step.run("validate-payment", () => validatePayment(event.data.paymentId));
await step.run("charge-customer", () => chargeCustomer(event.data));
await step.run("send-confirmation", () => sendEmail(event.data.email));

触发器和事件

事件触发器

触发器在 createFunction 第一个参数的 triggers 数组中定义:

// 单个事件触发器
inngest.createFunction(
  { id: "my-fn", triggers: [{ event: "user/signup" }] },
  async ({ event }) => { /* ... */ }
);

// 带条件过滤的事件
inngest.createFunction(
  { id: "my-fn", triggers: [{ event: "user/action", if: 'event.data.action == "purchase" && event.data.amount > 100' }] },
  async ({ event }) => { /* ... */ }
);

// 多个触发器(最多 10 个)
inngest.createFunction(
  {
    id: "my-fn",
    triggers: [
      { event: "user/signup" },
      { event: "user/login", if: 'event.data.firstLogin == true' },
      { cron: "0 9 * * *" } // 每天上午 9 点
    ]
  },
  async ({ event }) => { /* ... */ }
);

Cron 触发器

// 基本 cron
inngest.createFunction(
  { id: "my-fn", triggers: [{ cron: "0 */6 * * *" }] }, // 每 6 小时
  async ({ step }) => { /* ... */ }
);

// 带时区
inngest.createFunction(
  { id: "my-fn", triggers: [{ cron: "TZ=Europe/Paris 0 12 * * 5" }] }, // 巴黎时间周五中午 12 点
  async ({ step }) => { /* ... */ }
);

// 与事件结合
inngest.createFunction(
  {
    id: "my-fn",
    triggers: [
      { event: "manual/report.requested" },
      { cron: "0 0 * * 0" } // 每周日午夜
    ]
  },
  async ({ event, step }) => { /* ... */ }
);

函数调用

// 作为步骤调用另一个函数
const result = await step.invoke("generate-report", {
  function: generateReportFunction,
  data: { userId: event.data.userId }
});

// 使用返回的数据
await step.run("process-report", () => {
  return processReport(result);
});

幂等性策略

事件级幂等性(生产者端)

// 使用自定义 ID 防止重复事件
await inngest.send({
  id: `checkout-completed-${cartId}`, // 24 小时去重
  name: "cart/checkout.completed",
  data: { cartId, email: "user@example.com" }
});

函数级幂等性(消费者端)

const sendEmail = inngest.createFunction(
  {
    id: "send-checkout-email",
    triggers: [{ event: "cart/checkout.completed" }],
    // 每个 cartId 每 24 小时只运行一次
    idempotency: "event.data.cartId"
  },
  async ({ event, step }) => {
    // 此函数不会对同一 cartId 运行两次
  }
);

// 复杂幂等键
const processUserAction = inngest.createFunction(
  {
    id: "process-user-action",
    triggers: [{ event: "user/action.performed" }],
    // 每个用户 + 组织组合唯一
    idempotency: 'event.data.userId + "-" + event.data.organizationId'
  },
  async ({ event, step }) => {
    /* ... */
  }
);

取消模式

基于事件的取消

在表达式中,event = 原始触发事件,async = 被匹配的事件。有关完整细节,请参阅表达式语法参考

const processOrder = inngest.createFunction(
  {
    id: "process-order",
    triggers: [{ event: "order/created" }],
    cancelOn: [
      {
        event: "order/cancelled",
        if: "event.data.orderId == async.data.orderId"
      }
    ]
  },
  async ({ event, step }) => {
    await step.sleepUntil("wait-for-payment", event.data.paymentDue);
    // 如果收到 order/cancelled 事件,将被取消
    await step.run("charge-payment", () => processPayment(event.data));
  }
);

超时取消

const processWithTimeout = inngest.createFunction(
  {
    id: "process-with-timeout",
    triggers: [{ event: "long/process.requested" }],
    timeouts: {
      start: "5m", // 如果 5 分钟内未开始则取消
      finish: "30m" // 如果 30 分钟内未完成则取消
    }
  },
  async ({ event, step }) => {
    /* ... */
  }
);

处理取消清理

// 监听取消事件
const cleanupCancelled = inngest.createFunction(
  { id: "cleanup-cancelled-process", triggers: [{ event: "inngest/function.cancelled" }] },
  async ({ event, step }) => {
    if (event.data.function_id === "process-order") {
      await step.run("cleanup-resources", () => {
        return cleanupOrderResources(event.data.run_id);
      });
    }
  }
);

错误处理和重试

默认重试行为

  • 每个步骤 总共 5 次尝试(1 次初始 + 4 次重试)
  • 指数退避 带抖动
  • 每个步骤 独立的重试计数器

自定义重试配置

const reliableFunction = inngest.createFunction(
  {
    id: "reliable-function",
    triggers: [{ event: "critical/task" }],
    retries: 10 // 每个步骤最多 10 次重试
  },
  async ({ event, step, attempt }) => {
    // `attempt` 是函数级尝试计数器(从 0 开始)
    // 它跟踪当前执行步骤的重试次数,而不是整个函数
    if (attempt > 5) {
      // 当前步骤后续尝试的不同逻辑
    }
  }
);

不可重试错误

防止对重试不会成功的代码进行重试。

import { NonRetriableError } from "inngest";

const processUser = inngest.createFunction(
  { id: "process-user", triggers: [{ event: "user/process.requested" }] },
  async ({ event, step }) => {
    const user = await step.run("fetch-user", async () => {
      const user = await db.users.findOne(event.data.userId);

      if (!user) {
        // 不要重试 - 用户不存在
        throw new NonRetriableError("User not found, stopping execution");
      }

      return user;
    });

    // 继续处理...
  }
);

自定义重试时机

import { RetryAfterError } from "inngest";

const respectRateLimit = inngest.createFunction(
  { id: "api-call", triggers: [{ event: "api/call.requested" }] },
  async ({ event, step }) => {
    await step.run("call-api", async () => {
      const response = await externalAPI.call(event.data);

      if (response.status === 429) {
        // 根据 API 指定的时间重试
        const retryAfter = response.headers["retry-after"];
        throw new RetryAfterError("Rate limited", `${retryAfter}s`);
      }

      return response.data;
    });
  }
);

日志记录最佳实践

正确的日志设置

import winston from "winston";

// 配置日志记录器
const logger = winston.createLogger({
  level: "info",
  format: winston.format.json(),
  transports: [new winston.transports.Console()]
});

const inngest = new Inngest({
  id: "my-app",
  logger // 将日志记录器传递给客户端
});

// 或使用内置的 ConsoleLogger 进行简单的日志级别控制
import { ConsoleLogger, Inngest } from "inngest";

const inngest = new Inngest({
  id: "my-app",
  logger: new ConsoleLogger({ level: "debug" }) // "debug" | "info" | "warn" | "error"
});

⚠️ v4 重大变更: logLevel 选项已被移除。请使用 logger 选项配合 ConsoleLogger 或自定义日志记录器。

函数日志记录模式

const processData = inngest.createFunction(
  { id: "process-data", triggers: [{ event: "data/process.requested" }] },
  async ({ event, step, logger }) => {
    // ✅ 正确:在步骤内记录日志以避免重复
    const result = await step.run("fetch-data", async () => {
      logger.info("Fetching data for user", { userId: event.data.userId });
      return await fetchUserData(event.data.userId);
    });

    // ❌ 避免:在步骤外记录日志可能导致重复
    // logger.info("Processing complete"); // 这可能会运行多次!

    await step.run("log-completion", async () => {
      logger.info("Processing complete", { resultCount: result.length });
    });
  }
);

性能优化

检查点

检查点在 v4 中默认启用。它允许函数在执行期间定期持久化状态,减少步骤之间的延迟。

// 检查点在 v4 中默认启用
// 为无服务器平台配置 maxRuntime(设置为平台超时时间的 60-80%)
const realTimeFunction = inngest.createFunction(
  {
    id: "real-time-function",
    triggers: [{ event: "realtime/process" }],
    checkpointing: {
      maxRuntime: "50s", // 对于超时时间为 60s 的无服务器环境
    }
  },
  async ({ event, step }) => {
    // 步骤立即执行,并定期检查点
    const result1 = await step.run("step-1", () => process1(event.data));
    const result2 = await step.run("step-2", () => process2(result1));
    return { result2 };
  }
);

// 如果需要,禁用检查点
const legacyFunction = inngest.createFunction(
  {
    id: "legacy-function",
    triggers: [{ event: "legacy/process" }],
    checkpointing: false
  },
  async ({ event, step }) => { /* ... */ }
);

高级模式

条件步骤执行

const conditionalProcess = inngest.createFunction(
  { id: "conditional-process", triggers: [{ event: "process/conditional" }] },
  async ({ event, step }) => {
    const userData = await step.run("fetch-user", () => {
      return getUserData(event.data.userId);
    });

    // 条件步骤执行
    if (userData.isPremium) {
      await step.run("premium-processing", () => {
        return processPremiumFeatures(userData);
      });
    }

    // 始终运行
    await step.run("standard-processing", () => {
      return processStandardFeatures(userData);
    });
  }
);

错误恢复模式

const robustProcess = inngest.createFunction(
  { id: "robust-process", triggers: [{ event: "process/robust" }] },
  async ({ event, step }) => {
    let primaryResult;

    try {
      primaryResult = await step.run("primary-service", () => {
        return callPrimaryService(event.data);
      });
    } catch (error) {
      // 回退到辅助服务
      primaryResult = await step.run("fallback-service", () => {
        return callSecondaryService(event.data);
      });
    }

    return { result: primaryResult };
  }
);

常见错误避免

  1. ❌ 步骤外的非确定性代码
  2. ❌ 步骤外的数据库调用
  3. ❌ 步骤外的日志记录(导致重复)
  4. ❌ 部署后更改步骤 ID
  5. ❌ 未处理 NonRetriableError 情况
  6. ❌ 忽略关键函数的幂等性

下一步


本技能涵盖 Inngest 的持久化函数模式。有关事件发送和 Webhook 处理,请参阅 inngest-events 技能。