context
Run context carried onto every Datadog telemetry event.
Identity (username, pod_identifier) and run coordinates (dataset_name,
task_id, run_id) are known at different times and in different places from
where telemetry is configured: setup_datadog_telemetry runs at process start,
before login in several entry points, and its tags are fixed for the life of the
handler. This module holds those values instead, and DatadogLogsHandler.emit
reads them per record.
Two layers, because a contextvars value set after a worker thread copied its
context is invisible to that thread:
- The global layer is process-wide and lock-guarded, for values true of the whole process - who is logged in, which pod this is.
- The overlay is a
ContextVar, for values true of one run or one step.run_on_daemon_threadcopies the context, so a DAG step inherits the overlay its executor set.
Every read is total: telemetry must not raise into the code it observes.
Module
Functions
clear_global_telemetry_context
def clear_global_telemetry_context() ‑> None:Empty the process-wide telemetry context.
clear_telemetry_context
def clear_telemetry_context() ‑> None:Drop everything set_telemetry_context put in the current context.
The counterpart to setting without a scope: a process that reuses a context
across runs - and a test - needs a way back to empty. Values held by an
open telemetry_context_scope are restored by that scope as usual.
current_files_in_scope
def current_files_in_scope() ‑> int | None:Return the file count noted for the running DAG or wave, if any.
current_step_progress
def current_step_progress() ‑> StepProgress | None:Return the recorder bound to the current context, if any.
current_telemetry_context
def current_telemetry_context() ‑> dict[str, str | int | float | bool]:Return the context to attach to an event, overlay over global.
Returns A new dict; mutating it does not affect either layer. Empty if nothing has been set, or if reading either layer failed.
files_in_scope_slot
def files_in_scope_slot( ,) ‑> collections.abc.Iterator[FilesInScopeSlot]:Bind a fresh slot for the entry point about to run to note into.
The slot starts from the enclosing slot's value, so an entry point that notes nothing reports the denominator it was called with rather than nothing at all.
note_files_in_scope
def note_files_in_scope(count: int | None) ‑> None:Record how many files the running DAG or wave was given.
Written into the slot the timing decorator bound, not into a value scoped around the entry point: the timing event is emitted after the body returns, so a scope closed inside it would already have been restored. Writing into the innermost slot is what keeps a wave's count off the run's timing.
Arguments
count: Files in scope, orNonewhen the run has no file scope to report - an unwaved run does not narrow the datasource, so there is no cheap count to take and none is invented.
report_step_progress
def report_step_progress( *, files_processed: int | None = None, patients_processed: int | None = None,) ‑> None:Record how much work the running step did.
A no-op when no recorder is bound, so a step function stays callable outside a DAG run and in its own unit tests. Only the arguments given are written, so a step can report the two counts from different places.
Arguments
files_processed: Files this step actually read this run.patients_processed: Distinct patients this step handled this run.
set_global_telemetry_context
def set_global_telemetry_context(**fields: Any) ‑> None:Merge fields into the process-wide telemetry context.
Call this as soon as a value is known - in particular immediately before
setup_datadog_telemetry, so that events emitted from the moment the
handler exists already carry it.
Arguments
**fields: Context entries. ANonevalue is ignored rather than recorded, so a caller can pass an optional value unconditionally.
set_telemetry_context
def set_telemetry_context(**fields: Any) ‑> None:Add fields to the telemetry context without a scope to leave.
For a value that describes the whole of the work the current context is
doing, set where that work begins - a DAG run naming its dataset. Use
telemetry_context_scope instead wherever the value has an end as well as
a beginning; this one stays until the context it was set in goes away, or
until the next run overwrites it.
Arguments
**fields: Context entries.Nonevalues are ignored.
step_progress_scope
def step_progress_scope(progress: StepProgress) ‑> collections.abc.Iterator[None]:Make progress the recorder report_step_progress writes to.
Arguments
progress: The recorder to bind for the duration of the block.
telemetry_context_scope
def telemetry_context_scope(**fields: Any) ‑> collections.abc.Iterator[None]:Add fields to the telemetry context for the duration of the block.
Nested scopes merge, innermost winning, and the overlay is restored on exit whether the block returned or raised.
Arguments
**fields: Context entries.Nonevalues are ignored.
Classes
FilesInScopeSlot
class FilesInScopeSlot(value: int | None = None):Where one entry point's file denominator is held while it runs.
A slot per entry point rather than a single value per context: a wave runs inside the run that scheduled it, and both are timed, so a shared value would leave the run's timing reporting the last wave's count. The slot is an object the caller keeps a reference to, so the timing decorator can read it after the body has returned and the binding has been restored.
Attributes
value: Files this entry point was given, orNonewhen it has none to report.
Variables
- static
value : int | None
StepProgress
class StepProgress( files_processed: int | None = None, patients_processed: int | None = None,):How much work one DAG step invocation did.
Mutable by design. contextvars.copy_context() copies bindings, not the
objects they point at, so a step mutating this on its worker thread is seen
by the executor that created it - which is how a count reported inside a
step reaches the timing event emitted around it.
Both fields mean work done by this invocation and nothing else. A step
that cannot answer honestly leaves them None; a partition size or a
denominator does not belong here. files_in_scope on the executor's own
timing events is the field for a denominator.
Attributes
files_processed: Files this step actually read this run, excluding ones skipped because a cached result already covered them.patients_processed: Distinct patients this step handled this run.