Skip to main content

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: Supplies task_id and modeller_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; see TaskErrorEvent.
  • 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.