InfluxElixir (InfluxElixir v0.1.35)

Copy Markdown View Source

Elixir client library for InfluxDB v3 with v2 compatibility.

All public API operations go through this facade module. Delegates to the configured client implementation (InfluxElixir.Client.HTTP or InfluxElixir.Client.Local).

Named Connections

All facade functions accept either a keyword config or an atom name. When an atom is passed, it is resolved via InfluxElixir.Connection.fetch!/1 from the persistent-term registry (populated automatically by InfluxElixir.ConnectionSupervisor on startup).

# Using a named connection (registered at startup)
InfluxElixir.health(:trading)
InfluxElixir.write(:trading, "cpu value=1.0", database: "prices")

# Using a raw config (e.g. LocalClient in tests)
InfluxElixir.health(conn)

Configuration

# config/config.exs
config :influx_elixir, :client, InfluxElixir.Client.HTTP

config :influx_elixir, :connections,
  trading: [
    host: "influx-trading",
    token: "...",
    database: "prices"
  ]

# config/test.exs
config :influx_elixir, :client, InfluxElixir.Client.Local

Telemetry

write/3 and the query functions emit [:influx_elixir, :write | :query, :start | :stop | :exception] events; see InfluxElixir.Telemetry.

Summary

Functions

Adds a new named connection dynamically at runtime.

Returns the configured client implementation module.

Creates a bucket in InfluxDB v2 (backwards compatibility). opts may carry retention: in seconds and org_id:; see InfluxElixir.Admin.Buckets.create/3.

Creates a database in InfluxDB v3. opts may carry retention: as a duration string such as "30d"; see InfluxElixir.Admin.Databases.create/3.

Creates an API token in InfluxDB v3.

Deletes a bucket in InfluxDB v2 (backwards compatibility).

Deletes a database in InfluxDB v3.

Deletes an API token in InfluxDB v3.

Sends a SQL statement to /api/v3/query_sql as it is.

Forces an immediate flush of the batch writer for a connection.

Checks the health of an InfluxDB instance.

Lists all buckets in InfluxDB v2 (backwards compatibility).

Lists all databases in InfluxDB v3.

Constructs a new Point struct.

Executes a Flux query against InfluxDB v2 (backwards compatibility).

Executes an InfluxQL query against InfluxDB v3.

Executes a SQL query against InfluxDB v3.

Executes a streaming SQL query, returning a lazy Stream.

Removes a named connection dynamically at runtime.

Resolves a connection reference to a config keyword list.

Returns batch writer statistics for a connection.

Writes line protocol to InfluxDB using the configured client.

Functions

add_connection(name, opts)

@spec add_connection(
  atom(),
  keyword()
) :: Supervisor.on_start_child()

Adds a new named connection dynamically at runtime.

If the connection does not start — a batch_writer: option that InfluxElixir.Write.BatchWriter refuses, say — the error is returned and nothing is left behind: the name does not resolve and the client's connection state is released.

client()

@spec client() :: module()

Returns the configured client implementation module.

create_bucket(connection, bucket_name, opts \\ [])

@spec create_bucket(
  InfluxElixir.Client.connection(),
  binary(),
  keyword()
) :: :ok | {:error, term()}

Creates a bucket in InfluxDB v2 (backwards compatibility). opts may carry retention: in seconds and org_id:; see InfluxElixir.Admin.Buckets.create/3.

create_database(connection, db_name, opts \\ [])

@spec create_database(
  InfluxElixir.Client.connection(),
  binary(),
  keyword()
) :: :ok | {:error, term()}

Creates a database in InfluxDB v3. opts may carry retention: as a duration string such as "30d"; see InfluxElixir.Admin.Databases.create/3.

create_token(connection, description, opts \\ [])

@spec create_token(
  InfluxElixir.Client.connection(),
  binary(),
  keyword()
) :: {:ok, map()} | {:error, term()}

Creates an API token in InfluxDB v3.

delete_bucket(connection, bucket_name)

@spec delete_bucket(InfluxElixir.Client.connection(), binary()) ::
  :ok | {:error, term()}

Deletes a bucket in InfluxDB v2 (backwards compatibility).

delete_database(connection, db_name)

@spec delete_database(InfluxElixir.Client.connection(), binary()) ::
  :ok | {:error, term()}

Deletes a database in InfluxDB v3.

delete_token(connection, token_id)

@spec delete_token(InfluxElixir.Client.connection(), binary()) ::
  :ok | {:error, term()}

Deletes an API token in InfluxDB v3.

execute_sql(connection, sql, opts \\ [])

@spec execute_sql(
  InfluxElixir.Client.connection(),
  binary(),
  keyword()
) :: {:ok, map() | [map()]} | {:error, term()}

Sends a SQL statement to /api/v3/query_sql as it is.

InfluxDB 3 Core refuses every statement that changes data (verified): DELETE, INSERT and UPDATE are {:error, %{status: 400, body: "Error during planning: DML not supported: ..."}}, CREATE and DROP are ... DDL not supported: ..., anything else is a 405. A SELECT returns {:ok, rows} typed like query_sql/3 rows. Remove data with delete_database/2 instead. Client.Local answers the same way; under its :v3_enterprise profile it also runs DELETE FROM m [WHERE ...] and returns {:ok, %{"rows_affected" => n}}.

