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 toInfluxElixir.Write.Writer:database(String.t/0) - database every flush writes to. A:databasein:write_optswins.:batch_size(pos_integer/0) - maximum points per flush The default value is5000.:flush_interval_ms(pos_integer/0) - timer interval in milliseconds The default value is1000.:jitter_ms(non_neg_integer/0) - random jitter added to the flush timer and retry backoff The default value is0.:max_retries(non_neg_integer/0) - retry attempts for 5xx and transport errors. 4xx is never retried. The default value is3.:base_retry_delay_ms(non_neg_integer/0) - base of the exponential backoff: attempt N waits aboutbase * 2^NThe default value is100.:no_sync(boolean/0) - whentrue,write_sync/3behaves likewrite/3. Not InfluxDB 3'sno_syncwrite parameter: pass that, likeaccept_partial: false, in:write_opts. The default value isfalse.:write_opts(keyword/0) - options forInfluxElixir.Write.Writer.write/3on every flush (:database,:timeout,:precision, ...). APoint'sDateTimetimestamp 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) - aGenServername 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 is5000.
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
@type stat_key() :: :total_writes | :total_errors | :total_bytes
@type stats() :: %{required(stat_key()) => non_neg_integer()}
@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
@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.
@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 BatchWritertimeout-GenServer.callwait bound in ms (default:60_000)
Examples
iex> InfluxElixir.Write.BatchWriter.flush(pid)
:ok
@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
@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]
@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 BatchWriterpayload- aInfluxElixir.Write.Point.t()or pre-encoded binarytimeout-GenServer.callwait bound in ms (default:60_000)
Examples
iex> InfluxElixir.Write.BatchWriter.write(pid, "cpu value=1.0")
:ok
@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 BatchWriterpayload- aInfluxElixir.Write.Point.t()or pre-encoded binarytimeout-GenServer.callwait bound in ms (default:300_000). Pass:infinityfor unbounded blocking.
Examples
iex> InfluxElixir.Write.BatchWriter.write_sync(pid, "cpu value=1.0")
:ok