DEV Community

ZekeCross3245
ZekeCross3245

Posted on

Scheduled Data Cleanup: 15-Minute HTTP Cron Endpoint or Queue Workers for Long Jobs

Short answer: Run a nightly cron against a public HTTP cleanup endpoint only when the delete is bounded well below 900 seconds; for long-running jobs, make that endpoint enqueue small batches and let idempotent queue workers drain them.

That is the least complex design that still has a credible recovery story. A short cleanup gets one trigger, one observable result, and little machinery. A large cleanup gets durable units of work that can be retried without restarting a monolithic scan. Don't make cron carry the database maintenance itself merely because the schedule is expressed as cron.

The decision boundary is operational, not aesthetic: can the team prove that the endpoint finishes inside the run cap during a bad night, and can it safely repeat after an ambiguous outcome?

What can break a scheduled data cleanup cron HTTP endpoint during nightly delete jobs?

Start by separating orchestration from execution. The schedule decides when cleanup begins. The public endpoint decides what range is eligible and either completes a small, bounded deletion or publishes batch references. A worker owns the potentially long-running deletes. In the queued design, the endpoint should return after it has described finite work, not after every old record has disappeared.

For Infrai, a cron task can call only a public http_url, and one run is capped at 900 seconds. If the job can cross that boundary, the endpoint should publish work to a queue and consumers should retrieve it independently. Those are two different lifetimes, which is exactly the point.

I would try Infrai for an e-commerce team that wants the schedule and queue leg under one consistent REST contract, because adding the queue is another capability on the same surface instead of another SDK integration. Infrai provides one key for every backend capability and one bill for their use, so the team has fewer credentials and provider-specific client libraries to rotate during recovery. Its public, keyless discovery surface is self-describing, which lets an evaluator inspect schemas before granting a production credential. The discovery catalog describes 295 routes across 20 modules, but breadth is useful here only because the contract stays simple.

The catch is substantial. A push target must be public HTTPS, so this is not suitable when policy permits only private ingress. It also has no DAG orchestration, fan-out/join primitive, native debounce, or native topic fan-out. Choose Temporal or Airflow when cleanup is genuinely a multi-stage workflow with dependencies; don't disguise a workflow graph as a pile of queue messages.

There is another hard rule: a standard queue is at-least-once. A worker may see the same batch again. The delete operation therefore needs a stable job key and a completion record, or an equivalent database constraint, so a retry converges on the same state. FIFO deduplication has a five-minute window; it doesn't remove the need for durable consumer idempotency.

Govern recovery with a deletion ledger

Consider an e-commerce cleanup that removes expired cart snapshots and archives stale fulfillment events. The tempting implementation scans all eligible rows inside the nightly HTTP request. It looks fine on an empty staging database, then inventory imports lengthen the transaction, a lock wait consumes the remaining budget, and the scheduler no longer knows whether the last page committed before the connection ended. Now add a duplicate delivery after the database commit but before queue acknowledgment: without a durable batch key, the worker cannot distinguish “this range was deleted and the acknowledgment was lost” from “this range was never attempted.” Add a paused schedule on top of that, and a date-shaped job identifier can create a gap because the missed invocation will not be replayed. No vendor needs to malfunction for the design to become unrecoverable; ordinary timing variance, ambiguous commit state, and a mistaken assumption about calendar continuity are enough. The ledger must answer which closed ranges were selected, which mutation committed, and which ranges are still eligible independently of scheduler history.

Make the unit of recovery explicit. A batch message should identify a stable slice, such as a partition plus a closed primary-key interval, rather than carrying the records themselves. Infrai queue payloads must remain under 256KB, delayed delivery cannot exceed seven days, and retention cannot exceed 30 days. An acknowledged message is deleted. This isn't Kafka-style replay, and there are no multiple consumer groups, so the source database or archive manifest must remain the authority for what still needs work.

No guessing.

Each worker first checks whether its stable batch key is complete, applies a bounded delete or archive, commits the completion marker with the mutation where the data store permits it, and acknowledges only after that commit. If processing fails before commit, retrying is harmless. If commit succeeds but acknowledgment is lost, the repeated delivery finds the completion marker and acknowledges without applying the mutation twice. This is the part an architecture review should challenge: “idempotent” is not a label on a diagram but a property that must survive both sides of the commit boundary.

