Skip to content

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:

from realtime_engine.mqtt import MQTTConsumerStep
app.steps["worker"].add(MQTTConsumerStep)

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.