follow_up
Scheduling a metadata run's own replacement for later.
Two cases hand a metadata run's work forward to a later run of the same
deployment with the same parameters: a scan that must wait for indexing to
finish, and a run cut short because its datasource root became unreachable.
Both submit through schedule_follow_up.
Transient-failure follow-ups form a bounded chain, counted by the
transient_retry: tag. Each follow-up probes the root before doing anything
else and, if it is still unreachable, hands on to the next link at once, so a
retry into a dead share costs one stat. A follow-up that is not yet due is
invisible to dedup (is_pending_transient_retry), so it never blocks a fresh
submission; if one supersedes it, it fires later and finds nothing to do.
Module
Functions
current_flow_run_tags
def current_flow_run_tags() ‑> list[str]:Tags of the flow run this code is executing in, empty outside one.
schedule_follow_up
def schedule_follow_up( deployment: str, *, parameters: dict[str, Any], task_hash: str, run_type: str, delay: timedelta, tags: Sequence[str] = (),) ‑> None:Schedule a run of deployment delay from now, returning immediately.
Arguments
deployment: The"<flow>/<deployment>"name to run.parameters: The run's parameters, forwarded verbatim.task_hash: The(pod, datasource)the run belongs to.run_type: The run'srun_type:tag value.delay: How long from now the run is scheduled for.tags: Tags added to the task-hash and run-type tags.
Raises
Exception: Whatever the Prefect submission raised.
schedule_transient_retry
def schedule_transient_retry( deployment: str, *, parameters: dict[str, Any], task_hash: str, run_type: str, log: logging.Logger | logging.LoggerAdapter[Any],) ‑> bool:Schedule the next follow-up in this run's transient-failure chain.
The attempt is read off this run's own transient_retry: tag, so any
submitter's run starts a fresh chain. Never raises.
Arguments
deployment: The"<flow>/<deployment>"name to run.parameters: This run's parameters, forwarded verbatim.task_hash: The(pod, datasource)the run belongs to.run_type: The run'srun_type:tag value.log: The flow-run logger.
Returns
Whether a follow-up was scheduled. False when the chain is exhausted
or the submission failed; the nightly refresh is then what picks the
work up.