Running aggregates
The aggregate middleware keeps one state per key inside
the process and updates it with every message: counters, sums, moving averages. The message
goes on with the state written into its payload. Use it for values that depend on earlier
messages, such as “how many readings did this sensor send so far” or “what is this sensor’s
usual reading”. No database is asked, so it is much faster than a lookup that
lets a database keep the counter.
Where it fits
mq-bridge moves messages, and a few middlewares keep something about them between
messages: deduplication the keys it has seen, aggregate a value
computed from them. It is meant for a value the pipeline needs while the message is on its
way: to filter, route or alert on it, to hand it to the receiver, or to thin out what reaches
an expensive sink. It is kept small on purpose: one state per key, updated message by
message, in the order the messages arrive.
| You need | Use |
|---|---|
a running value per key that a filter, a switch or the receiver decides on | aggregate |
| the same, surviving restarts or shared by several instances | aggregate with a store |
| a counter that other applications update too | a writing lookup |
| totals for reports only, in a database that aggregates on its own | write the raw rows and let the sink aggregate |
| time windows, event time and late data, joins between streams, exactly-once state | a stream processor, with mq-bridge in front of it or behind it |
Letting the database keep the counter is the obvious alternative, and the slow one: it
costs a read and a write per message and key. In our measurements a lookup that updated
five keys per message reached about 17k msg/s on PostgreSQL and about 2k on MongoDB. With
single_writer, aggregate writes a changed state once per flush instead of once per
message and reaches 240k on the same PostgreSQL (see Performance).
Windows, event time and joins are out of scope. They need the state and the input position committed together, which a bridge between arbitrary endpoints cannot promise; see How this compares to stream processors.
A counter and a sum per key
sensor_stats:
input:
kafka: { topic: "readings", url: "localhost:9092" }
middlewares:
- aggregate:
consistency: single_writer
key: "${payload:sensor_id}"
into: sensor
expression: "{ total: (state.total ?? 0) + reading, count: (state.count ?? 0) + 1 }"
output:
kafka: { topic: "readings.enriched", url: "localhost:9092" }
{"sensor_id": 7, "reading": 20} leaves as
{"sensor_id": 7, "reading": 20, "sensor": {"total": 20, "count": 1}}. The next message for
sensor 7 continues from that state.
consistency: single_writer without a store says that the states may live in this process
only: they are gone after a restart and one instance runs the route. Leaving out both is a
startup error, so that states in memory are never the result of a forgotten store.
A CSV sink writes that nested result as the columns sensor.total and sensor.count, not as
JSON text in one cell; see CSV.
expression returns the new state. It reads payload fields by name, metadata as meta.<key>,
and state, the value it returned for the previous message with the same key. For a new key
state is null, so state.total ?? 0 gives the starting value. The language is the one
transform and filter use.
The usual form avg * 0.99 + reading * 0.01 starts at the first reading and needs hundreds
of messages to forget it. Keep the weighted sum and the weight, and divide when writing:
- aggregate:
consistency: single_writer
key: "${payload:sensor_id}"
into: sensor.avg_reading
expression: "{ s: (state.s ?? 0) * 0.99 + reading, w: (state.w ?? 0) * 0.99 + 1 }"
output: "state.s / state.w"
output shapes what the message carries; the stored state stays {s, w}. After a first
reading of 500 followed by readings of 50, this reads 52.6 after 100 messages. The usual form
still reads 216.
Compare a message with what came before
With emit: previous the message carries the state from before its own update. That is what
a rule such as “this reading is five times the sensor’s average” needs:
- aggregate:
consistency: single_writer
key: "${payload:sensor_id}"
into: sensor.avg_before
emit: previous
expression: "{ s: (state.s ?? 0) * 0.99 + reading, w: (state.w ?? 0) * 0.99 + 1 }"
output: "state.s / state.w"
For the first message of a sensor sensor.avg_before is null.
A fixed factor is a rough rule. With the built-in aggregates (see below) the message can
carry the average and the standard deviation from before its own update, and a filter or
switch behind it can ask how many deviations the reading is away:
- aggregate:
consistency: single_writer
key: "${payload:sensor_id}"
into: before
emit: previous
fields:
avg: ema(reading, 0.01)
sd: ema_stddev(reading, 0.01)
stddev and variance cover all readings of the key; the ema_ forms weight them like
ema, so old readings fade and the result follows a sensor that drifts. Both are sample
values (divided by n − 1) and null for the first message of a key, which has no spread yet.
Several keys from one message
entries updates one state per key, each with its own expression:
- aggregate:
consistency: single_writer
entries:
- key: "${payload:sensor_id}"
into: stats.sensor
expression: "{ sum: (state.sum ?? 0) + reading, n: (state.n ?? 0) + 1 }"
- key: "${payload:site_id}"
into: stats.site
expression: "{ max: max([state.max ?? reading, reading]) }"
Built-in aggregates: the fast path
For counters, sums and averages you do not need an expression. fields names the results and
what each one computes:
- aggregate:
consistency: single_writer
entries:
- key: "${payload:sensor_id}"
into: sensor
fields:
n: count
total: sum(reading)
avg: ema(reading, 0.01)
- key: "${payload:site_id}"
into: site
fields: { high: max(reading), usual: mean(reading) }
Available are count, sum, min, max, last, mean, stddev, variance, and
ema(path, alpha), ema_stddev(path, alpha), ema_variance(path, alpha). ema is the
moving average from above: it has no start bias, so the first message reads its own value. The value may be a number or a numeric string, as a CSV source delivers it. fields
is several times faster than an expression; use an expression where you need logic that the
built-ins do not cover.
A moving average by time
ema(reading, 0.01) counts messages: a sensor that reports every second forgets a value a
thousand times faster than one that reports every quarter of an hour. Give a half-life
instead of alpha and the weight follows the time between two messages of a key:
- aggregate:
consistency: single_writer
key: "${payload:sensor_id}"
into: sensor
time: measured_at
fields:
avg: ema(reading, 5m)
spread: ema_stddev(reading, 1h)
A reading counts half after five minutes, a quarter after ten. time names the payload
field with the message’s time: epoch seconds or milliseconds, or an RFC 3339 text. Without
time the clock of the process is used, and a replay then computes different values.
A message that is older than the newest one its state has seen is detected, not
corrected: it is folded as if it had arrived at that newest time and carries the metadata
mqb.aggregate.late, which names the entries concerned. Filter or route on it where late
data must not count:
- filter: 'meta["mqb.aggregate.late"] == null'
Messages with missing fields
A message that lacks a field an entry reads, or a key, cannot be folded. on_error says
what happens then:
on_error | Effect |
|---|---|
drop (default) | On an input the message is logged, acked and dropped; an HTTP caller still gets 202. On an output it fails. |
fail | On an input the message is nacked: an HTTP caller gets 500, a broker delivers it again. |
skip | Only the entries that cannot be computed are left out. The message goes on and carries mqb.aggregate.skipped with their into paths. |
- aggregate:
consistency: single_writer
on_error: skip
entries:
- { key: "${payload:sensor_id}", into: sensor, fields: { avg: "ema(reading, 0.01)" } }
- { key: "${payload:site_id}", into: site, fields: { high: max(reading) } }
A message without site_id leaves with sensor filled, without site, and with the
metadata mqb.aggregate.skipped: site.
Try a configuration without changing the states
read_only: true reads the states from the store, computes and writes the result into the
message, and stores nothing. Use it for a dry run against production states, or for a second
route that only reads what another one maintains. With emit: previous the message carries
the stored state as it is; with emit: updated what it would become.
- aggregate:
store: "postgres://localhost/telemetry/sensor_states"
read_only: true
key: "${payload:sensor_id}"
into: sensor
fields: { n: count, avg: "ema(reading, 0.01)" }
Keep the states across restarts
Without store the states are gone after a restart. Give the middleware a database and they
are kept there:
- aggregate:
store: "postgres://localhost/telemetry"
key: "${payload:sensor_id}"
into: sensor
fields: { n: count, avg: "ema(reading, 0.01)" }
PostgreSQL, SQLite (sqlite:///var/lib/mqb/agg.db?mode=rwc) and MongoDB
(mongodb://localhost/telemetry) work. The table or collection is created on start and is
named mqb_aggregate_<route> unless the URL names one.
By default (consistency: shared) several instances of the route can run against the same
store: every batch loads its states, updates them and writes them back, and a state that
another instance changed in between is computed again. That is safe and bound by the
database, about two round trips per batch, so use a large batch_size (4096 to 8192).
If exactly one instance runs the route, consistency: single_writer is several times faster:
the states stay in memory and are written to the store continuously. A batch is still only
acked once its states are stored, so no update is lost on a crash.
How it works internally
Folding a batch
The middleware works on whole batches. For each message, in order, and for each entry:
- The
keytemplate is rendered from the message. - The state of that key is looked up in a hash map; every entry has its own map.
- The new state is computed, by the expression or by the built-in aggregates.
- The state (or the previous one, or what
outputreturns) is written into the payload atinto. Fields the middleware does not write are copied byte for byte.
A message that fails in any step changes no state. Without a store this is everything:
the maps live in the process and the message goes on at once.
What a store holds
One row per key, in one table for all entries of the middleware:
| Column | Content |
|---|---|
agg_key | <into>:<key>, for example sensor:7; at most 512 characters |
state | the state as JSON text |
version | a counter, raised by one with every write |
MongoDB holds the same as a document {_id, state, v}. Every write names the version the
state was computed from and only takes effect while the row still has that version. This is
optimistic concurrency: nobody locks a state while computing, and a writer that was overtaken
notices it when it writes.
shared: the store owns the states
Per batch:
- The distinct keys of the batch are collected and loaded with one query (
WHERE agg_key IN (…), 1000 keys per statement), each with its version. - The batch is folded in memory, as above.
- Every changed state is written back with the version it was loaded with. PostgreSQL and SQLite use one transaction of upserts, so the states of a batch land completely or not at all. MongoDB uses one bulk command; each state is atomic, the batch is not.
- If a write is refused, another instance changed that key in between. The states are loaded again and the affected messages are folded again from them. MongoDB repeats only the refused keys; PostgreSQL and SQLite roll back and repeat the whole batch. After 32 attempts the batch fails and is redelivered.
- Only then do the messages go on. Nothing stays in memory for the next batch.
One instance runs its batches against the store one after the other, so it does not collide
with itself. The cost is two round trips per batch whatever its size, which is why large
batches pay off and why fields versus expression hardly matters here.
Between instances there is no order: of two batches that change the same key, the one that writes first wins and the other computes again on top of it. Counters and sums come out right either way; a moving average depends on which batch was first.
single_writer: memory owns the states
- A key is loaded from the store the first time a batch needs it. After that its state
stays in memory, until
max_keyspushes it out. - The batch is folded in memory, the changed keys are marked, and the batch gets a number.
- The messages go on at once; they do not wait for the store.
- A background task writes the marked states. As soon as one write returns it takes everything changed in the meantime, so under load one write covers many batches, and a key that changed a hundred times is written once. This is where the speed comes from.
- The ack of a batch (on an output: the return of
send) waits until a write has covered its number.
The writes carry versions here too. A refused write can only mean that somebody else writes the same states, and is treated like a failed write: the states in memory are dropped, every batch not yet stored is nacked, and the redelivered batches load their states again from the store. So a second instance is noticed, but not prevented; nothing fences it off.
What happens on a crash or a retry
Storing a state and acking a message are two steps, not one transaction. A message is never
acked before its state is stored, so no update is lost. The other way round is possible: the
state is stored, then the process dies or a later step fails, and the message is delivered
again. It then counts twice. A deduplication in front narrows that
window, but it has a store of its own and does not close it completely.
How this compares to stream processors
The building blocks are common ones: keyed state folded per message, as in Kafka Streams or
Flink; versioned conditional writes for shared; write-behind with one write for many
batches and an ack that waits for it for single_writer. Three things those systems do are
not done here, and they decide whether the middleware fits:
- No exactly-once. Kafka Streams and Flink commit the state and the input position together. Here a redelivered message counts twice, as described above.
- No key ownership. They assign each key to exactly one worker and move the state when
workers come and go. Here you choose:
sharedlets every instance write every key and pays a database round trip per batch,single_writeris fast and leaves it to you to run one instance. - No windows and no expiry by time. There is no “sum of the last hour” and no state
that ends with a window or a time to live. A moving average stands in for a window, by
message count or by time. States
leave memory only when
max_keysis reached, and rows in a store are kept for good. For fixed windows you can put the window into the key (key: "${payload:sensor_id}:${payload:hour}") and use astorewithshared, which keeps nothing in memory; the rows of past windows stay until you delete them.
If you need one of these, compute the aggregate in the database with a
lookup, or use a stream processor for that step.
What to know before relying on it
- Without
store, state lives in memory only. A restart starts from empty states and two instances do not share them. That is why it has to be asked for withconsistency: single_writer. - Memory is bounded. An entry keeps at most
max_keysstates in memory, one million unless you set it. See Memory for what happens beyond that. single_writerwith two instances is wrong. The second writer is noticed as a failed write and its batches are redelivered, but do not rely on that: run one instance.- A replayed message counts twice. Where that matters, let
deduplicationsee the message first: on an input, list it afteraggregate, because a message passes the last middleware of an input first. - The table is named after the route. Renaming a route, or a consumer that has no
id, starts with empty states. Name the table in the URL (postgres://host/db/sensor_states) to keep them. - Order matters for a moving average. Messages of a batch are folded in order. With a
route
concurrencyabove 1 the order across batches is not guaranteed. Withtime, a message older than its state is markedmqb.aggregate.late, not reordered. - A message that cannot be folded fails alone and changes no state: its payload is not a
JSON object, its key has no value, a field is missing, or an expression fails. On an input
it is logged and dropped, and the source is told it succeeded; on an output only that
message fails, so a following
dlqcan keep it.on_error: failorskipchanges that, see Messages with missing fields. - Field order. The fields of a written state object are not in the order the expression lists them.
Performance and tuning
Measured on one Apple M1 core, one route worker, states in memory unless a store is named. Read the numbers as proportions; absolute values depend on the machine and payload.
| Setup | msg/s |
|---|---|
fields, one key, three aggregates | 1.4M to 1.8M |
fields, five keys, three moving averages each | 700k |
expression, one key | 330k to 370k |
expression, five keys | 100k |
PostgreSQL store, single_writer, five keys, fields | 240k |
PostgreSQL store, single_writer, five keys, expression | 68k |
PostgreSQL store, shared, five keys | 17k |
The update itself costs about 0.3 µs per message with fields and about 2 µs with an
expression. The same counter, sum and moving average built from Redpanda Connect’s branch,
cache and mapping processors cost about 8 µs per message on the same machine.
What to change, in the order it pays off:
- Use
fieldswhere the built-ins cover it. It is four to seven times faster than an expression, because no expression runs and the payload is not parsed into a tree. Keep expressions for the entries that need logic. - With a store, pick the consistency first.
single_writeris about 14 times faster thanshared, but only correct with one instance. Undersharedthe database sets the pace, andfieldsversusexpressionhardly matters. - With a store, raise
batch_sizeto 4096 or 8192. A batch costs about two round trips undersharedwhatever its size. Without a store the batch size has little effect. - Update fewer keys per message. The cost grows with the number of
entries, not with the number of aggregates in one entry: five keys cost about three times as much as one. - Do not raise
concurrencyfor speed if the result depends on order, as a moving average does. Counters and sums stay correct; undershared, workers that update the same keys collide and compute those states again. - Mind the neighbours. A
deduplicationkeyed on a payload field adds about 0.5 µs per message, onmessage_idabout 0.2 µs. Atransformbefore or after parses the payload again; letaggregatewrite the fields you need instead of reshaping afterwards.
Memory
An entry keeps at most max_keys states in memory, one million unless you set it. When it is
full, the half that was used least recently is dropped, so a key that keeps coming stays and
a key that went quiet goes. What “dropped” means depends on where the states live:
| Setup | A dropped state |
|---|---|
no store (single_writer) | is forgotten; its key starts again from an empty state |
store, single_writer | is only dropped from memory and loaded again when its key returns |
store, shared | does not occur; nothing is kept in memory between batches |
Without a store a warning is logged the first time states are forgotten. If the number of
keys is bounded, such as sensors or customers, and below the limit, nothing is ever dropped.
If it is unbounded, such as a session id, the limit is what keeps the process alive for
weeks: set max_keys to what you can afford, or use a store so that no state is lost.
- aggregate:
consistency: single_writer
key: "${payload:session_id}"
into: session
max_keys: 200000
fields: { n: count }
A fields state takes roughly 100 to 300 bytes, so the default limit stays within a few
hundred megabytes per entry; an expression state with a large object takes more. max_keys: 0 removes the limit. Rows in a store are never removed; delete them yourself if the key is
unbounded.