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

1""" 

2The Scheduler service. 

3 

4This service schedules flow runs from deployments with active schedules. 

5""" 

6 

7from __future__ import annotations 

8 

9import datetime 

10import logging 

11from datetime import timedelta 

12from typing import Any, Sequence 

13from uuid import UUID 

14 

15import sqlalchemy as sa 

16from docket import Perpetual 

17from sqlalchemy.ext.asyncio import AsyncSession 

18 

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 

27 

28logger: logging.Logger = get_logger(__name__) 

29 

30 

31class TryAgain(Exception): 

32 """Internal control-flow exception used to retry the Scheduler's main loop""" 

33 

34 

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. 

43 

44 The query gets the IDs of any deployments where ANY active schedule has: 

45 

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 

50 

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. 

54 

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") 

61 

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) 

69 

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 ) 

88 

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 ) 

102 

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 ) 

117 

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 

128 

129 

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 

168 

169 

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 

210 

211 

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. 

225 

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 

238 

239 db = provide_database_interface() 

240 total_inserted_runs = 0 

241 last_id = None 

242 

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 ) 

248 

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) 

252 

253 result = await session.execute(query) 

254 deployment_ids = result.scalars().unique().all() 

255 

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 

269 

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) 

277 

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] 

284 

285 logger.info(f"Scheduled {total_inserted_runs} runs.") 

286 

287 

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. 

301 

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. 

304 

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 

316 

317 db = provide_database_interface() 

318 total_inserted_runs = 0 

319 last_id = None 

320 

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 ) 

326 

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) 

329 

330 result = await session.execute(query) 

331 deployment_ids = result.scalars().unique().all() 

332 

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 

345 

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) 

352 

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] 

357 

358 logger.info(f"Scheduled {total_inserted_runs} runs.")