Coverage for /usr/local/lib/python3.12/site-packages/prefect/server/services/scheduler.py: 78%
106 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 02:04 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 02:04 +0000
1"""
2The Scheduler service.
4This service schedules flow runs from deployments with active schedules.
5"""
7from __future__ import annotations
9import datetime
10import logging
11from datetime import timedelta
12from typing import Any, Sequence
13from uuid import UUID
15import sqlalchemy as sa
16from docket import Perpetual
17from sqlalchemy.ext.asyncio import AsyncSession
19import prefect.server.models as models
20from prefect.logging import get_logger
21from prefect.server.database import PrefectDBInterface, provide_database_interface
22from prefect.server.schemas.states import StateType
23from prefect.server.services.perpetual_services import perpetual_service
24from prefect.settings.context import get_current_settings
25from prefect.types._datetime import now
26from prefect.utilities.collections import batched_iterable
28logger: logging.Logger = get_logger(__name__)
31class TryAgain(Exception):
32 """Internal control-flow exception used to retry the Scheduler's main loop"""
35def _get_select_deployments_to_schedule_query(
36 db: PrefectDBInterface,
37 deployment_batch_size: int,
38 min_runs: int,
39 min_scheduled_time: datetime.timedelta,
40) -> sa.Select[tuple[UUID]]:
41 """
42 Returns a sqlalchemy query for selecting deployments to schedule.
44 The query gets the IDs of any deployments where ANY active schedule has:
46 - EITHER:
47 - fewer than `min_runs` auto-scheduled runs for that schedule
48 - OR the max scheduled time for that schedule is less than
49 `min_scheduled_time` in the future
51 This per-schedule check ensures that high-frequency schedules get
52 re-evaluated even when other schedules on the same deployment still
53 have runs far in the future.
55 Each schedule's runs are identified via the `created_by` JSON column
56 which stores `{"id": "<schedule_id>", "type": "SCHEDULE", ...}` on
57 every auto-scheduled flow run. An expression index on
58 `(created_by->>'id')` keeps these correlated subqueries fast.
59 """
60 right_now = now("UTC")
62 # Use type_coerce to bypass the Pydantic TypeDecorator so SQLAlchemy
63 # emits a bare `created_by->>'id'` (Postgres) / `json_extract(created_by, '$.id')`
64 # (SQLite) without an extra CAST wrapper. This is required for
65 # PostgreSQL to match the expression index on `(created_by->>'id')`.
66 schedule_id_match = sa.type_coerce(db.FlowRun.created_by, sa.JSON)[
67 "id"
68 ].as_string() == sa.cast(db.DeploymentSchedule.id, sa.String)
70 per_schedule_run_count = (
71 sa.select(sa.func.count())
72 .select_from(db.FlowRun)
73 .where(
74 db.FlowRun.deployment_id == db.DeploymentSchedule.deployment_id,
75 db.FlowRun.state_type == StateType.SCHEDULED,
76 db.FlowRun.next_scheduled_start_time >= right_now,
77 # `== true` (not `.is_(True)`) so the predicate matches the
78 # `auto_scheduled = true` clause of the partial index
79 # `ix_flow_run__schedule_id_scheduler`. PostgreSQL's predicate
80 # implication prover does not treat `IS true` as implying
81 # `= true`, so `.is_(True)` silently disqualifies the index.
82 db.FlowRun.auto_scheduled == sa.true(),
83 schedule_id_match,
84 )
85 .correlate(db.DeploymentSchedule)
86 .scalar_subquery()
87 )
89 per_schedule_max_time = (
90 sa.select(sa.func.max(db.FlowRun.next_scheduled_start_time))
91 .select_from(db.FlowRun)
92 .where(
93 db.FlowRun.deployment_id == db.DeploymentSchedule.deployment_id,
94 db.FlowRun.state_type == StateType.SCHEDULED,
95 db.FlowRun.next_scheduled_start_time >= right_now,
96 db.FlowRun.auto_scheduled == sa.true(),
97 schedule_id_match,
98 )
99 .correlate(db.DeploymentSchedule)
100 .scalar_subquery()
101 )
103 any_schedule_needs_runs = (
104 sa.select(sa.literal(1))
105 .select_from(db.DeploymentSchedule)
106 .where(
107 db.DeploymentSchedule.deployment_id == db.Deployment.id,
108 db.DeploymentSchedule.active.is_(True),
109 sa.or_(
110 per_schedule_run_count < min_runs,
111 per_schedule_max_time < right_now + min_scheduled_time,
112 ),
113 )
114 .correlate(db.Deployment)
115 .exists()
116 )
118 query = (
119 sa.select(db.Deployment.id)
120 .where(
121 db.Deployment.paused.is_not(True),
122 any_schedule_needs_runs,
123 )
124 .order_by(db.Deployment.id)
125 .limit(deployment_batch_size)
126 )
127 return query
130def _get_select_recent_deployments_to_schedule_query(
131 db: PrefectDBInterface,
132 deployment_batch_size: int,
133 loop_seconds: float,
134) -> sa.Select[tuple[UUID]]:
135 """
136 Returns a sqlalchemy query for selecting recently updated deployments to schedule.
137 """
138 query = (
139 sa.select(db.Deployment.id)
140 .where(
141 sa.and_(
142 db.Deployment.paused.is_not(True),
143 # use a slightly larger window than the loop interval to pick up
144 # any deployments that were created *while* the scheduler was
145 # last running (assuming the scheduler takes less than one
146 # second to run). Scheduling is idempotent so picking up schedules
147 # multiple times is not a concern.
148 db.Deployment.updated
149 >= now("UTC") - datetime.timedelta(seconds=loop_seconds + 1),
150 (
151 # Only include deployments that have at least one
152 # active schedule.
153 sa.select(db.DeploymentSchedule.deployment_id)
154 .where(
155 sa.and_(
156 db.DeploymentSchedule.deployment_id == db.Deployment.id,
157 db.DeploymentSchedule.active.is_(True),
158 )
159 )
160 .exists()
161 ),
162 )
163 )
164 .order_by(db.Deployment.id)
165 .limit(deployment_batch_size)
166 )
167 return query
170async def _collect_flow_runs(
171 db: PrefectDBInterface,
172 session: AsyncSession,
173 deployment_ids: Sequence[UUID],
174 max_scheduled_time: datetime.timedelta,
175 min_scheduled_time: datetime.timedelta,
176 min_runs: int,
177 max_runs: int,
178) -> list[dict[str, Any]]:
179 """Collect flow runs to schedule from a list of deployments."""
180 runs_to_insert: list[dict[str, Any]] = []
181 for deployment_id in deployment_ids:
182 right_now = now("UTC")
183 # guard against erroneously configured schedules
184 try:
185 runs_to_insert.extend(
186 await models.deployments._generate_scheduled_flow_runs(
187 db,
188 session=session,
189 deployment_id=deployment_id,
190 start_time=right_now,
191 end_time=right_now + max_scheduled_time,
192 min_time=min_scheduled_time,
193 min_runs=min_runs,
194 max_runs=max_runs,
195 )
196 )
197 except Exception:
198 logger.exception(
199 f"Error scheduling deployment {deployment_id!r}.",
200 )
201 finally:
202 connection = await session.connection()
203 if connection.invalidated: 203 ↛ anywhereline 203 didn't jump anywhere: it always raised an exception.
204 # If the error we handled above was the kind of database error that
205 # causes underlying transaction to rollback and the connection to
206 # become invalidated, rollback this session.
207 await session.rollback()
208 raise TryAgain()
209 return runs_to_insert
212@perpetual_service(
213 enabled_getter=lambda: get_current_settings().server.services.scheduler.enabled,
214)
215async def schedule_deployments(
216 perpetual: Perpetual = Perpetual(
217 automatic=True,
218 every=timedelta(
219 seconds=get_current_settings().server.services.scheduler.loop_seconds
220 ),
221 ),
222) -> None:
223 """
224 Main scheduler - schedules flow runs from deployments with active schedules.
226 Schedule flow runs by:
227 - Querying for deployments with active schedules
228 - Generating the next set of flow runs based on each deployment's schedule
229 - Inserting all scheduled flow runs into the database
230 """
231 settings = get_current_settings().server.services.scheduler
232 deployment_batch_size = settings.deployment_batch_size
233 max_runs = settings.max_runs
234 min_runs = settings.min_runs
235 max_scheduled_time = settings.max_scheduled_time
236 min_scheduled_time = settings.min_scheduled_time
237 insert_batch_size = settings.insert_batch_size
239 db = provide_database_interface()
240 total_inserted_runs = 0
241 last_id = None
243 while True:
244 async with db.session_context(begin_transaction=False) as session:
245 query = _get_select_deployments_to_schedule_query(
246 db, deployment_batch_size, min_runs, min_scheduled_time
247 )
249 # use cursor based pagination
250 if last_id: 250 ↛ 251line 250 didn't jump to line 251 because the condition on line 250 was never true
251 query = query.where(db.Deployment.id > last_id)
253 result = await session.execute(query)
254 deployment_ids = result.scalars().unique().all()
256 # collect runs across all deployments
257 try:
258 runs_to_insert = await _collect_flow_runs(
259 db,
260 session=session,
261 deployment_ids=deployment_ids,
262 max_scheduled_time=max_scheduled_time,
263 min_scheduled_time=min_scheduled_time,
264 min_runs=min_runs,
265 max_runs=max_runs,
266 )
267 except TryAgain:
268 continue
270 # bulk insert the runs based on batch size setting
271 for batch in batched_iterable(runs_to_insert, insert_batch_size): 271 ↛ 272line 271 didn't jump to line 272 because the loop on line 271 never started
272 async with db.session_context(begin_transaction=True) as session:
273 inserted_runs = await models.deployments._insert_scheduled_flow_runs(
274 session=session, runs=list(batch)
275 )
276 total_inserted_runs += len(inserted_runs)
278 # if this is the last page of deployments, exit the loop
279 if len(deployment_ids) < deployment_batch_size: 279 ↛ 283line 279 didn't jump to line 283 because the condition on line 279 was always true
280 break
281 else:
282 # record the last deployment ID
283 last_id = deployment_ids[-1]
285 logger.info(f"Scheduled {total_inserted_runs} runs.")
288@perpetual_service(
289 enabled_getter=lambda: get_current_settings().server.services.scheduler.enabled,
290)
291async def schedule_recent_deployments(
292 perpetual: Perpetual = Perpetual(
293 automatic=True,
294 every=timedelta(
295 seconds=get_current_settings().server.services.scheduler.recent_deployments_loop_seconds
296 ),
297 ),
298) -> None:
299 """
300 Recent deployments scheduler - schedules deployments that were updated very recently.
302 This scheduler runs on a tight loop and ensures that runs from newly-created or
303 updated deployments are rapidly scheduled without waiting for the main scheduler.
305 Note that scheduling is idempotent, so it's okay for this scheduler to attempt
306 to schedule the same deployments as the main scheduler.
307 """
308 settings = get_current_settings().server.services.scheduler
309 deployment_batch_size = settings.deployment_batch_size
310 max_runs = settings.max_runs
311 min_runs = settings.min_runs
312 max_scheduled_time = settings.max_scheduled_time
313 min_scheduled_time = settings.min_scheduled_time
314 insert_batch_size = settings.insert_batch_size
315 loop_seconds = settings.recent_deployments_loop_seconds
317 db = provide_database_interface()
318 total_inserted_runs = 0
319 last_id = None
321 while True:
322 async with db.session_context(begin_transaction=False) as session:
323 query = _get_select_recent_deployments_to_schedule_query(
324 db, deployment_batch_size, loop_seconds
325 )
327 if last_id: 327 ↛ 328line 327 didn't jump to line 328 because the condition on line 327 was never true
328 query = query.where(db.Deployment.id > last_id)
330 result = await session.execute(query)
331 deployment_ids = result.scalars().unique().all()
333 try:
334 runs_to_insert = await _collect_flow_runs(
335 db,
336 session=session,
337 deployment_ids=deployment_ids,
338 max_scheduled_time=max_scheduled_time,
339 min_scheduled_time=min_scheduled_time,
340 min_runs=min_runs,
341 max_runs=max_runs,
342 )
343 except TryAgain:
344 continue
346 for batch in batched_iterable(runs_to_insert, insert_batch_size): 346 ↛ 347line 346 didn't jump to line 347 because the loop on line 346 never started
347 async with db.session_context(begin_transaction=True) as session:
348 inserted_runs = await models.deployments._insert_scheduled_flow_runs(
349 session=session, runs=list(batch)
350 )
351 total_inserted_runs += len(inserted_runs)
353 if len(deployment_ids) < deployment_batch_size: 353 ↛ 356line 353 didn't jump to line 356 because the condition on line 353 was always true
354 break
355 else:
356 last_id = deployment_ids[-1]
358 logger.info(f"Scheduled {total_inserted_runs} runs.")