Skip to content

Python API

Generated from the docstrings. Everything here is importable from polars_telemetry or the module named in its heading.

Instrumenting

install

install(
    config: Config | None = None,
    exporter: Exporter | Sequence[Exporter] | None = None,
) -> Installation | None

Instrument every polars query this process runs.

Parameters:

  • config (Config | None, default: None ) –

    What to record. Defaults to Config().

  • exporter (Exporter | Sequence[Exporter] | None, default: None ) –

    Where queries go: one exporter or several. Defaults to OTelExporter.

Returns:

  • Installation | None –

    What was installed, including what the probe found about this polars; None when this polars cannot be instrumented at all.

Enabling polars' monitoring sets its engine affinity to "streaming", so this changes how queries execute and never happens on import. Calling it again while installed logs a warning and changes nothing.

uninstall

uninstall() -> None

Stop instrumenting, and hand queries back to Polars Cloud if it was there.

The engine affinity goes back to what it was before install(), unless the application chose another engine in the meantime.

Config

Config(
    node_metrics: bool = True,
    include_plan: bool = False,
    call_site: bool = True,
    redaction: Redaction | None = None,
)

What to record about each query. Every field has a working default.

node_metrics

node_metrics: bool = True

Read per-node counters once when the query ends.

polars exposes cumulative counters and no per-node timestamps, so these are exact totals with no timing. Disable to emit the query span alone.

include_plan

include_plan: bool = False

Attach the full plan and its counters to the span as JSON.

Off by default: it is kilobytes per span and identical for every run of a shape. Turn it on when you want the topology, which nothing else carries.

call_site

call_site: bool = True

Record the file, line and function that ran the query.

Costs well under a microsecond. Turn it off to keep source paths out of telemetry you do not control.

redaction

redaction: Redaction | None = None

What to mask before any exporter receives a query; None masks nothing.

Redaction() masks literal values. One exporter can be given its own with redacted(). Metrics never carry literals, whatever this says.

Masking

Redaction

Redaction(
    strings: bool = True,
    numbers: bool = True,
    temporal: bool = True,
    paths: bool = False,
    call_site: bool = False,
    labels: bool = False,
    custom: Callable[[str], str] | None = None,
)

What to mask before a query reaches an exporter.

The default masks literal values: text, numbers, and dates and times. Column names, operators and the plan's shape are always kept, so a masked query still shows which filter was slow.

Examples:

>>> Redaction()  # literal values
>>> Redaction(numbers=False)  # keep numbers, mask the rest
>>> Redaction(paths=True, call_site=True, labels=True)  # as little as possible

strings

strings: bool = True

Quoted text such as "Brand#12" becomes "<str>". Column and alias names are kept.

numbers

numbers: bool = True

Numbers such as 60.0 or 1.0000e-9 become <num>.

temporal

temporal: bool = True

Dates, datetimes, times and durations become <date>, <datetime>, <time> and <duration>.

paths

paths: bool = False

File paths that are scanned or written become <path>.

call_site

call_site: bool = False

Drop the file, line and function that ran the query.

labels

labels: bool = False

Drop labels set with label().

custom

custom: Callable[[str], str] | None = None

Your own rule, applied to every expression and error message after the masks above.

masks

masks: tuple[str, ...]

What this masks, by field name, for a reader to show.

redacted

redacted(
    exporter: Exporter, redaction: Redaction | None
) -> Redacted

Give one exporter its own redaction, in place of Config.redaction.

Examples:

>>> polars_telemetry.install(
...     Config(redaction=Redaction()),
...     exporter=[
...         redacted(OTelExporter(config), Redaction(paths=True, call_site=True)),
...         redacted(FileExporter("profiles/full.jsonl"), None),
...     ],
... )

Parameters:

  • exporter (Exporter) –

    Any exporter.

  • redaction (Redaction | None) –

    What to mask for this exporter; None sends it everything, whatever Config.redaction says.

Returns:

  • Redacted –

    The exporter, to pass to install().

Labelling and scoping

label

label(name: str) -> Iterator[None]

Label every query run inside the block.

Examples:

>>> with label("etl"), label("customers"):
...     frame.collect()  # labelled "etl/customers"

Parameters:

  • name (str) –

    Any non-empty text. Nested labels are joined with /.

Raises:

  • ValueError –

    If name is empty or not a string.

The label goes on the query's span as polars.query.label and into its profile, but never onto metrics: free-form values would make unbounded metric series. Each thread has its own labels. Queries run with collect_async() or collect_batches() carry none: polars reports those from its own threads.

profile

profile(config: Config | None = None) -> Iterator[Session]

Collect every query run inside the block.

Examples:

>>> with profile() as session:
...     frame.collect()
>>> session.slowest.call_site

Parameters:

  • config (Config | None, default: None ) –

    Used when nothing is installed yet; an existing installation keeps its own. The session masks what the installation masks and what config.redaction adds.

Queries still reach any exporters already installed, and blocks may nest. If nothing was installed, instrumentation is installed for the block and removed after it. A block collects every query that finishes while it is open, including ones other threads ran.

Session

Session()

The queries that ran inside a profile() block, in the order they finished.

Iterate it, index it or take its len() like a list of Query.

slowest

slowest: Query | None

The query with the longest wall time, or None if none ran.

wall_ms

wall_ms: float

Summed wall time. Not elapsed time: queries may overlap.

profiles

profiles() -> list[dict[str, Any]]

Each query as a profile document, the format the viewer reads.

Literals are already masked if the block's config asked for it.

write

write(path: str | Path) -> Path

Write the queries to a .jsonl file the viewer can open.

