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.
- get_many(entries: Iterable[tuple[int, str, str]]) dict[tuple[int, str, str], list[Job]]
Fetch the jobs indexed under several assets or sensors, in a fixed number of Redis round trips.
Each entry is an asset or sensor id, a queue and an asset or sensor type, as passed to get(). However many entries and jobs there are, this pings Redis once, reads every index in one pipeline and every job in another, so that a page listing the jobs of many sensors does not pay a round trip per sensor. A job indexed under several entries is fetched once. An ID whose job can no longer be found, because its RQ job expired, is removed from its index.
Exceptions
- exception flexmeasures.data.services.job_map.NoRedisConfigured(message='Redis not configured')