flexmeasures.data.services.job_map
Index RQ jobs by asset or sensor, and look them up again.
Classes
- class flexmeasures.data.services.job_map.JobMap(connection: Redis)
Map assets or sensors and queues to RQ job IDs in Redis.
Job creation code adds IDs to Redis sets. Listing code uses those IDs to retrieve jobs from RQ, which stores the job records separately. An expired job’s ID remains in the map until a read removes it.
- Each index key contains the queue, entity type, and entity ID:
forecasting:sensor:1 (forecasting jobs can be stored by sensor only)
scheduling:sensor:2
scheduling:asset:3
get() fetches all matching jobs. get_enqueued_at() reads only the fields needed to check that each job exists and sort it by enqueue time. This lets a paginated listing fetch full records only for the requested page with fetch_jobs().
- __init__(connection: Redis)
- fetch_jobs(job_ids: list[str]) list[Job | None]
Fetch the given jobs, in order, with None for each job that has expired meanwhile.
- get(asset_or_sensor_id: int, queue: str, asset_or_sensor_type: str) list[Job]
Fetch the current jobs from the Redis ID index.
- get_enqueued_at(asset_or_sensor_id: int, queue: str, asset_or_sensor_type: str) list[tuple[str, datetime | None]]
List the current job IDs from the Redis ID index, each with the time its job was enqueued.
Only these fields are read, rather than the full jobs, so that all jobs can be sorted cheaply. A job that was created but not enqueued yet (e.g. one waiting on another job) comes with None.
Exceptions
- exception flexmeasures.data.services.job_map.NoRedisConfigured(message='Redis not configured')