InfluxElixir.Write.BatchWriter (InfluxElixir v0.1.35)

Copy Markdown View Source

GenServer-based batch writer with configurable flush intervals, batch sizes, retry with exponential backoff, and backpressure.

Points or pre-encoded line protocol strings are buffered in memory and flushed either when the buffer reaches batch_size or when the flush_interval_ms timer fires — whichever comes first.

Options

Validated by start_link/1, which returns {:error, %NimbleOptions.ValidationError{}} for an unknown key or a value of the wrong type — a batch_size: 0 used to refuse every write as :buffer_full, and a misspelt key was silently ignored.

  • :connection (term/0) - Required. connection term passed to InfluxElixir.Write.Writer

  • :database (String.t/0) - database every flush writes to. A :database in :write_opts wins.

  • :batch_size (pos_integer/0) - maximum points per flush The default value is 5000.

  • :flush_interval_ms (pos_integer/0) - timer interval in milliseconds The default value is 1000.

  • :jitter_ms (non_neg_integer/0) - random jitter added to the flush timer and retry backoff The default value is 0.

  • :max_retries (non_neg_integer/0) - retry attempts for 5xx and transport errors. 4xx is never retried. The default value is 3.

  • :base_retry_delay_ms (non_neg_integer/0) - base of the exponential backoff: attempt N waits about base * 2^N The default value is 100.

  • :no_sync (boolean/0) - when true, write_sync/3 behaves like write/3. Not InfluxDB 3's no_sync write parameter: pass that, like accept_partial: false, in :write_opts. The default value is false.

  • :write_opts (keyword/0) - options for InfluxElixir.Write.Writer.write/3 on every flush (:database, :timeout, :precision, ...). A Point's DateTime timestamp is encoded in its :precision; an integer timestamp must already be in that unit. The default value is [].

  • :client (atom/0) - client module to write with instead of the configured one (InfluxElixir.Client.impl/0)

  • :name (term/0) - a GenServer name to register the writer under

  • :shutdown - how long a supervisor waits for the final flush when it stops the writer (OTP's default for workers); see "Shutdown" The default value is 5000.

Backpressure

While any batch is being retried, the automatic flushes (batch size reached, timer fired) wait for every retry chain to finish instead of starting another against a server that is already failing; writes keep buffering meanwhile. Once the buffer holds 10 * batch_size entries, write/3 and write_sync/3 return {:error, :buffer_full} until the last chain ends. An explicit flush/2, and write_sync/3 without :no_sync, always flush immediately — and may start a chain of their own. When the last chain ends the deferred buffer is flushed if it has reached batch_size; otherwise the timer takes it.

Shutdown

The writer traps exits, so a supervisor stopping it — application shutdown, InfluxElixir.remove_connection/1 — runs terminate/2, which writes every batch still being retried and then the buffer, once each and in the order they were written, and answers any write_sync/3 caller waiting on them. A write that fails there is logged and dropped. The supervisor kills the writer if that takes longer than :shutdown.

Retry Policy

Only 5xx and network errors are retried using asynchronous exponential backoff with optional jitter. 4xx errors are discarded and logged. Retries are non-blocking — the GenServer continues to accept messages between retry attempts. A write_sync/3 caller whose batch is being retried is answered with that chain's final result.

Stats

Call stats/1 to retrieve a map with :total_writes, :total_errors, and :total_bytes counters.

Summary

Functions

The child spec a supervisor starts the writer with; :shutdown in opts sets how long it waits for the final flush.

Forces an immediate flush of the buffer.

Starts a BatchWriter GenServer linked to the calling process.

Returns the current stats map.

Buffers a point (or line protocol binary) for writing.

Synchronously writes a point and waits for the next flush to complete.

Types

stat_key()

@type stat_key() :: :total_writes | :total_errors | :total_bytes

stats()

@type stats() :: %{required(stat_key()) => non_neg_integer()}

t()

@type t() :: %InfluxElixir.Write.BatchWriter{
  base_retry_delay_ms: non_neg_integer(),
  batch_size: pos_integer(),
  buffer: [binary()],
  buffer_size: non_neg_integer(),
  chains: %{required(integer()) => {binary(), GenServer.from() | nil}},
  connection: term(),
  database: binary() | nil,
  flush_interval_ms: pos_integer(),
  jitter_ms: non_neg_integer(),
  max_retries: non_neg_integer(),
  no_sync: boolean(),
  pending_sync: GenServer.from() | nil,
  stats: stats(),
  timer_ref: reference() | nil,
  write_opts: keyword()
}

Functions

child_spec(init_arg)

@spec child_spec(keyword()) :: Supervisor.child_spec()

The child spec a supervisor starts the writer with; :shutdown in opts sets how long it waits for the final flush.

flush(server, timeout \\ 60000)

@spec flush(GenServer.server(), timeout()) :: :ok

Forces an immediate flush of the buffer.

Bounded by timeout (default: 60000 ms). Returns :ok once the underlying HTTP write completes (or schedules a retry). Retries scheduled by do_flush are asynchronous and do NOT extend the caller's wait.

Parameters

  • server - PID or registered name of the BatchWriter
  • timeout - GenServer.call wait bound in ms (default: 60_000)

Examples

iex> InfluxElixir.Write.BatchWriter.flush(pid)
:ok

start_link(opts)

@spec start_link(keyword()) ::
  GenServer.on_start() | {:error, NimbleOptions.ValidationError.t()}

Starts a BatchWriter GenServer linked to the calling process.

Options

See module documentation for available options.

Examples

iex> {:ok, pid} = InfluxElixir.Write.BatchWriter.start_link(
...>   connection: conn,
...>   database: "mydb"
...> )
iex> is_pid(pid)
true

stats(server)

@spec stats(GenServer.server()) :: {:ok, stats()}

Returns the current stats map.

Keys

  • :total_writes - total successful write operations
  • :total_errors - total failed write operations
  • :total_bytes - total bytes flushed

Examples

iex> {:ok, stats} = InfluxElixir.Write.BatchWriter.stats(pid)
iex> Map.keys(stats)
[:total_bytes, :total_errors, :total_writes]

write(server, payload, timeout \\ 60000)

@spec write(
  GenServer.server(),
  InfluxElixir.Write.Point.t() | binary(),
  timeout()
) :: :ok | {:error, :buffer_full | term()}

Buffers a point (or line protocol binary) for writing.

Returns immediately after buffering unless the buffer reaches batch_size, in which case handle_call triggers a synchronous do_flush that calls the HTTP client. The timeout argument bounds the GenServer.call/3 wait — defaults to 60000 ms, generous enough to cover a single HTTP write at Client.HTTP's 30s default.

Parameters

  • server - PID or registered name of the BatchWriter
  • payload - a InfluxElixir.Write.Point.t() or pre-encoded binary
  • timeout - GenServer.call wait bound in ms (default: 60_000)

Examples

iex> InfluxElixir.Write.BatchWriter.write(pid, "cpu value=1.0")
:ok

write_sync(server, payload, timeout \\ 300_000)

@spec write_sync(
  GenServer.server(),
  InfluxElixir.Write.Point.t() | binary(),
  timeout()
) :: :ok | {:error, term()}

Synchronously writes a point and waits for the next flush to complete.

Blocks until the buffered data has been flushed and the write is confirmed. When no_sync: true is configured, behaves identically to write/3.

The caller waits through the full retry chain. The default timeout of 300000 ms covers up to max_retries + 1 HTTP writes at the default 30s HTTP timeout plus exponential backoff. Override for endpoints with longer expected tail latencies.

Parameters

  • server - PID or registered name of the BatchWriter
  • payload - a InfluxElixir.Write.Point.t() or pre-encoded binary
  • timeout - GenServer.call wait bound in ms (default: 300_000). Pass :infinity for unbounded blocking.

Examples

iex> InfluxElixir.Write.BatchWriter.write_sync(pid, "cpu value=1.0")
:ok