Short answer: schedule one durable cleanup run, feed bounded PostgreSQL delete batches through a queue, and admit each batch only while the weekly fintech digest remains inside its latency budget. The worker must be idempotent because both the cron trigger and queue delivery can repeat; the database, not the scheduler or broker, should decide whether an old log row is still eligible.
This design deliberately optimizes for foreground latency before cleanup speed. Retention work can finish later. A customer digest cannot recover the hour in which it was supposed to arrive.
Calculate the drain window before writing the trigger
The useful starting question is not which cron package to install. It is whether the cleanup service rate can exceed the arrival rate of expiring logs without taking database capacity needed by the digest path. In this fintech system, active customers receive a weekly digest, while delivery logs age toward a retention boundary. Those workloads share PostgreSQL but have different deadlines, so treating them as equal consumers is already a design error.
Estimate eligible rows per week, observed rows removed per batch, safe batches per minute during quiet periods, and the hours in which digest latency leaves no cleanup capacity. Those values produce a drain window. If the projected window reaches the next retention boundary, the design needs a higher safe service rate through batch tuning, measured concurrency, or a different physical layout; adding a queue without this arithmetic merely gives the backlog somewhere else to wait.
Then define an admission rule from application measurements: when digest latency or database contention crosses the team's agreed limit, stop publishing new cleanup batches; allow the current short transaction to finish, and resume later. This exposes the real latency-versus-cost trade-off. A single slow consumer costs less operationally and applies predictable pressure, while extra consumers shorten drain time by spending more compute and database capacity.
Capture eligible_before once when the scheduled run is created. Do not recalculate it in every worker. A moving cutoff lets a long run quietly expand its own scope, which makes counts difficult to reconcile and can pull newly eligible rows into a run whose safety checks were made against an earlier state. Give the run a stable run_id, a logical schedule period, the captured cutoff, and progress timestamps. A unique constraint on the logical period makes a repeated cron trigger point to the same unit of intent.
Cron is allowed to be boring.
The trigger should create or find that run and enqueue its identifier. It should not execute an unbounded DELETE, calculate thousands of row identifiers in application memory, or assume that firing once means execution once. GitHub's scheduling documentation offers a useful reminder even if that service is not used here: scheduled work can be delayed under high load, and some queued jobs may be dropped. The general architectural lesson is narrower than any product choice — a timer requests work, while durable state proves what work exists.
How can Node.js implement cron-triggered idempotent PostgreSQL cleanup batches?
Use a small message containing the stable run_id. A worker loads the captured boundary, begins a transaction, selects no more than the configured batch size from rows that are still eligible, deletes that claimed set, and commits. Only after commit should the delivery be acknowledged. If the batch was full, another message can continue the run; if it was partial, the worker can mark the run drained after checking the same durable state.
The application may be Node.js, but the contract is language-independent. The following Python keeps that contract visible without pretending that a framework creates correctness:
from dataclasses import dataclass
from datetime import datetime
from typing import Protocol
@dataclass(frozen=True)
class CleanupRun:
run_id: str
eligible_before: datetime
batch_size: int
class CleanupStore(Protocol):
def delete_eligible_batch(self, run: CleanupRun) -> int:
"""Claim and delete at most one batch in a single transaction."""
...
def mark_drained(self, run_id: str) -> None:
...
class Queue(Protocol):
def publish(self, run_id: str) -> None:
...
def process_batch(run: CleanupRun, store: CleanupStore, queue: Queue) -> None:
deleted = store.delete_eligible_batch(run)
if deleted == run.batch_size:
queue.publish(run.run_id)
else:
store.mark_drained(run.run_id)
The critical method is delete_eligible_batch, not process_batch. Its selection and deletion must share one transaction and one eligibility predicate. That predicate needs every preservation rule that can veto deletion, including legal holds or an audit flag; checking such a rule before the transaction opens creates a race. Foreign keys also belong in the design review. A repeatable worker does not make an undeclared cascade policy safe.
There is a subtle edge at the end of a run: a full batch says only that more rows may exist. Publishing another message is harmless because the next transaction re-evaluates eligibility. A partial batch says the current query found fewer rows than its limit, but completion still deserves a durable transition rather than an inference from an empty queue. Queues are delivery mechanisms, not ledgers.
Keep acknowledgements on the far side of commit. RabbitMQ's acknowledgement documentation explains that unacknowledged deliveries are automatically requeued when their channel or connection closes. That behavior supports at-least-once delivery. It does not make the database transaction exactly-once, and it cannot tell whether a process lost its connection one instruction after commit.
So design for the replay.
If a committed message arrives again, the old rows from that transaction no longer satisfy the query because they no longer exist; the worker takes the next eligible batch or observes that the run is drained. Avoid side effects that are not covered by this invariant. If removing a log also requires deleting an object from separate storage, PostgreSQL cannot atomically commit both systems. Record durable cleanup intent, make the external deletion safe to repeat, and mark the intent complete only after the two states have converged.
Keep retention authority in durable state
The word "idempotent" is too often attached to a function name and left there. A testable governance invariant is stricter: replaying a trigger or delivery must never broaden the captured retention scope, delete a preserved row, or lose track of a run that still has eligible work. The ledger should retain who or what created the run, the logical period, the immutable cutoff, its latest progress time, and its terminal state, so operators can distinguish authorized deletion from transport activity.
Walk the system across boundaries where ownership changes. A trigger can be replayed before publication; the unique logical-period constraint returns the existing run. A worker can stop before commit; PostgreSQL rolls the transaction back, and the delivery remains available for retry. Commit can succeed before acknowledgement; redelivery advances the same run under the same predicate. Publication of the continuation can fail after a batch is acknowledged; a reconciler must find a non-drained run with stale progress and publish its run_id again. That last case is why chaining messages alone is insufficient. The run ledger, not queue depth, says whether retention has converged.
A malformed message is different from a replay. Validate its shape, require the referenced run to exist, classify terminal application errors such as RUN_NOT_FOUND separately from retryable delivery interruption, and route exhausted messages to an operator-visible quarantine. Retrying invalid data forever increases queue age without increasing safety. Do not use a broker redelivery count as a substitute for run progress, because one counts transport attempts and the other describes durable business state.
Testing should deliberately interrupt the worker immediately before commit, immediately after commit, and before acknowledgement. Then replay the same run and compare its final eligible-row count with a clean execution. Also test concurrent workers against rows with preservation flags changing inside transactions. The important assertion is not that every message runs once; it is that every committed delete obeys the same captured boundary and preservation predicate.
I'm not sure what batch size is safe for an unseen production schema. Nobody should be. Row width, index maintenance, query plans, foreign keys, storage throughput, and the digest workload all change the answer. Start with a deliberately small configured limit, observe it under representative traffic, and alter one control at a time.
Measure the digest budget, then stage the rollout
Evaluation begins with the least expensive mechanism that could meet both deadlines. Queues add a worker, a ledger, reconciliation, and more observability. Those costs are justified only when a direct scheduled statement cannot stay within the foreground latency budget or finish retention work by its own deadline.
| Design | Foreground latency risk | Operating cost | Recovery boundary | Use it when |
|---|---|---|---|---|
| One scheduled delete | Concentrated in one statement | Lowest | Rerun the statement under a verified predicate | The table is small and load tests show no digest impact |
| One queued batch worker | Bounded and easy to pause | Moderate | Ledger plus redelivery resumes progress | Cleanup may yield to the weekly digest |
| Several batch workers | Higher lock and I/O pressure | Higher | Claims must remain transactionally isolated | One worker misses a measured retention deadline |
| Time partition removal | Low deletion work when boundaries align | Higher schema and lifecycle complexity | Partition state becomes the unit of recovery | Retention and preservation rules match whole partitions |
For the weekly digest, one queued worker is the conservative baseline when a direct delete has already failed the latency test. The catch is that this design is not suitable when the expiring volume grows faster than a serial consumer can remove it before the next deadline. Add limited concurrency only after measuring that gap; if preservation rules align with complete time ranges, partitioning may deserve evaluation instead. Conversely, stick with one scheduled delete when representative tests prove the table is small enough, because operating a queue and reconciler for a trivial workload is wasted cost.
Observe digest p95 and p99 latency beside cleanup transaction duration, rows deleted per batch, age of the oldest eligible row, queue age, redelivery count, and stale run count. The oldest eligible row shows whether retention is converging. A worker can report perfect health while falling farther behind every week.
Roll out compactly: first create read-only runs that count candidates under the exact cutoff and preservation predicate; next enable one worker for a narrow historical slice; then replay triggers and deliveries while testing each failure boundary; finally overlap cleanup with a digest run at the lowest concurrency. The kill switch should prevent new batches from being published while allowing the active database transaction to resolve. Change batch size or concurrency, never both in the same experiment — otherwise the latency result cannot identify its cause.
Your mileage may vary, but the decision rule should not: protect the customer deadline, prove retention convergence, and buy faster deletion only when a measured deadline requires it.
References
- RabbitMQ, "Consumer Acknowledgements and Publisher Confirms": https://www.rabbitmq.com/docs/confirms
- GitHub Docs, "Events that trigger workflows": https://docs.github.com/en/actions/using-workflows/events-that-trigger-workflows
Further reading
The acknowledgement and scheduling references above are the primary sources for the delivery and trigger behavior discussed here.
Top comments (0)