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 throughreport_step_progress.Nonewhen 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, orNoneif it returned. Defaults toNoneso 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.