Parameters:

  • path (str | Path) –

    Where to write. An existing file is replaced.

Returns:

  • Path –

    The path written.

Exporters

OTelExporter

OTelExporter(config: Config)

One span per query and per-node metrics, through OpenTelemetry.

Uses the global tracer and meter providers, so the application's OpenTelemetry SDK decides where they go. Without an SDK installed, both are no-ops and nothing is sent.

Parameters:

  • config (Config) –

    Shared with install(); node_metrics and include_plan apply here. Masking happens before the query arrives; see Config.redaction.

config

config: Config

The configuration this exporter was made with.

DogStatsdExporter

DogStatsdExporter(
    client: StatsdClient | None = None,
    *,
    metric_names: Mapping[str, str | None]
    | Callable[[str], str | None]
    | None = None,
    tag_names: Mapping[str, str | None] | None = None,
    tag_labels: bool = False,
    distributions: bool = True,
)

Per-query and per-node metrics as DogStatsD, with tags.

Totals are sent as counts; times and ratios as distributions, or as histograms with distributions=False.

Parameters:

  • client (StatsdClient | None, default: None ) –

    A datadog.DogStatsd. Defaults to datadog.statsd, the client datadog.initialize() configures. Turn on its buffering and background sender, so sending happens off the query's thread.

  • metric_names (Mapping[str, str | None] | Callable[[str], str | None] | None, default: None ) –

    Renames metrics, from their names in Spans and metrics. A mapping or a function; a name mapped to None is not sent.

  • tag_names (Mapping[str, str | None] | None, default: None ) –

    Renames tag keys: engine, fingerprint, node_kind and direction. A key mapped to None is not sent, such as {"fingerprint": None} to keep one series per query shape off a bill.

  • tag_labels (bool, default: False ) –

    Also tag every metric with the query's label. Labels are free-form, so only do this when yours come from a small, fixed set.

  • distributions (bool, default: True ) –

    Send times and ratios as distributions (|d), which Datadog aggregates across hosts. False sends histograms (|h), which Telegraf summarises per flush; Telegraf keeps only one sample of a distribution per flush.

close

close() -> None

Send what the client still holds. Called by uninstall() and at exit.

FileExporter

FileExporter(
    path: str | Path, *, max_bytes: int = DEFAULT_MAX_BYTES
)

Append one profile per query to a .jsonl file the viewer can open.

Parameters:

  • path (str | Path) –

    The session file. Its directory is created if missing.

  • max_bytes (int, default: DEFAULT_MAX_BYTES ) –

    64 MiB by default. When the file would grow past this, it moves to <name>.1, replacing the previous one, and a new file starts. At most about twice this is on disk. A profile is never split, so a file can run over by one record.

path

path: Path

The file being written.

ConsoleExporter

ConsoleExporter(stream: TextIO | None = None)

Print a short summary of each query: totals, call site, slowest nodes.

Parameters:

  • stream (TextIO | None, default: None ) –

    Where to write. Defaults to standard error.

Exporter

Anything with an export(query) method can receive queries.

An exporter that holds data, such as a buffer, may also have a close() method. uninstall() calls it, and so does the process on exit.

export

export(query: Query) -> None

Handle one finished query.

Called on the thread that ran the query, once per query, after it has finished, so the time spent here is added to the caller's. An exception is logged and counted; after five, this exporter stops receiving queries and the others carry on.

What an exporter receives

Query

Query(
    query_id: UUID,
    wall_ms: float,
    plan: dict[int, PlanNode],
    logical: dict[int, PlanNode] = dict(),
    metrics: dict[int, NodeMetrics] = dict(),
    call_site: CallSite | None = None,
    label: str | None = None,
    engine: str | None = None,
    polars_version: str = "",
    fingerprint: str = "",
    diagnostics: Diagnostics | None = None,
    failed: str | None = None,
    redaction: Redaction | None = None,
    started_unix_ns: int = 0,
)

One finished query, as every exporter receives it.

query_id

query_id: UUID

polars' id for the query: a UUIDv7, so ids sort by start time.

wall_ms

wall_ms: float

Time from start to finish, in milliseconds.

plan

plan: dict[int, PlanNode]

Physical plan, by node id; empty when the query did not run on the streaming engine. Node ids here are what metrics key on.

logical

logical: dict[int, PlanNode] = field(default_factory=dict)

IR plan. Carries the user's own column names; the physical plan rewrites group-by keys and aggregations to _POLARS_TMP_N.

call_site

call_site: CallSite | None = None

Where in the caller's code the query ran.

label

label: str | None = None

What the application called it, via polars_telemetry.label().

engine

engine: str | None = None

The polars engine that ran it; None when it failed before planning.

polars_version

polars_version: str = ''

The polars that ran it. Carried here so nothing downstream imports polars.

fingerprint

fingerprint: str = ''

The plan shape, hashed once on arrival; empty on a Query built by hand.

diagnostics

diagnostics: Diagnostics | None = None

Derived once on arrival, so every exporter reports the same figures.

failed

failed: str | None = None

polars' error message when the query failed, else None.

redaction

redaction: Redaction | None = None

What was masked before this reached the exporter, else None.

started_unix_ns

started_unix_ns: int = 0

When the query started, in nanoseconds since the Unix epoch.

cpu_ms

cpu_ms: float

Summed node self time. Exceeds wall time on a parallel query.

result_rows

result_rows: int | None

Rows reaching the sink, when the sink reported any.

hottest

hottest: tuple[PlanNode, NodeMetrics] | None

The node with the most self time -- usually the whole answer.