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 oftake_pending;Nonemerges 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