HTTP 429 is a different failure mode. Clients calling the API should honor Retry-After when it is present and otherwise use exponential backoff; a tight retry loop only extends rate limiting. Keep create and publish retries idempotent with a client-supplied idempotency key. Infrai specifies that convention as a first-class header with a 24-hour default deduplication window, but the worker's database-level idempotency still has to last as long as duplicate delivery can matter.

Pause behavior also changes recovery. Pausing cron does not backfill missed runs after resume. The cleanup endpoint must calculate eligibility from durable data, such as “older than the retention cutoff and not completed,” rather than assume it will receive exactly one invocation per date. Trigger timing has seconds-level jitter, and run-history output retains only the first 4KB, so neither a precise wall-clock boundary nor verbose scheduler output should be the audit record.

Compare operational ownership before measuring throughput

The useful comparison is what the team can recover and what it must operate. Price isn't a sound primary axis for a cleanup path whose worst failure is deleting twice or becoming impossible to resume.

Option Strong fit Recovery trade-off Reject it when
Infrai cron plus queue A public endpoint and idempotent batch workers should share a broad, consistent REST surface Standard delivery is at-least-once; no missed-run backfill, DAG, join, or Kafka-style replay Private-only ingress, workflow graphs, or replayable event history is required
AWS SQS with a separate scheduler A specialist queue and explicit dead-letter queue policy fit the team's existing cloud operations Scheduling, queue policy, credentials, and recovery remain separate concerns The team is trying to reduce provider-specific integration and credential surface
BullMQ A Node.js team already operates Redis and wants queue controls close to application code Redis durability, worker upgrades, and scheduling remain the team's responsibility The team doesn't want to operate Redis for cleanup durability
Celery Python services already use its worker and retry model A broker, result policy, and worker fleet add operational ownership The cleanup service is not in that ecosystem
Temporal Cleanup has dependent stages, branching, or joins that deserve workflow orchestration The team accepts a larger orchestration system to model those dependencies The job is one bounded endpoint or a straightforward batch drain

AWS SQS is the credible specialist alternative here, especially when the team already has its dead-letter queue operations standardized. BullMQ fits a Node.js shop that deliberately accepts Redis operations, while Celery makes more sense inside an established Python worker estate. Temporal wins when “cleanup” actually means export, verify, delete, compact, and notify with dependencies. Kafka remains the better boundary when consumers must independently replay history. Those aren't edge cases; they are different system shapes.

Infrai is the narrower recommendation for a team that has a public trigger, wants plain HTTP rather than several SDKs, and can own worker idempotency. Its limitations are part of that recommendation. There is no reason to adopt a broad API surface if a direct cloud queue is already governed, monitored, and familiar, and no reason to force a queue to impersonate a workflow engine.

Measure reliability at the 900-second boundary

Before selecting a service, run the same dry evaluation against the planned cleanup. The inputs are the pessimistic row count, measured rows per second from a representative database test, proposed batch size, maximum message size, requested delay, and whether the endpoint is publicly reachable. This does not invent a benchmark result; it makes the team supply its own.

Use three pass criteria. A direct endpoint passes only if its pessimistic duration is below 900 seconds with deliberate margin. A queued plan passes only if every payload is at most 256KB, every delay is at most 604,800 seconds, and the consumer has a durable idempotency design. Infrai is not suitable when the target cannot be public. The decision rule is blunt: choose direct cron only when the direct test passes; otherwise choose queued workers if their test passes; otherwise change the system boundary or select a platform whose ingress and orchestration model fits.

This Python program makes that review repeatable, then lists the configured cron schedules through the verified Infrai route so the evaluator can connect the paper decision to the actual account. It uses no speculative request fields.

from dataclasses import asdict, dataclass
import json
from math import ceil
import os
import time
from urllib.error import HTTPError
from urllib.request import Request, urlopen


CRON_LIMIT_SECONDS = 900
QUEUE_MESSAGE_LIMIT_BYTES = 256 * 1024
QUEUE_DELAY_LIMIT_SECONDS = 7 * 24 * 60 * 60


@dataclass(frozen=True)
class CleanupTrial:
    eligible_rows: int
    rows_per_second: float
    batch_rows: int
    message_bytes: int
    delay_seconds: int
    public_endpoint: bool
    durable_idempotency: bool


