Process large datasets without starting over
Use steps, fan-out, flow control, and waits to run jobs that take hours or days, survive crashes and deploys, and never redo finished work.
If you run nightly imports, backfills, data pipelines, or any job that works through more records than fit in one request, you've probably seen it fail partway through and start over from zero. A pod runs out of memory, a deploy restarts the process, or the database chokes under load, and the job's progress is gone with it. These four primitives let a job commit progress as it goes, split into pieces that fail independently, and pace itself against your systems, so you don't have to write checkpointing or retry logic by hand.
§Which primitive do you need?
| Symptom | Primitive | What it does |
|---|---|---|
| A crash or deploy restarts the whole job from the beginning | Steps | Saves each step's result so a retry resumes where it failed |
| One job loops over every record, hits limits, or runs out of memory | Fan-out | Splits the work into chunks, each running as its own function run |
| Parallel chunks overwhelm your database or an upstream API | Concurrency and throttle | Caps how much work runs at once and how fast new runs start |
| The job pauses for an approval or webhook and ties up a worker | Waits | Holds the run's state without holding compute |
Most large jobs need the first three together. Add waits when the job depends on something outside it.
Steps vs. fan-out: when it's not obvious
Steps divide one run into units that each commit. Fan-out divides the job into many runs. A single function can have at most 1,000 steps, so if the number of steps grows with your dataset, fan out. Each run should do a bounded amount of work no matter how big the input gets.
§Primitive 1: Steps
You need this if: Your job fails partway through and the only safe recovery is to run it again from the start.
Without steps: Progress lives in the process. When the process dies, the job can't tell what finished, so it redoes everything, including work that already wrote to your database.
01import { inngest } from "./client";0203export const processRecordPage = inngest.createFunction(04 {05 id: "process-record-page",06 triggers: [{ event: "records/page.ready" }],07 retries: 5,08 },09 async ({ event, step }) => {10 const { page, pageSize } = event.data;1112 // Load, transform, and stage this page. If a later step fails,13 // this one is not re-run: its result is memoized.14 const transformed = await step.run("transform-page", async () => {15 return await transformAndStage(page, pageSize); // returns a count16 });1718 // Match the staged records. A failure here retries only this step.19 const matched = await step.run("match-page", async () => {20 return await matchStaged(page, pageSize); // returns a count21 });2223 return { page, transformed, matched };24 }25);Adapt it:
- Make each step a unit of work that's safe to retry by itself. If a step writes to an external system, make the write idempotent (upsert, or pass an idempotency key) so a retry after a partial write doesn't duplicate data.
- Return small values from steps, like counts or IDs, never the dataset itself. Step results are stored as run state, so write the data to your own storage and pass references.
retriesapplies to each step independently. A step that exhausts its retries fails the run; steps that already completed stay completed.- Deploying mid-run is safe: completed steps are matched by ID and skipped on the new code. Renaming a step's ID makes it a new step, so it runs again. See Versioning.
§Primitive 2: Fan-out
You need this if: Your dataset is large or growing, and a single run that loops over every record would exceed step limits, run out of memory, or let one bad record block everything.
Without fan-out: One giant run holds the whole job. A single failure stalls all of it, and no worker can process the data without loading more of it than it should.
01import { inngest } from "./client";0203const PAGE_SIZE = 500;04const MAX_EVENTS_PER_SEND = 5_000; // max events per send0506export const nightlyImport = inngest.createFunction(07 {08 id: "nightly-import",09 triggers: [{ cron: "0 1 * * *" }],10 },11 async ({ step }) => {12 // Count only. The coordinator never loads the records themselves.13 const total = await step.run("count-records", async () => {14 return await countPendingRecords();15 });1617 const pages = Math.ceil(total / PAGE_SIZE);18 const events = Array.from({ length: pages }, (_, page) => ({19 name: "records/page.ready",20 data: { page, pageSize: PAGE_SIZE },21 }));2223 // Each send is its own memoized step, so a crash here won't re-send24 // pages that were already dispatched.25 for (let i = 0; i < events.length; i += MAX_EVENTS_PER_SEND) {26 await step.sendEvent(27 `fan-out-pages-${i}`,28 events.slice(i, i + MAX_EVENTS_PER_SEND)29 );30 }3132 return { total, pages };33 }34);Each records/page.ready event starts a run of processRecordPage from Primitive 1. Every page retries on its own, and a failure on one page never touches the others.
Adapt it:
- Size
PAGE_SIZEso one page comfortably finishes within a step. Smaller pages mean cheaper retries and less memory per worker; larger pages mean fewer runs. - Put identifiers or ranges in the event, not record data. Event payloads should stay small.
- A single
step.sendEventcall accepts up to 5,000 events, so chunk larger sends as shown. - If you need to act after every page finishes, have each worker report completion (for example, by incrementing a counter in your database) and trigger a final step when the count reaches
pages.
§Primitive 3: Concurrency and throttle
You need this if: Your fan-out produces hundreds or thousands of runs, and running them all at once would exhaust database connections or trip an upstream API's rate limit.
Without flow control: Every page starts at once, your database slows down for everything else, and the upstream API answers with 429s that turn into retries that make it worse.
This is the same worker from Primitive 1, with flow control added to its configuration.
01export const processRecordPage = inngest.createFunction(02 {03 id: "process-record-page",04 triggers: [{ event: "records/page.ready" }],05 retries: 5,06 concurrency: [07 { limit: 50 }, // at most 50 steps executing at once, across all pages08 ],09 throttle: {10 limit: 100, // page runs started per minute11 period: "1m", // stays under the upstream API limit12 },13 },14 async ({ event, step }) => {15 // ...same steps as Primitive 116 }17);Adapt it:
- Use concurrency for what your database or worker pool can handle at once. Use throttle for a published rate limit, with 10-20% headroom.
- Add a
key(for example,event.data.tenantId) if one tenant's import shouldn't starve another's. - Runs over the limit queue and start as capacity frees up, so the job still finishes, just at a pace your systems can take. For the full set of flow control options, see Flash sales and bursty workflows.
§Primitive 4: Waits
You need this if: Part of the job depends on something outside it, like a human approval, a webhook, or another system finishing, and today a worker sits idle until it arrives.
Without waits: A paused job is still a running job. It holds a thread, a connection, or a worker slot for as long as it waits, and enough of them stall the queue behind them.
01import { inngest } from "./client";0203export const reviewImport = inngest.createFunction(04 {05 id: "review-import",06 triggers: [{ event: "records/import.finished" }],07 },08 async ({ event, step }) => {09 const report = await step.run("build-report", async () => {10 return await buildImportReport(event.data.importId); // returns { id, ... }11 });1213 // Holds state, not compute. Returns null if the timeout passes first.14 const approval = await step.waitForEvent("wait-for-approval", {15 event: "records/import.approved",16 timeout: "3d",17 match: "data.importId",18 });1920 if (!approval) {21 await step.run("escalate", async () => {22 await escalate(report.id);23 });24 return;25 }2627 await step.run("publish-import", async () => {28 await publishImport(event.data.importId);29 });30 }31);Adapt it:
matchpairs the waiting run with the right event. Use a field both events share, like an import or order ID.- Always handle the timeout (
null) case. If you don't, the code after the wait runs as if the approval arrived. - Use
step.sleeporstep.sleepUntilinstead when the wait is for a time, not an event.
§How long can a job run?
Keep each step short and let the run be long. A step can run for up to 2 hours (your hosting platform may cap it lower). step.sleep can pause a run for up to a year (30 days on the Free plan). A function run can last up to 366 days, depending on plan. See Usage limits for the current numbers.