Celery workers, queues & beat¶
Databús runs three Celery processes: two workers consuming different queues, and one beat scheduler. This page describes their roles, queue routing, and the beat schedule.
Worker topology¶
┌─────────────────────────────────────┐ ┌─────────────────────────────────────┐
│ realtime-engine (worker) │ │ schedule-engine (worker) │
│ queue: realtime_engine │ │ queue: schedule_engine │
│ │ │ │
│ Tasks: │ │ Tasks: │
│ • process_position_update │ │ • build_vehicle_positions │
│ • run_lifecycle_event │ │ • build_trip_updates │
│ • scan_stale_runs │ │ • build_alerts (stub, unscheduled) │
│ • fetch_positions │ │ • build_schedule │
│ │ │ │
│ Bootstep: MQTTConsumerStep │ │ No bootstep active │
│ (gated by MQTT_CONSUMER_ENABLED) │ │ │
└─────────────────────────────────────┘ └─────────────────────────────────────┘
┌─────────────────────────────────────┐
│ scheduler (Celery beat) │
│ Fires periodic tasks on schedule │
│ No task execution, no queues │
└─────────────────────────────────────┘
All three processes load backend/databus/celery.py as the Celery app. The queue each task lands on is declared via the queue= argument to @shared_task.
Queue routing¶
| Queue | Consumed by | Tasks |
|---|---|---|
realtime_engine |
realtime-engine worker |
process_position_update, run_lifecycle_event, scan_stale_runs, fetch_positions |
schedule_engine |
schedule-engine worker |
build_vehicle_positions, build_trip_updates, build_alerts, build_schedule |
Queue separation ensures that a spike in MQTT telemetry (many process_position_update tasks) does not starve GTFS-RT building tasks, and vice versa.
The MQTT consumer bootstep¶
The MQTTConsumerStep bootstep is registered globally in databus/celery.py:
This registers the step on every worker that loads the Celery app. The step checks MQTT_CONSUMER_ENABLED at startup and only activates on the realtime-engine worker. Other workers log a single INFO line and skip activation.
See Telemetry ingestion for the full bootstep description.
Beat schedule¶
The beat schedule is defined in code in backend/databus/celery.py, not in the django_celery_beat database admin.
Beat schedule is not in the Django admin
AGENTS.md suggests that schedules are managed through django_celery_beat and the Django admin. This is incorrect. The schedule is static, defined in app.conf.beat_schedule in databus/celery.py. To change a cadence, edit that file and redeploy the scheduler service.
app.conf.beat_schedule = {
"build-vehicle-positions-every-15s": {
"task": "schedule_engine.tasks.build_vehicle_positions",
"schedule": timedelta(seconds=15),
},
"build-trip-updates-every-15s": {
"task": "schedule_engine.tasks.build_trip_updates",
"schedule": timedelta(seconds=15),
},
"scan-stale-runs-every-30s": {
"task": "realtime_engine.tasks.scan_stale_runs",
"schedule": timedelta(seconds=30),
},
"fetch-positions": {
"task": "realtime_engine.tasks.fetch_positions",
"schedule": timedelta(seconds=10),
# A task that couldn't even start within its own 10s cycle is stale
# by the time a worker slot frees up -- revoke it instead of letting
# queued fetch_positions runs pile up behind a slow/unreachable
# source (see fetch_positions' soft_time_limit for the in-flight
# bound on runs that DO start).
"options": {"expires": 10},
},
"build-schedule-daily": {
"task": "schedule_engine.tasks.build_schedule",
"schedule": timedelta(days=1),
},
}
| Entry | Task | Queue | Cadence | Purpose |
|---|---|---|---|---|
build-vehicle-positions-every-15s |
build_vehicle_positions |
schedule_engine |
15 s | Rebuild vehicle_positions.{pb,json} |
build-trip-updates-every-15s |
build_trip_updates |
schedule_engine |
15 s | Rebuild trip_updates.{pb,json}, push WebSocket heartbeat |
scan-stale-runs-every-30s |
scan_stale_runs |
realtime_engine |
30 s | Detect run_tracking_lost / run_tracking_expired |
fetch-positions |
fetch_positions |
realtime_engine |
10 s (expires=10) |
Poll ACTIVE HTTP-source sensors for in-service vehicles and publish their readings over MQTT. Gated on vehicle:<id>:current_run (the same in-service test the MQTT consumer uses), not on runs:in_progress — see the docstring in backend/realtime_engine/tasks.py. A cycle that can't even start within its own 10 s window is revoked (expires=10) instead of queuing behind a slow/unreachable source. |
build-schedule-daily |
build_schedule |
schedule_engine |
1 day | Export the current GTFS Schedule to a zip via feed.schedule.exporter.publish_gtfs_zip |
build_alerts exists but is not scheduled
schedule_engine.tasks.build_alerts is a real Celery task (queue schedule_engine) but is deliberately not in app.conf.beat_schedule. Its docstring says why: it's a placeholder that returns a fixed string and writes no feed file — the ServiceAlert feed builder isn't implemented yet. It can still be invoked manually (e.g. from the Django admin or a shell), but nothing fires it periodically.
Celery task monitoring: Flower¶
The task-monitoring compose service runs Flower, a real-time Celery dashboard.
| Environment | URL |
|---|---|
| Development | http://localhost:5555 |
| Production | https://${FLOWER_DOMAIN} (e.g. https://tasks.databus.simovilab.com) |
Flower connects to RabbitMQ and shows:
- Active, scheduled, and reserved tasks per worker.
- Task history with arguments, result, and execution time.
- Worker status and queue depths.
Compose service → Celery process mapping¶
| Compose service | Celery command | Build target |
|---|---|---|
realtime-engine |
celery -A databus worker -Q realtime_engine --loglevel=info |
realtime-engine |
schedule-engine |
celery -A databus worker -Q schedule_engine --loglevel=info |
schedule-engine |
scheduler |
celery -A databus beat --loglevel=info |
scheduler |
Build targets are defined in backend/Dockerfile. Each target installs the same Python environment but sets a different CMD.
Related pages¶
- Telemetry ingestion — the
MQTTConsumerStepbootstep. - Server-side processing —
process_position_updatein detail. - GTFS Realtime publishing — the
schedule_enginetasks. - Stale-run scanning —
scan_stale_runslogic. - Configuration & env vars —
MQTT_CONSUMER_ENABLEDand related settings.