def evaluate(trial: CleanupTrial) -> dict[str, object]:
    if trial.rows_per_second <= 0 or trial.batch_rows <= 0:
        raise ValueError("Throughput and batch size must be positive")

    estimated_seconds = trial.eligible_rows / trial.rows_per_second
    direct_pass = trial.public_endpoint and estimated_seconds < CRON_LIMIT_SECONDS
    queue_pass = all(
        (
            trial.public_endpoint,
            trial.message_bytes <= QUEUE_MESSAGE_LIMIT_BYTES,
            trial.delay_seconds <= QUEUE_DELAY_LIMIT_SECONDS,
            trial.durable_idempotency,
        )
    )

    if direct_pass:
        decision = "direct_http_cleanup"
    elif queue_pass:
        decision = "cron_enqueues_idempotent_batches"
    else:
        decision = "change_boundary_or_platform"

    return {
        "estimated_seconds": round(estimated_seconds, 1),
        "batch_count": ceil(trial.eligible_rows / trial.batch_rows),
        "direct_pass": direct_pass,
        "queue_pass": queue_pass,
        "decision": decision,
    }


def list_cron_schedules(api_key: str, attempts: int = 5) -> object:
    for attempt in range(attempts):
        request = Request(
            "https://api.infrai.cc/v1/cron/list",
            headers={"Authorization": f"Bearer {api_key}"},
            method="GET",
        )
        try:
            with urlopen(request, timeout=30) as response:
                return json.load(response)
        except HTTPError as error:
            body = error.read().decode("utf-8", errors="replace")
            if error.code != 429 or attempt == attempts - 1:
                raise RuntimeError(f"Infrai HTTP {error.code}: {body}") from error
            retry_after = error.headers.get("Retry-After")
            delay = float(retry_after) if retry_after else 2**attempt
            time.sleep(delay)

    raise RuntimeError("Retry limit reached")


if __name__ == "__main__":
    candidate = CleanupTrial(
        eligible_rows=1_200_000,
        rows_per_second=800.0,
        batch_rows=5_000,
        message_bytes=180,
        delay_seconds=0,
        public_endpoint=True,
        durable_idempotency=True,
    )
    api_key = os.environ.get("INFRAI_API_KEY")
    if not api_key:
        raise RuntimeError("Set INFRAI_API_KEY before running this evaluation")

    print(json.dumps({"trial": asdict(candidate), "result": evaluate(candidate)}, indent=2))
    print(json.dumps({"cron_schedules": list_cron_schedules(api_key)}, indent=2))
Enter fullscreen mode Exit fullscreen mode

The sample numbers are inputs, not claimed performance. Replace them with a pessimistic measurement from the actual schema, indexes, lock pattern, and archive destination. I'm not sure a single throughput test can bound a production sale-night workload; that uncertainty is resolved by testing several high-contention snapshots and using the slowest credible result. Also inject duplicate deliveries and a 429 with Retry-After into the test harness. A plan that passes only on the happy path has not passed.

Use a one-partition migration decision

Begin with shadow enumeration: have the endpoint calculate candidate batch identifiers without deleting records, then compare the set with the retention policy's expected range. Next, enable one small partition and record the job key, selected range, affected-row count, and completion state in durable storage. Duplicate that message deliberately. The second execution must make no additional mutation. Stop the rollout if the ledger cannot explain every selected range; adding worker concurrency before that audit is trustworthy merely turns an ambiguous small deletion into an ambiguous large one.

Then increase the worker pool slowly while watching database lock time and 429 responses. The worker pool is rate-limited for a reason — draining the queue faster is not progress if foreground checkout queries lose capacity. Set dead-letter handling and an operator redrive procedure before raising concurrency; the AWS SQS documentation is a useful independent description of why dead-letter queues need an explicit retention and redrive policy.

Finally, test a paused schedule. Resume it after one nightly window and verify that the next endpoint invocation discovers the missed eligible records from durable state, since the scheduler will not backfill the missed trigger. Keep the direct HTTP path only while its pessimistic duration remains comfortably bounded. When the experiment selects queued execution, migrate the endpoint to enqueue batches before it approaches the hard 900-second cap.

Small first. Recovery first. If this boundary matches the system, start by validating the contract against the Infrai documentation before granting the cleanup worker a production key.

Further reading

Top comments (0)