Workers and CronJobs Beyond the API
Separating collection, forecasting, and push jobs from requests while preventing duplicate runs
The Hangangjari backend runs user-facing APIs separately from background workers. They share code, but their jobs are different.
The API has to answer immediately when a user opens the app. Public-data collection, forecast calculation, and push candidate selection can be slow or fail and need different intervals, so separate workers handle them.
Early on, I also considered adding a scheduler to one backend process. In the operating path, the API had to answer short requests immediately, while workers had to retry failed work and leave run records.
Workers are split to contain failure. If parking status checks briefly fail, the API returns the last checked value. If the push outbox backs up, the home-summary response should remain unaffected.
Work Prepared Outside the API
There are currently four always-running worker types. They prepare values before any user waits on an API call.
flowchart TB Scheduler["Per-worker APScheduler"] --> Parking["worker-parking\nStatus polling"] Scheduler --> Outing["worker-outing\nEvents / notices / city context / translation"] Scheduler --> Forecast["worker-forecast\nForecast generation / first-screen warming / metric rollup"] Scheduler --> Push["worker-push\nCandidates / preference indexes / outbox / maintenance"] Parking --> PG["Postgres"] Parking --> Redis["Redis"] Outing --> PG Forecast --> PG Forecast --> Redis Push --> PG Push --> Redis Push --> APNs["APNs"]
Occasional reference-data alignment and validation jobs are treated separately as CronJobs. If constant polling and reference-data sync live in the same picture, it becomes unclear which jobs run frequently and which run rarely.
flowchart LR Cron["K8s CronJob"] --> Master["Parking reference-data sync"] Cron --> Facility["Facility sync"] Cron --> Transit["Transit dataset sync"] Cron --> Backtest["Forecast backtest"] Master --> PG Facility --> PG Transit --> PG Backtest --> PG
One detail can be confusing. The name worker-parking may suggest that master sync and status polling both run in the always-on scheduler. In the current code, the always-on scheduler registers status polling. Parking master sync is separated into a job entrypoint and K8s CronJob schedule.
Separate Intervals and Responsibilities
In code, work is split like this.
| Worker | Main jobs |
|---|---|
| worker-parking | status polling |
| worker-outing | event sync, notice sync, realtime context sync, translation sync |
| worker-forecast | forecast generation, home-summary precompute, metric rollup |
| worker-push | candidate build, preference index rebuild, outbox drain, maintenance |
| K8s CronJob | reference-data sync, facility sync, dataset sync, backtest |
This keeps one failure from pulling down other screens. If parking collection briefly fails, the API can return the last checked state. If push delivery is delayed, the parking screen’s read path stays separate.
Worker names provide the first routing clue during incident response. An operator checks freshness for unstable parking collection and queues plus delivery records for a notification backlog.
Preventing Duplicate Runs
A simplified parking status check looks like this.
sequenceDiagram
autonumber
participant Scheduler as Scheduler
participant Job as Status polling job
participant Lock as DB lock
participant Source as Parking API
participant PG as Postgres
participant Redis as Status cache
participant Fact as Notification fact
Scheduler->>Job: Run at scheduled interval
Job->>Lock: Try duplicate-run lock
alt Lock acquired
Job->>Source: Request source parking status
Source-->>Job: Return source rows
Job->>PG: Store status snapshot and collection record
Job->>Redis: Refresh latest parking status cache
Job->>Fact: Record notification candidate facts if changed
else Duplicate execution
Job-->>Scheduler: Skip this run
end
Workers use scheduler settings, locks, and run records to detect overlapping execution. If the same job runs twice at once, duplicate rows, duplicate notifications, and cache races can happen.
Storage-level locks prevent the same job from being processed concurrently.
Refresh Cadence by Data Source
Hangangjari has a schedule catalog for data sources. It records the source, owning job category, refresh interval, and stale threshold.
Operators use this catalog to check each source’s interval and stale threshold. API freshness display and operational alerts use the same criteria.
A value checked every 30 seconds and one checked once a day need different freshness copy. The catalog keeps that distinction by grouping collection intervals and stale thresholds in one place.
flowchart LR Catalog["Collection schedule catalog"] --> SchedulerJobs["Always-on scheduler jobs"] Catalog --> CronJobs["Operational CronJobs"] SchedulerJobs --> Ingestion["Ingestion use cases"] CronJobs --> Ingestion Ingestion --> Runs["ingestion_runs\nsuccess / row_count / schema_hash"] Runs --> Health["Source status and freshness"] Health --> API["Outing/home-summary freshness"]
Shared Rules for Background Jobs
- Do not put ingestion in the API startup hook.
- Make jobs idempotent where possible.
- Record success, failure, row count, and schema hash for each source.
- Prevent duplicate execution with scheduler settings and storage locks.
- Track push delivery through outbox claims and delivery attempts.
- Source failures should appear as freshness for that area, not as a whole-app outage.
Results of Splitting the API and Workers
The API returns the last checked value quickly, while workers collect new values on source-specific schedules. Separate execution lifecycles made collection failure and user-request latency independently traceable.
Comments
No comments yet. Be the first to leave one.
Pending review