
inngest-steps
当需要实现必须能在进程重启后继续存在的延迟(例如24小时购物车放弃提醒、定时跟进)、等待人工审批或带超时的外部事件(审核关卡、webhook回调、异步API完成)、轮询外部服务且崩溃后不丢失状态、调用其他函数并等待其结果、记忆化昂贵操作以避免重试时重复执行,或在工作流中并行运行异步任务时使用。涵盖Inngest步骤方法:step.run、step.sleep、step.waitForEvent、step.waitForSignal、step.sendEvent、step.invoke、step.ai,以及循环和并行执行模式。
当需要实现必须能在进程重启后继续存在的延迟(例如24小时购物车放弃提醒、定时跟进)、等待人工审批或带超时的外部事件(审核关卡、webhook回调、异步API完成)、轮询外部服务且崩溃后不丢失状态、调用其他函数并等待其结果、记忆化昂贵操作以避免重试时重复执行,或在工作流中并行运行异步任务时使用。涵盖Inngest步骤方法:step.run、step.sleep、step.waitForEvent、step.waitForSignal、step.sendEvent、step.invoke、step.ai,以及循环和并行执行模式。
Inngest 步骤
使用 Inngest 的步骤方法构建健壮、持久的工作流。每个步骤都是一个独立的 HTTP 请求,可以独立重试和监控。
这些技能专注于 TypeScript。 对于 Python 或 Go,请参考 Inngest 文档 获取语言特定指导。核心概念适用于所有语言。
核心概念
🔄 关键:每个步骤都会从头重新执行你的函数。 将所有非确定性代码(API 调用、数据库查询、随机性)放在步骤内部,绝不要放在外部。
📊 步骤限制: 每个函数最多 1,000 个步骤,总步骤数据 4MB。
// ❌ 错误 - 会执行 4 次
export default inngest.createFunction(
{ id: "bad-example", triggers: [{ event: "test" }] },
async ({ step }) => {
console.log("这行会输出 4 次!"); // 在步骤外 = 不好
await step.run("a", () => console.log("a"));
await step.run("b", () => console.log("b"));
await step.run("c", () => console.log("c"));
}
);
// ✅ 正确 - 每个只输出一次
export default inngest.createFunction(
{ id: "good-example", triggers: [{ event: "test" }] },
async ({ step }) => {
await step.run("log-hello", () => console.log("hello"));
await step.run("a", () => console.log("a"));
await step.run("b", () => console.log("b"));
await step.run("c", () => console.log("c"));
}
);
step.run()
执行可重试的代码作为一个步骤。每个步骤 ID 可以重复使用 - Inngest 自动处理计数器。
// 基本用法
const result = await step.run("fetch-user", async () => {
const user = await db.user.findById(userId);
return user; // 始终返回有用的数据
});
// 同步代码也可以
const transformed = await step.run("transform-data", () => {
return processData(result);
});
// 副作用(无需返回值)
await step.run("send-notification", async () => {
await sendEmail(user.email, "Welcome!");
});
✅ 应该:
- 将所有非确定性逻辑放在步骤内部
- 返回后续步骤有用的数据
- 在循环中重用步骤 ID(计数器自动处理)
❌ 不应该:
- 不必要地将确定性逻辑放在步骤中
- 忘记每个步骤 = 单独的 HTTP 请求
step.sleep()
暂停执行而不消耗计算时间。
// 持续时间字符串
await step.sleep("wait-24h", "24h");
await step.sleep("short-delay", "30s");
await step.sleep("weekly-pause", "7d");
// 在工作流中使用
await step.run("send-welcome", () => sendEmail(email));
await step.sleep("wait-for-engagement", "3d");
await step.run("send-followup", () => sendFollowupEmail(email));
step.sleepUntil()
休眠到指定的日期时间。
const reminderDate = new Date("2024-12-25T09:00:00Z");
await step.sleepUntil("wait-for-christmas", reminderDate);
// 从事件数据中获取
const scheduledTime = new Date(event.data.remind_at);
await step.sleepUntil("wait-for-scheduled-time", scheduledTime);
step.waitForEvent()
🚨 关键:waitForEvent 只捕获在此步骤执行之后发送的事件。
- ❌ 在 waitForEvent 运行之前发送的事件 → 不会被捕获
- ✅ 在 waitForEvent 运行之后发送的事件 → 会被捕获
- 始终检查
null返回值(表示超时,事件从未到达)
// 带超时的基本事件等待
const approval = await step.waitForEvent("wait-for-approval", {
event: "app/invoice.approved",
timeout: "7d",
match: "data.invoiceId" // 简单匹配
});
// 基于表达式的匹配(CEL 语法)
const subscription = await step.waitForEvent("wait-for-subscription", {
event: "app/subscription.created",
timeout: "30d",
if: "event.data.userId == async.data.userId && async.data.plan == 'pro'"
});
// 处理超时
if (!approval) {
await step.run("handle-timeout", () => {
// 审批从未到达
return notifyAccountingTeam();
});
}
✅ 应该:
- 使用唯一 ID 进行匹配(userId、sessionId、requestId)
- 始终设置合理的超时时间
- 处理 null 返回值(超时情况)
- 与 Realtime 一起用于人工参与流程
❌ 不应该:
- 期望在此步骤之前发送的事件被处理
- 在生产中使用时不设置超时
表达式语法
在表达式中,event = 原始触发事件,async = 新匹配的事件。有关完整语法、运算符和模式,请参阅表达式语法参考。
step.waitForSignal()
等待唯一信号(不是事件)。更适合 1:1 匹配。
const taskId = "task-" + crypto.randomUUID();
const signal = await step.waitForSignal("wait-for-task-completion", {
signal: taskId,
timeout: "1h",
onConflict: "replace" // 必需:"replace" 覆盖待处理信号,"fail" 抛出错误
});
// 通过 Inngest API 或 SDK 在其他地方发送信号
// POST /v1/events 使用匹配 taskId 的信号
何时使用:
- waitForEvent:多个函数可能处理同一个事件
- waitForSignal:精确的 1:1 信号到特定函数运行
step.sendEvent()
扇出到其他函数,不等待结果。
// 触发其他函数
await step.sendEvent("notify-systems", {
name: "user/profile.updated",
data: { userId: user.id, changes: profileChanges }
});
// 一次发送多个事件
await step.sendEvent("batch-notifications", [
{ name: "billing/invoice.created", data: { invoiceId } },
{ name: "email/invoice.send", data: { email: user.email, invoiceId } }
]);
何时使用: 你想触发其他函数但不需要在当前函数中获取它们的结果。
step.invoke()
调用其他函数并处理它们的结果。非常适合组合。
const computeSquare = inngest.createFunction(
{ id: "compute-square", triggers: [{ event: "calculate/square" }] },
async ({ event }) => {
return { result: event.data.number * event.data.number };
}
);
// 调用并使用结果
const square = await step.invoke("get-square", {
function: computeSquare,
data: { number: 4 }
});
console.log(square.result); // 16,完全类型化!
// 对于跨应用调用(当无法直接导入函数时):
import { referenceFunction } from "inngest";
const externalFn = referenceFunction({
appId: "other-app",
functionId: "other-fn"
});
const result = await step.invoke("call-external", {
function: externalFn,
data: { key: "value" }
});
警告:v4 破坏性变更: 字符串函数 ID(例如 function: "my-app-other-fn")在 step.invoke() 中不再支持。使用导入的函数引用或 referenceFunction() 进行跨应用调用。
非常适合:
- 将复杂工作流分解为可组合的函数
- 在多个工作流中重用逻辑
- Map-reduce 模式
模式
带步骤的循环
重用步骤 ID - Inngest 自动处理计数器。
const allProducts = [];
let cursor = null;
let hasMore = true;
while (hasMore) {
// 相同的 ID "fetch-page" 被重用 - 计数器自动处理
const page = await step.run("fetch-page", async () => {
return shopify.products.list({ cursor, limit: 50 });
});
allProducts.push(...page.products);
if (page.products.length < 50) {
hasMore = false;
} else {
cursor = page.products[49].id;
}
}
await step.run("process-products", () => {
return processAllProducts(allProducts);
});
并行执行
使用 Promise.all 进行并行步骤。在 v4 中,并行步骤执行默认优化
// 创建步骤而不等待
const sendEmail = step.run("send-email", async () => {
return await sendWelcomeEmail(user.email);
});
const updateCRM = step.run("update-crm", async () => {
return await crmService.addUser(user);
});
const createSubscription = step.run("create-subscription", async () => {
return await subscriptionService.create(user.id);
});
// 全部并行运行
const [emailId, crmRecord, subscription] = await Promise.all([
sendEmail,
updateCRM,
createSubscription
]);
// 在 v4 中,并行步骤默认优化
export default inngest.createFunction(
{
id: "parallel-heavy-function",
triggers: [{ event: "process/batch" }]
},
async ({ event, step }) => {
const results = await Promise.all(
event.data.items.map((item, i) =>
step.run(`process-item-${i}`, () => processItem(item))
)
);
}
);
// ⚠️ Promise.race() 在 v4 优化并行下的行为:
// 所有 promise 在 race 解决之前都会完成。使用 group.parallel() 实现真正的竞争:
const winner = await group.parallel(async () => {
return Promise.race([
step.run("fast-service", () => callFastService()),
step.run("slow-service", () => callSlowService())
]);
});
// 如果需要禁用优化并行:
// 在客户端级别:new Inngest({ id: "app", optimizeParallelism: false })
// 在函数级别:{ id: "fn", optimizeParallelism: false, triggers: [...] }
有关并发和限流选项,请参阅 inngest-flow-control。
分块作业
非常适合带并行步骤的批处理。
export default inngest.createFunction(
{ id: "process-large-dataset", triggers: [{ event: "data/process.large" }] },
async ({ event, step }) => {
const chunks = chunkArray(event.data.items, 10);
// 并行处理块
const results = await Promise.all(
chunks.map((chunk, index) =>
step.run(`process-chunk-${index}`, () => processChunk(chunk))
)
);
// 合并结果
await step.run("combine-results", () => {
return aggregateResults(results);
});
}
);
关键陷阱
🔄 函数重新执行: 步骤外的代码在每次步骤执行时都会运行
⏰ 事件时机: waitForEvent 只捕获在步骤运行之后发送的事件
🔢 步骤限制: 每个函数最多 1,000 个步骤,每个步骤输出 4MB,每次函数运行总计 32MB
📨 HTTP 请求: 在 v4 中默认启用检查点,减少 HTTP 开销。对于无服务器平台,在客户端配置 maxRuntime
🔁 步骤 ID: 可以在循环中重用 - Inngest 处理计数器
⚡ 并行: 使用 Promise.all 进行并行步骤(v4 默认优化)。注意 Promise.race() 会等待所有 promise 完成——使用 group.parallel() 实现真正的竞争语义
常见用例
- 人工参与流程: waitForEvent + Realtime UI
- 多步骤入门: 步骤之间 sleep,waitForEvent 等待用户操作
- 数据处理: 并行步骤处理分块工作
- 外部集成: step.run 用于可靠的 API 调用
- AI 工作流: step.ai 用于持久的 LLM 编排
- 函数组合: step.invoke 构建复杂工作流
记住:步骤使你的函数持久、可观察和可调试。拥抱它们!





