lifecycle
Telemetry for pod and task lifecycle events.
Module
Functions
report_datasets_connected
def report_datasets_connected(datasource_types: Mapping[str, str]) ‑> None:Send a DatasetConnectedEvent per datasource. Never raises.
Arguments
datasource_types: Datasource class name, keyed by datasource name.
report_pod_error
def report_pod_error( phase: str, *, exc: BaseException | None = None, error_type: str | None = None, process_name: str | None = None, attempts: int | None = None,) ‑> None:Send a PodErrorEvent at ERROR level. Never raises.
Arguments
phase:"init","startup"or"process_spawn".exc: The exception, when in hand; its type and frames are sent.error_type: The error's type name, when there is no exception object.process_name: The background process that failed to spawn.attempts: How many spawn attempts were made.
report_recovery
def report_recovery( *, resubmitted: int, backlog_cancelled: int = 0, refresh_runs_ended: int = 0, slots_released: int = 0, failed_steps: list[str],) ‑> None:Send a RecoverySummaryEvent; WARNING when a step failed. Never raises.
Arguments
resubmitted: Replacement runs submitted for lost ones.backlog_cancelled: Queued duplicate runs cancelled.refresh_runs_ended: Wedged refresh runs ended.slots_released: Stranded concurrency slots released.failed_steps: Reconciliation steps that raised.
report_schema_error
def report_schema_error(dataset_id: str, stage: str, exc: BaseException) ‑> None:Send a SchemaGenerationErrorEvent with exc's frames. Never raises.
Arguments
dataset_id: The datasource whose schema failed.stage:"generation","upload"or"partial_upload".exc: The failure.
Classes
TaskTelemetry
class TaskTelemetry(identity: _TaskIdentity, pod_name: str):Lifecycle events for one worker task.
Identity is read from identity when each event is sent, so it reflects
the task as it is then. Every emission is swallowed on failure — telemetry
must not fail a task.
Create telemetry for the task identity describes.
Arguments
identity: Suppliestask_idandmodeller_name.pod_name: The pod running the task.
Methods
accepted
def accepted(self, datasource_type: str) ‑> None:Send TaskAcceptedEvent.
Arguments
datasource_type: Class name of the datasource the task runs on.
completed
def completed(self) ‑> None:Send TaskCompleteEvent.
failed
def failed( self, outcome: str, error_type: str, *, level: int = 40, exc: BaseException | None = None,) ‑> None:Send TaskErrorEvent for a task that did not complete.
Arguments
outcome: How the task ended; seeTaskErrorEvent.error_type: The exception type name, or a label where there is none.level: Log level, which Datadog reads as the event's severity.exc: The exception, when in hand, so its frames are attached.