Skip to main content

datasource

Aggregated telemetry for per-file datasource outcomes.

A datasource can skip hundreds of thousands of files, and does so again on every listing pass for task-scoped skips. Rather than one Datadog log per file, DatasourceTelemetry counts per-file outcomes in memory and sends one summary log per kind of outcome each time a window closes:

  • skipped files, as FileSkipSummaryEvent;
  • files whose matching series were cut to one row, as FileMultiSeriesSummaryEvent.

Classes​

DatasourceTelemetry​

class DatasourceTelemetry(datasource: _NamedDatasource):

Per-file outcome telemetry for one datasource.

Recording is cheap: no filesystem access and no emission unless the window has outlived skipped_file_telemetry_interval_seconds. The owner calls flush where a window naturally ends (a listing pass, a task). Nothing is sent for an empty window, and every failure is swallowed — telemetry must not break file processing.

A copy (pickled into a worker process, or deep-copied) starts with an empty window, so it reports only its own outcomes. A worker process that does not live for a whole pass or task should hand its counts back with take_pending for the owning process to merge, rather than flush: one summary per short-lived worker is one summary per file.

Create telemetry for datasource.

Arguments

  • datasource: The datasource reported on. Its type and current name are read when a summary is sent, so a rename after construction is picked up.

Variables​

  • pending : int - Outcomes of any kind recorded since the last flush.

Methods​


flush​

def flush(self, trigger: SummaryTrigger) ‑> None:

Send the window's summaries and start a new window.

Arguments

  • trigger: What closed the window.

merge​

def merge(self, pending: PendingFileTelemetry | None) ‑> None:

Add counts taken from another process's copy to this window.

Arguments

  • pending: The output of take_pending; None merges nothing.

record_multi_series_reduction​

def record_multi_series_reduction(    self, *, rows_before: int, laterality: str, series_protocol: str,) ‑> None:

Count one file whose matching series were cut to its first row.

Arguments

  • rows_before: Matching rows the file produced.
  • laterality: The configured laterality.
  • series_protocol: The configured series protocol.

record_skip​

def record_skip(    self, filename: str | os.PathLike[str], reason: Enum, source: SkipSource,) ‑> None:

Count one skipped file, when enable_skipped_file_telemetry is set.

Arguments

  • filename: The skipped file.
  • reason: Why it was skipped.
  • source: "datasource" for cached skips, "task" for task-scoped ones.

take_pending​

def take_pending(self) ‑> PendingFileTelemetry:

Remove and return everything counted since the last flush.

Returns The pending counts, for merge in another process.

PendingFileTelemetry​

class PendingFileTelemetry(    skips: ForwardRef('_SkipCounts'), multi_series: ForwardRef('_MultiSeriesCounts'),):

Counts taken from one process's DatasourceTelemetry, to merge elsewhere.

Picklable, so a worker process can return what it counted to the process that owns the datasource instead of sending its own summaries.

Variables​

  • multi_series : bitfount.telemetry.datasource._MultiSeriesCounts - Alias for field number 1
  • skips : bitfount.telemetry.datasource._SkipCounts - Alias for field number 0