Inngest
具备自动重试、事件扇出和本地开发服务器的持久多步骤后台函数——由 @nebutra/event-bus 提供支持。
概述
Inngest 是 Nebutra 的主要后台任务系统。它将每个函数步骤视为持久检查点——如果某个步骤失败,只有该步骤会被重试,而不是整个任务。这使其非常适合多步骤工作流,如入职流程、报表生成流水线和通知序列。
Nebutra 通过 @nebutra/event-bus 封装 Inngest,在整个 Monorepo 中预配置客户端并统一暴露。
配置步骤
INNGEST_EVENT_KEY="" # 来自 Inngest 控制台 → Event Keys
INNGEST_SIGNING_KEY="" # 来自 Inngest 控制台 → Signing KeysAPI 网关已暴露 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 }) => {
// 函数体
}
);配置选项
| 选项 | 类型 | 默认值 | 说明 |
|---|---|---|---|
id | string | 必填 | 唯一函数标识符 |
retries | number | 3 | 失败步骤的最大重试次数 |
concurrency.limit | number | 无限制 | 最大并发执行数 |
throttle.limit | number | — | 单位时间内的最大执行次数 |
throttle.period | string | — | 例如 "1m"、"1h" |
timeouts.finish | string | "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/scheduled | Cron 触发 |
事件名称使用小写字母和正斜杠分隔符。避免使用 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/inngestDev 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?
最后更新于