flush(connection_name, timeout \\ 60000)

@spec flush(atom(), timeout()) :: :ok | {:error, :no_batch_writer}

Forces an immediate flush of the batch writer for a connection.

Returns :ok on success or {:error, :no_batch_writer} if no batch writer is configured for the given connection.

The optional timeout is forwarded to BatchWriter.flush/2 and bounds the underlying GenServer.call. See BatchWriter.flush/2 for the default and semantics.

health(connection)

@spec health(InfluxElixir.Client.connection()) :: {:ok, map()} | {:error, term()}

Checks the health of an InfluxDB instance.

list_buckets(connection)

@spec list_buckets(InfluxElixir.Client.connection()) ::
  {:ok, [map()]} | {:error, term()}

Lists all buckets in InfluxDB v2 (backwards compatibility).

list_databases(connection)

@spec list_databases(InfluxElixir.Client.connection()) ::
  {:ok, [map()]} | {:error, term()}

Lists all databases in InfluxDB v3.

point(measurement, fields, opts \\ [])

Constructs a new Point struct.

Parameters

  • measurement - measurement name
  • fields - field key-value pairs
  • opts - optional :tags and :timestamp

Examples

InfluxElixir.point("cpu", %{"value" => 0.64},
  tags: %{"host" => "server01"}
)

query_flux(connection, flux, opts \\ [])

Executes a Flux query against InfluxDB v2 (backwards compatibility).

query_influxql(connection, influxql, opts \\ [])

Executes an InfluxQL query against InfluxDB v3.

query_sql(connection, sql, opts \\ [])

Executes a SQL query against InfluxDB v3.

Supports transport: :http | :flight option for transport selection (InfluxElixir.Client.HTTP only; Client.Local is in-memory and ignores it).

Common Options

  • :database — overrides the connection-level default database.
  • :timeout — per-call receive timeout in milliseconds. Both the HTTP and Flight transports honour this; default is 30_000 ms.
  • :pool_timeout — per-call Finch pool checkout timeout in milliseconds (HTTP transport; default 5_000). Applies before :timeout.
  • :params — map of $name => value placeholder substitutions (HTTP transport only; Flight returns {:error, :params_unsupported_over_flight}).
  • :transport — :http (default) or :flight (Arrow Flight gRPC).
  • :flight_port — gRPC port for :flight; falls back to the connection's :flight_port, then 443. Pass tls: false for plaintext ports.

query_sql_stream(connection, sql, opts \\ [])

@spec query_sql_stream(
  InfluxElixir.Client.connection(),
  binary(),
  keyword()
) :: Enumerable.t()

Executes a streaming SQL query, returning a lazy Stream.

Use for large result sets to avoid loading all rows into memory. On the HTTP transport the JSONL response is decoded incrementally with back-pressure, so only one chunk plus a partial line is held in memory regardless of result size.

Because the return value is an Enumerable.t(), errors cannot be returned as an {:error, reason} tuple. Instead, failure classes — a missing database, a non-success HTTP status, or a transport error — are raised as an InfluxElixir.StreamError when the stream is enumerated. This mirrors the {:error, reason} contract of query_sql/3: an outage surfaces as an exception, never as a silent empty result.

Options

  • :database — overrides the connection-level default database.
  • :params — map of $name => value placeholder substitutions.
  • :timeout — per-call receive timeout in milliseconds.

remove_connection(name)

@spec remove_connection(atom()) :: :ok | {:error, term()}

Removes a named connection dynamically at runtime.

The connection's batch writer, if any, writes what it still holds before it stops (see "Shutdown" in InfluxElixir.Write.BatchWriter).

resolve_connection(name)

Resolves a connection reference to a config keyword list.

Accepts either an atom name (looked up via Connection.fetch!/1) or a keyword/map config (returned as-is).

Examples

resolve_connection(:trading)
resolve_connection(host: "localhost", token: "t")

stats(connection_name)

@spec stats(atom()) :: {:ok, map()} | {:error, :no_batch_writer}

Returns batch writer statistics for a connection.

Returns {:ok, stats_map} or {:error, :no_batch_writer} if no batch writer is configured for the given connection.

write(connection, line_protocol, opts \\ [])

Writes line protocol to InfluxDB using the configured client.

Goes through InfluxElixir.Write.Writer, so payloads over 1 KB are gzipped and a [:influx_elixir, :write, ...] telemetry span is emitted.

Options

  • :database — overrides the connection-level default database.
  • :precision — the unit of the timestamps (default nanoseconds); see InfluxElixir.Client.Local for the spellings each server accepts.
  • :accept_partial — InfluxDB 3 only. false makes the write all-or-nothing: the first bad line (a parse error or a schema conflict, in line order) rejects the payload, nothing is stored, and the error is {:error, %{status: 400, body: json}} with "line protocol parsing error" and that one line under "data". Default true: good lines are stored and the bad ones reported as a partial write.
  • :no_sync — InfluxDB 3 only. true acknowledges the write before it is persisted to the write-ahead log: faster, and a query right after it may not see the points yet (verified).