Background jobs

Inngest

具备自动重试、事件扇出和本地开发服务器的持久多步骤后台函数——由 @nebutra/event-bus 提供支持。

概述

Inngest 是 Nebutra 的主要后台任务系统。它将每个函数步骤视为持久检查点——如果某个步骤失败,只有该步骤会被重试,而不是整个任务。这使其非常适合多步骤工作流,如入职流程、报表生成流水线和通知序列。

Nebutra 通过 @nebutra/event-bus 封装 Inngest,在整个 Monorepo 中预配置客户端并统一暴露。

配置步骤

INNGEST_EVENT_KEY=""      # 来自 Inngest 控制台 → Event Keys
INNGEST_SIGNING_KEY=""    # 来自 Inngest 控制台 → Signing Keys

API 网关已暴露 POST /api/v1/inngest。这是 Inngest 调用以执行函数的 URL。在 Inngest 控制台Apps → Add App 中注册:

https://api.yourdomain.com/api/v1/inngest

本地开发时,请改用 Inngest Dev Server(见下文)。

// packages/integrations/event-bus/src/functions/welcome-email.ts
import { inngest } from "@nebutra/event-bus";
import { db } from "@nebutra/db";
import { email } from "@nebutra/email";

export const sendWelcomeEmail = inngest.createFunction(
  { id: "send-welcome-email", retries: 3 },
  { event: "user/signed-up" },
  async ({ event, step }) => {
    const user = await step.run("fetch-user", async () => {
      return db.user.findUnique({ where: { id: event.data.userId } });
    });

    await step.run("send-email", async () => {
      return email.send({
        to: user.email,
        template: "welcome",
        data: { userName: user.name },
      });
    });
  }
);
import { inngest } from "@nebutra/event-bus";

// 在服务端代码的任意位置调用
await inngest.send({
  name: "user/signed-up",
  data: { userId: "user_abc123" },
});

函数定义

基本函数

import { inngest } from "@nebutra/event-bus";

export const myFunction = inngest.createFunction(
  {
    id: "my-function",           // 唯一标识符——在控制台和日志中使用
    retries: 3,                  // 每个步骤的最大重试次数(默认:3)
    concurrency: {
      limit: 10,                 // 该函数的最大并发执行数
    },
  },
  { event: "resource/action" },  // 触发事件名称
  async ({ event, step }) => {
    // 函数体
  }
);

配置选项

选项类型默认值说明
idstring必填唯一函数标识符
retriesnumber3失败步骤的最大重试次数
concurrency.limitnumber无限制最大并发执行数
throttle.limitnumber单位时间内的最大执行次数
throttle.periodstring例如 "1m""1h"
timeouts.finishstring"1h"函数最大总执行时长

步骤函数

步骤是 Inngest 函数的核心构建块。每个 step.run() 调用都具备:

  • 记忆化 — 函数重试时,已完成的步骤不会重新执行
  • 独立重试 — 只有失败的步骤重试,而非整个函数
  • 检查点 — 步骤之间的进度会被保存
export const processOrder = inngest.createFunction(
  { id: "process-order", retries: 5 },
  { event: "order/placed" },
  async ({ event, step }) => {
    // 步骤 1:验证库存
    const inventory = await step.run("check-inventory", async () => {
      return checkStock(event.data.items);
    });

    if (!inventory.available) {
      // 提前返回是合法的——标记函数为已完成
      return { status: "out-of-stock" };
    }

    // 步骤 2:扣款(失败时独立重试)
    const charge = await step.run("charge-payment", async () => {
      return chargeCard(event.data.paymentMethodId, event.data.total);
    });

    // 步骤 3:发送确认邮件
    await step.run("send-confirmation", async () => {
      return email.send({
        to: event.data.customerEmail,
        template: "order-confirmation",
        data: { orderId: event.data.orderId, charge },
      });
    });

    return { status: "processed", chargeId: charge.id };
  }
);

步骤间等待

// 等待 24 小时后发送跟进邮件
await step.sleep("wait-one-day", "24h");

await step.run("send-followup", async () => {
  return email.send({ to: user.email, template: "day-one-followup" });
});

等待外部事件

// 暂停函数,等待特定事件到来(最长 7 天)
const approval = await step.waitForEvent("wait-for-approval", {
  event: "approval/granted",
  match: "data.requestId",    // 匹配 event.data.requestId === 当前 event.data.requestId
  timeout: "72h",
});

if (!approval) {
  // 超时——未收到审批
  await step.run("notify-expired", () => notifyRequestExpired(event.data.requestId));
}

重试配置

每个步骤独立重试,默认使用指数退避算法。

inngest.createFunction(
  {
    id: "resilient-function",
    retries: 5,                // 每个失败步骤最多重试 5 次
  },
  { event: "data/sync" },
  async ({ event, step, attempt }) => {
    // `attempt` 从 0 开始(第一次尝试)
    console.log(`第 ${attempt + 1} 次尝试`);

    await step.run("sync-data", async () => {
      return syncToExternalApi(event.data);
    });
  }
);

若要让步骤永久失败而不重试,抛出 NonRetriableError

import { NonRetriableError } from "inngest";

await step.run("validate", async () => {
  if (!isValid(data)) {
    throw new NonRetriableError("数据永久无效——跳过重试");
  }
});

事件命名规范

所有事件名称使用 资源/动作 规范:

事件触发来源
user/signed-up认证回调
user/deleted账户删除
invoice/payment-failed计费 Webhook
invoice/paid计费 Webhook
quota/threshold-reached计量流水线
export/requested控制台操作
report/scheduledCron 触发

事件名称使用小写字母和正斜杠分隔符。避免使用 job/run 等通用名称。

发送带数据的事件

import { inngest } from "@nebutra/event-bus";

// 发送单个事件
await inngest.send({
  name: "invoice/payment-failed",
  data: {
    tenantId: "org_123",
    invoiceId: "inv_456",
    amount: 99.99,
    customerId: "cust_789",
  },
});

// 批量发送事件(并行处理)
await inngest.send([
  { name: "user/signed-up", data: { userId: "user_1" } },
  { name: "user/signed-up", data: { userId: "user_2" } },
  { name: "user/signed-up", data: { userId: "user_3" } },
]);

Inngest Dev Server

本地开发时,运行 Inngest Dev Server 在本地执行函数,无需连接 Inngest Cloud:

npx inngest-cli@latest dev -u http://localhost:3001/api/v1/inngest

Dev Server 功能:

  • 监听 http://localhost:8288
  • 从接口 URL 自动发现函数
  • 提供 UI 用于触发事件和查看函数运行情况
  • 显示逐步执行追踪

在浏览器中打开 http://localhost:8288 访问控制台。

在本地 .env 中设置 INNGEST_DEV=true 可跳过与 Dev Server 配合时的签名验证。

扇出模式

一个事件可以同时触发多个函数。注册多个监听相同事件的函数:

// 两个函数都在同一事件上触发
export const sendWelcomeEmail = inngest.createFunction(
  { id: "send-welcome-email" },
  { event: "user/signed-up" },
  async ({ event, step }) => { /* 发送邮件 */ }
);

export const provisionWorkspace = inngest.createFunction(
  { id: "provision-workspace" },
  { event: "user/signed-up" },
  async ({ event, step }) => { /* 创建默认工作空间 */ }
);

export const trackSignUp = inngest.createFunction(
  { id: "track-sign-up" },
  { event: "user/signed-up" },
  async ({ event, step }) => { /* 发送分析事件 */ }
);

发送单个 user/signed-up 事件时,三个函数并行运行。

相关文档

How is this guide?

目录