Clustering
View SourceAurora Meter counts in ETS on every node and still gives you one cluster-wide number. This guide explains how, what it guarantees, and what to configure.
How it works
Every counter row on a node is {key, value, pending_flush, pending_gossip}:
valueis this node's view of the cluster-wide totalpending_flushis what this node has added since its last database flushpending_gossipis what this node has added since its last PubSub tick
Two loops keep the nodes in agreement.
Delta gossip (fast). Every :broadcast_interval (1 s by default) the
broadcaster takes each touched counter's pending_gossip and publishes one
{:aurora_meter, :deltas, node, [{key, delta}]} message on the
"aurora_meter:cluster" topic. Every other node adds those deltas to its own
value. A LiveView on node B therefore sees node A's increments within about a
second.
Delta flush and total announcement (durable). Every :flush_interval (5 s
by default) the flusher takes each dirty counter's pending_flush and writes it
to Postgres as a delta (value = value + Δ, AuroraMeter.Storage.add_counters/1).
Because every node writes only what it added, the row is the true cluster total
no matter how many nodes flush it. Postgres returns that total; the flushing node
re-bases its value on it (total + pending_flush) and publishes
{:aurora_meter, :totals, node, [{key, total}]} so every other node re-bases
too. Anything gossip missed (a dropped message, a node that seeded from the
database while another had unflushed increments) heals on the next announcement.
Announced totals are only ever applied forward: if a late announcement carries a smaller total than a node's current base, it is ignored, and the node's own next flush re-bases it unconditionally.
Guarantees
- Correct totals in the database. After every node has flushed, the row
equals the sum of every increment on every node. A hard crash loses at most
one
:flush_intervalof one node's increments, as before. - Convergent reads.
AuroraMeter.usage/2on any node is the true total minus, at most, what the other nodes added in the last:broadcast_interval(one:flush_intervalif PubSub dropped a message). - Bounded overshoot on hard limits.
reserve/3andwith_quota/4enforce against the local view, so a burst spread across N nodes can exceed a hard cap by what the other N−1 nodes admitted within one:broadcast_interval. Lower the interval to tighten it (cost: one PubSub message per node per tick), or mark the feature durable and reconcile invoices from the event log. - Idempotent flushes. A flush with nothing pending writes nothing. Re-basing never double counts: it is an absolute correction against a snapshot, so a bump that lands mid-rebase is preserved exactly.
- Failure safe. If Postgres is unreachable the flusher puts the taken deltas
back, re-marks the keys dirty, logs, emits
[:aurora_meter, :flush, :error]and tries again next interval. Nothing is lost while the node stays up.
Requirements
- A distributed
Phoenix.PubSub(the defaultPhoenix.PubSub.PG2adapter with Distributed Erlang, or the Redis adapter). This is the same PubSub you already configure for LiveView; nothing extra is needed. On a single node, or with a non-distributed PubSub, every message is local and dropped, and behaviour is exactly what it was. - Schema version 2. 0.3 adds no migration.
Configuration
config :aurora_meter,
cluster_sync: true, # default; false = per-node counters, cluster-wide tenant broadcasts
broadcast_interval: 1_000, # gossip and LiveView update cadence
flush_interval: 5_000 # durable write and total-announcement cadenceWith cluster_sync: true, tenant usage broadcasts
({:aurora_meter, :usage, ...}) are node-local: each node informs its own
LiveViews from its own converged view, so a browser never receives two slightly
different numbers from two nodes. With cluster_sync: false they fan out
cluster-wide as in 0.2.
Rolling upgrades from 0.2
0.2 nodes write absolute values and can overwrite a 0.3 node's deltas while both versions run. Prefer a full restart; if you must roll, expect per-node behaviour until the last 0.2 node is gone, after which the next announcements re-base everything.
Testing it
AuroraMeter.Test.simulate_node/3 applies deltas as if another node had
gossiped them; AuroraMeter.Test.simulate_flush/2 applies totals as if another
node had flushed. See the testing guide.
Telemetry
[:aurora_meter, :cluster, :apply] fires for every batch applied from another
node with %{count} and %{kind: :deltas | :totals, origin: node}.
[:aurora_meter, :flush] now also carries delta_sum, and
[:aurora_meter, :broadcast] carries deltas (how many were gossiped).