InfluxElixir.Client.Local.Flux (InfluxElixir v0.1.35)

Copy Markdown View Source

The Flux subset InfluxElixir.Client.Local answers, with InfluxDB 2's semantics (verified against influxdb:2.7; see docs/design/2026-09-24_local-flux-pipeline.md).

A query is a pipeline, from(bucket: "b") |> range(...) |> ..., and every stage is applied — or the query is refused. Before this module the double matched a few regexes anywhere in the text and ignored the rest, so |> mean() returned the raw rows.

Supported stages:

  • from(bucket: "b") — first; the bucket must exist (404 otherwise)
  • range(start: s[, stop: s]) — required, as on the engine; s is Unix seconds, an RFC3339 time, a negative duration (-1h, -30m, -7d, -10s, -2w) or now(); stop defaults to now. Every row carries _start and _stop.
  • filter(fn: (r) => ...) — r.key or r["key"] compared with == != < <= > >= against a string, number or boolean, combined with and, or, not and parentheses. A key the row lacks never matches.
  • first() last() min() max() — the selected row per table
  • mean() sum() count() — one row per table, without _time (mean is a float; count counts rows)
  • limit(n: N[, offset: M]) — per table
  • yield(name: "x") — names the result

Tables are the series — measurement, tag set, field — numbered from 0 in that order, rows in time order, as the engine numbers them.

Summary

Types

A comparison operand: a row key compared against a literal.

A parsed query.

One pipeline stage after from.

Functions

Parses a Flux query. now_ns is the instant now() and a relative range resolve against. Returns {:error, message} for syntax the double does not model.

Runs a parsed query over the bucket's points (%{measurement, tags, fields, timestamp} maps, duplicates already merged).

Types

predicate()

@type predicate() ::
  {:cmp, binary(), binary(), term()}
  | {:and, predicate(), predicate()}
  | {:or, predicate(), predicate()}
  | {:not, predicate()}

A comparison operand: a row key compared against a literal.

query()

@type query() :: %{bucket: binary(), stages: [stage()]}

A parsed query.

stage()

@type stage() ::
  {:range, integer(), integer()}
  | {:filter, predicate()}
  | {:selector, :first | :last | :min | :max}
  | {:aggregate, :mean | :sum | :count}
  | {:limit, non_neg_integer(), non_neg_integer()}
  | {:yield, binary()}

One pipeline stage after from.

Functions

parse(flux, now_ns)

@spec parse(binary(), integer()) :: {:ok, query()} | {:error, binary()}

Parses a Flux query. now_ns is the instant now() and a relative range resolve against. Returns {:error, message} for syntax the double does not model.

run(map, points)

@spec run(query(), [map()]) :: {:ok, [map()]} | {:error, binary()}

Runs a parsed query over the bucket's points (%{measurement, tags, fields, timestamp} maps, duplicates already merged).