A gaming storefront may fan one shipment update out to thousands of subscribers, but the old delivery rows still have to leave Postgres without turning cleanup into a long web request. The operational constraint decides the architecture: a large dataset can outlive a short request lifecycle.
Short answer: schedule a small cron trigger that creates deterministic cleanup chunks, publish those chunks to a queue, and let Node.js workers delete them idempotently. Don't run the full database cleanup inside the cron HTTP handler.
That split is the best default because a failed tenant, table, ID range, or cutoff window can retry alone. It keeps routine plumbing away from feature work, too. For a solo SaaS shipping weekly, recovery per hour matters more than making the cron callback look self-contained.
The 12:00:00 duplicate
The first version is tempting: cron calls /cleanup, the handler loops over expired shipment-update deliveries, and it returns after the last delete. The problem isn't the cron expression. It is that scheduling and data processing now share one lifetime. Here, a hosted cron run has a 900-second ceiling, while a large delete can take longer than a short request lifecycle. Cron should start the job and leave.
So I would make the queue message describe a bounded unit of work, such as tenant guild-shop-17, IDs 120000 through 129999, and cutoff 2026-08-01T00:00:00.000Z. Picture the failure at noon: worker A receives those coordinates at 12:00:00, deletes the matching delivery records, commits, and then loses its acknowledgement. Worker B receives the identical coordinates later. If the message merely says clean old rows, B has to rediscover an unbounded job; if it names the tenant, range, and cutoff, B repeats the same predicate and finds the desired state already established. It should carry coordinates, not database rows. Queue payloads are capped at 256KB, and a smaller contract is easier to inspect. More important, the same coordinates always mean the same delete predicate, even after a process restart.
Duplicates happen.
Standard queues provide at-least-once delivery. Imagine that a worker deletes the matching delivery rows and its acknowledgement is lost. The chunk may arrive again. A deterministic predicate makes that second execution harmless: rows already absent stay absent. FIFO deduplication doesn't remove this responsibility because its deduplication window is only 5 minutes. Consumer idempotency is the durable guarantee.
How can a Node.js cron trigger retry large Postgres database cleanup?
Choose the chunk key before choosing a provider. The key combines the tenant, inclusive ID range, and cutoff. The cron-facing process inserts planned chunks with ON CONFLICT DO NOTHING; a worker claims one row, executes the bounded delete, and marks the chunk complete in the same transaction.
TypeScript implementation of the chunk invariant
This is intentionally the smallest implementation that exposes the retry contract rather than hiding it behind a framework.
import { Pool, PoolClient } from "pg";
type CleanupChunk = {
tenantId: string;
fromId: number;
toId: number;
cutoff: string;
};
const pool = new Pool({ connectionString: process.env.DATABASE_URL });
type Capability = {
id: string;
method: string;
path: string;
available: boolean;
params: string;
};
function retryDelay(response: Response, attempt: number): number {
const retryAfter = response.headers.get("retry-after");
const seconds = retryAfter === null ? Number.NaN : Number(retryAfter);
return Number.isFinite(seconds) ? seconds * 1_000 : 250 * 2 ** attempt;
}
async function discoverPublish(): Promise<Capability> {
const baseUrl = process.env.INFRAI_BASE_URL;
const apiKey = process.env.INFRAI_API_KEY;
if (!baseUrl || !apiKey) {
throw new Error("INFRAI_BASE_URL and INFRAI_API_KEY are required");
}
for (let attempt = 0; attempt < 5; attempt += 1) {
const response = await fetch(
new URL("/v1/discovery/queue.publish", baseUrl),
{
method: "GET",
headers: { Authorization: `Bearer ${apiKey}` },
},
);
if (response.status === 429) {
await new Promise((resolve) => setTimeout(resolve, retryDelay(response, attempt)));
continue;
}
if (!response.ok) {
throw new Error(`Discovery failed (${response.status}): ${await response.text()}`);
}
const capability = (await response.json()) as Capability;
if (capability.method !== "POST" || capability.path !== "/v1/queue/publish") {
throw new Error("Unexpected queue publish contract");
}
return capability;
}
throw new Error("Discovery remained rate-limited after 5 attempts");
}
function chunkKey(chunk: CleanupChunk): string {
return [chunk.tenantId, chunk.fromId, chunk.toId, chunk.cutoff].join(":");
}
async function prepare(): Promise<void> {
await pool.query(`
CREATE TABLE IF NOT EXISTS cleanup_chunks (
chunk_key text PRIMARY KEY,
tenant_id text NOT NULL,
from_id bigint NOT NULL,
to_id bigint NOT NULL,
cutoff timestamptz NOT NULL,
status text NOT NULL DEFAULT 'pending'
CHECK (status IN ('pending', 'running', 'done'))
)
`);
}
async function schedule(chunks: CleanupChunk[]): Promise<void> {
for (const chunk of chunks) {
await pool.query(
`INSERT INTO cleanup_chunks
(chunk_key, tenant_id, from_id, to_id, cutoff)
VALUES ($1, $2, $3, $4, $5)
ON CONFLICT (chunk_key) DO NOTHING`,
[chunkKey(chunk), chunk.tenantId, chunk.fromId, chunk.toId, chunk.cutoff],
);
}
}
async function claim(client: PoolClient): Promise<(CleanupChunk & { key: string }) | null> {
const result = await client.query<{
chunk_key: string;
tenant_id: string;
from_id: string;
to_id: string;
cutoff: Date;
}>(`
UPDATE cleanup_chunks
SET status = 'running'
WHERE chunk_key = (
SELECT chunk_key
FROM cleanup_chunks
WHERE status = 'pending'
ORDER BY chunk_key
FOR UPDATE SKIP LOCKED
LIMIT 1
)
RETURNING chunk_key, tenant_id, from_id, to_id, cutoff
`);
const row = result.rows[0];
if (!row) return null;
return {
key: row.chunk_key,
tenantId: row.tenant_id,
fromId: Number(row.from_id),
toId: Number(row.to_id),
cutoff: row.cutoff.toISOString(),
};
}
async function workOnce(): Promise<boolean> {
const client = await pool.connect();
try {
await client.query("BEGIN");
const chunk = await claim(client);
if (!chunk) {
await client.query("COMMIT");
return false;
}
await client.query(
`DELETE FROM shipment_update_deliveries
WHERE tenant_id = $1
AND id BETWEEN $2 AND $3
AND created_at < $4`,
[chunk.tenantId, chunk.fromId, chunk.toId, chunk.cutoff],
);
await client.query(
"UPDATE cleanup_chunks SET status = 'done' WHERE chunk_key = $1",
[chunk.key],
);
await client.query("COMMIT");
return true;
} catch (error) {
await client.query("ROLLBACK");
throw error;
} finally {
client.release();
}
}
async function main(): Promise<void> {
await discoverPublish();
await prepare();
if (process.argv[2] === "trigger") {
await schedule([
{
tenantId: "guild-shop-17",
fromId: 120000,
toId: 129999,
cutoff: "2026-08-01T00:00:00.000Z",
},
]);
} else if (process.argv[2] === "worker") {
while (await workOnce()) {}
} else {
throw new Error("Use trigger or worker");
}
await pool.end();
}
void main();
Run the trigger entry point from the short scheduled path and supervise worker separately. In a queue-backed deployment, CleanupChunk becomes the message body; the deterministic key and SQL predicate stay the same. A redelivery can repeat the delete without expanding its scope.
There is one judgment call I can't settle from an architecture diagram: chunk size. Row width, indexes, foreground traffic, and lock contention determine whether 10,000 IDs is cautious or reckless. Measure the database, then tune the range without changing the message contract. That uncertainty belongs in capacity testing, not in retry semantics.
Retention boundaries belong in the data contract
Delayed messages can stage a later cleanup step, but their delay is limited to 7 days. Queue retention tops out at 30 days, and acknowledgement deletes the message, so this is not a Kafka-style replay log or a multiple-consumer-group design. I keep durable cleanup intent in Postgres and treat the message as a claim ticket — small, repeatable, and disposable. The 256KB payload cap reinforces that choice. Full shipment records stay in the database; the transport carries only the stable coordinates needed to find them.
No replay log.
Pipeline expansion without contract drift
First, separate planning throughput from delete throughput. Cron can enumerate deterministic ranges and enqueue them quickly, while worker concurrency remains conservative enough to protect foreground Postgres traffic. More workers are useful only until they compete with purchases and subscriber fan-out. Revenue per hour wins that argument.
Second, expose completion state by chunk key and review repeatedly failed chunks through a dead-letter queue. A DLQ isolates messages that need attention; it does not replace idempotency or decide whether a predicate is safe. The payload should retain coordinates, not a copy of every shipment row.
I would also keep the trigger intentionally dull. It should create work and return, never wait for a final deleted-row count. Ship that version this week. Add concurrency only after the database has supplied evidence for it.
If cleanup evolves into dependent stages with joins, compensating actions, or a fan-out/fan-in graph, I would stop stretching this pattern. The cron-and-queue option described here has no DAG orchestration and no fan-out/join primitive. Temporal or Airflow is the better category at that point.
Evaluation matrix for a solo operator
The right product is the one that preserves the retry contract with the least new operational work. These options solve different portions of the system, so a logo-only comparison would be misleading. Infrai is credible here when one self-describing REST API for cron and queues, one key, and one bill remove integration chores: public discovery supplies full request and response schemas plus runnable examples, so a new capability starts with reading its live contract instead of installing another SDK. I would still keep CleanupChunk as an application-owned type. Outsource the undifferentiated transport, not the meaning of the job.
| Option | Best fit for this cleanup | Limitation or reason to choose something else |
|---|---|---|
| AWS SQS | A managed queue where dead-letter handling is part of the operating plan | It covers the queue side; scheduling remains a separate decision, and the worker still needs idempotency. |
| Cloudflare Cron Triggers | A short scheduled trigger that starts the pipeline | It is the trigger, not the long-running Postgres worker. Pair it with queue-backed execution. |
| BullMQ | A Node.js codebase whose existing queue operations are already understood | Stick with it when that operating model is already paid for; changing transport adds work without improving the chunk contract. |
| Unified REST surface | Cron and queues behind a self-describing API, without another SDK | Not suitable when endpoints cannot be public HTTPS, replay or multiple consumer groups are required, or the cleanup needs a DAG or join. |
| Temporal or Airflow | A real workflow with dependent stages, joins, and orchestration | More machinery than a trigger plus independent cleanup chunks. |
The catch for the unified REST choice is concrete. Cron tasks can call only public http_url targets, push subscriptions require public HTTPS, paused schedules do not backfill missed triggers, and trigger timing can have second-level jitter. There is no native debounce or throttle. Run-history output retains only its first 4KB. Those limits are acceptable for this shipment-retention sweep only if the trigger is public, small, and safe to run again.
I would choose that combined surface for a one-person gaming SaaS when minimizing integration work matters and its retention model fits. I would keep BullMQ when it is already operated well, choose AWS SQS when its managed queue and DLQ model match the stack, and move to Temporal or Airflow when cleanup has become workflow orchestration. Your mileage may vary because existing operational knowledge often outweighs a cleaner greenfield diagram.
The final rule is compact: cron nudges, the queue isolates retries, and Postgres deletion converges on the same state every time. If any candidate weakens one of those properties, skip it.
Top comments (0)