Skip to main content

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's run_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's run_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.