Skip to main content

execution

Execution-time telemetry for the DAG paths.

DAG step and entry-point timings are emitted as DAGExecutionTimeEvent log events through telemetry_logger. Every timing carries how the work ended - status and, when it did not succeed, error_type - since the timing is emitted from a finally and a duration alone cannot tell a step that completed from one that crashed halfway.

The outcome also picks the record's level, because Datadog reads the level as the log's severity and that is what an error count is built from. An error is logged as one; a cancellation is a warning, since the pod cancels in-flight work on every redeploy and those must not read as failures. The exception is passed as exc_info either way, so the handler can attach frames to both.

Volume divides by who can answer: a step reports its own files_processed and patients_processed through report_step_progress, while the executor's entry points carry files_in_scope, the run's denominator.

These log events are the only form these timings take; the legacy execution_time_s custom metric is no longer sent for them.

Module​

Functions​

record_dag_step_time​

def record_dag_step_time(    start: float,    *,    step_name: str,    task_hash: str | None = None,    files_processed: int | None = None,    patients_processed: int | None = None,    exc: BaseException | None = None,) ‑> None:

Emit the execution time of one DAG step, and how it ended.

Arguments

  • start: time.perf_counter() captured before the step began.
  • step_name: The step's name.
  • task_hash: Task hash of the run the step belongs to, when known.
  • files_processed: Files this invocation read, when the step reported it through report_step_progress. None when it reported nothing — a step that cannot answer leaves the field empty rather than having a number guessed for it.
  • patients_processed: Patients this invocation handled, same contract.
  • exc: The exception that ended the step, or None if it returned. Defaults to None so that a caller which does not track the outcome keeps working; callers that can tell should pass it.

track_dag_execution_time​

def track_dag_execution_time(func: F) ‑> F:

Decorator emitting a DAG entry point's execution time.

Emits a DAGExecutionTimeEvent naming the function and, for a method, its class. Supports sync and async callables.

Arguments

  • func: The function to wrap.

Returns The wrapped function.