Event-Driven Analytics Architecture Explained

Learn event-driven analytics: producers, a broker, consumers, schemas, delivery guarantees and replay, with a small runnable Node.js example for Acme Shop.

Your collector works. It validates an event, inserts a row into PostgreSQL and returns 202. Then the product team asks for a Slack alert when a purchase fails, the data team wants a copy in a warehouse, and someone wants a live “visitors now” counter. Each request adds another line to the collector’s request handler, and each line can slow down or break event capture. That is the moment event-driven analytics starts to pay for itself.

The idea is simple. Instead of the collector calling every downstream system, it appends the event to a durable log. Independent consumers read that log at their own pace. Capture stays fast, and new uses of the data cost you a new consumer, not a change to the collector.

This article explains the pattern as an analytics problem: producers, a broker, consumers, event schemas, delivery guarantees and replay. You get a diagram, a comparison of broker options, and a small runnable Node.js slice that shows an append-only log, two consumers and a replay.

It builds on the async ingestion article, which gave you an in-process buffer. This article is the conceptual step beyond that buffer: the full event-driven analytics architecture. Moving the data into a warehouse is a separate topic, covered in the warehouse migration article.

Executive Summary: Event-driven analytics separates capturing an event from using it. Producers write events to a durable, ordered log, and consumers read that log independently, tracking their own position. Delivery is at-least-once in practice, so every consumer must be idempotent, and the event_id you already generate makes that straightforward. The reward is cheap fan-out, isolated failures and the ability to replay history into a fixed or brand new consumer.

The Problem at Scale: A Collector That Knows Too Much

Here is the shape most collectors take after a year of feature requests.

// ILLUSTRATIVE: the coupled collector
app.post('/v1/events', async (req, res) => {
  const event = validate(req.body);
  await insertEvent(event);          // PostgreSQL
  await updateDailyRollup(event);    // more SQL
  await notifySlackIfFailure(event); // third-party HTTP call
  await pushToWarehouse(event);      // another HTTP call
  res.status(202).end();
});

Every call in that handler is a failure point. If Slack is slow, event capture is slow. If the warehouse call throws, you either lose the event or return an error and trigger client retries, which then duplicate the row. The latency of your slowest dependency becomes the latency of your collector.

In my experience building event pipelines, the worst version of this is the partial failure. The insert succeeds, the rollup update fails, and the client retries. Now the raw table has two rows and the rollup has one, and the numbers disagree for a month before anyone finds the cause. The root problem is that one request is trying to do four jobs atomically across four systems.

Event-driven design fixes this by splitting the work. The collector does one thing, which is to accept and durably record the event. Everything else happens later, elsewhere, and can fail without affecting capture.

The Event-Driven Analytics Pattern in One Diagram

  PRODUCERS            BROKER (durable log)            CONSUMERS
  ---------            ---------------------            ---------

  browser SDK  --\                                   /-> writer     -> events table
                  \     +-------------------+       /
  collector  -----+-->  | topic: events     |  ----+--> rollup     -> daily_pageviews
                  /     | offsets 0,1,2,3.. |       \
  server jobs  --/      | retained 7 days   |        \-> alerts     -> Slack webhook
                        +-------------------+         \
                                                       \-> warehouse -> ClickHouse/BigQuery

  Each consumer keeps its own offset. Slow consumers never block fast ones.

Read the diagram from left to right. Producers append. The broker keeps events in order and does not delete them when a consumer reads them. Each consumer remembers how far it has read, and moves forward at its own speed.

Components and Their Responsibilities

Producers

A producer is anything that creates events: the browser SDK, the Express collector, a cron job that emits a nightly summary. The collector remains the trust boundary. It still validates every event, enforces the CORS allow-list and the payload size limit, and only then publishes to the broker. Never let browsers publish to a broker directly.

The broker

The broker stores events in an ordered, append-only log and hands them to consumers. Two properties define it. Events are retained after being read, and every event has a position (an offset). Those two properties are what make replay possible.

The word “topic” names a stream of related events, such as events. Large brokers split a topic into partitions, so many consumers can read in parallel. Order is guaranteed only within a partition. The Apache Kafka introduction explains topics, partitions and the per-partition ordering guarantee.

Consumers

A consumer reads events and performs one job: write a row, update a rollup, call a webhook, load a warehouse batch. Consumers that share a job form a consumer group, and the broker splits partitions among the group members. A different job gets a different group, so each group sees every event.

Keep each consumer small and single-purpose. A consumer that writes the table and sends alerts brings back the coupling you removed.

Event Schemas: The Contract Between Teams

Once producers and consumers are decoupled, the event shape is the only thing they share. Treat it as an API. The canonical Acme Shop event already has the right base: event_id, event_name, occurred_at, the identifiers and a free-form properties object.

Add one field to the message when you adopt a broker: schema_version, an integer. This extends the event definition used since the collector article, but only on the log. The collector stamps it after validation, and the writer consumer drops it before the insert, so the events table keeps the shape from the schema article. Validation still happens at the edge, so the log contains only events that passed it.

Follow three evolution rules:

  • Adding an optional field is safe. Old consumers ignore it.
  • Renaming or removing a field is a breaking change. Publish both versions during a migration.
  • Never change the meaning of a field in place. A revenue field that switches from cents to dollars corrupts every historical aggregate silently.

A mistake I have seen in production is a producer changing a property type, such as cart_value going from a string in one release to a number in the next. The consumer that cast it to a number started writing nulls. Consumers should be strict about what they accept and send bad events to a dead-letter location, not drop them.

Delivery Guarantees and the Idempotent Consumer

Brokers offer three delivery models, and the differences are not academic.

Model What happens on failure Analytics consequence
At-most-once Event may be lost, never repeated Undercounted revenue and funnels
At-least-once Event is never lost, may be repeated Duplicates unless the consumer dedupes
Exactly-once Each event affects state once Only within a bounded system, with caveats

For analytics, choose at-least-once and make consumers idempotent. Exactly-once claims usually hold inside one system, such as a broker transaction, and stop at the boundary where you write to your own database or call a webhook. Idempotent consumers work everywhere.

The client already generates a UUID per event, as covered in the async article. That event_id is your idempotency key. A consumer that inserts into PostgreSQL uses it to ignore repeats.

-- Assumes events.event_id has a unique constraint, as in the schema article
INSERT INTO events (event_id, event_name, occurred_at, anonymous_id,
                    user_id, session_id, page_url, referrer, user_agent, properties)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10)
ON CONFLICT (event_id) DO NOTHING;

The unique constraint comes from the PostgreSQL schema design article. With ON CONFLICT DO NOTHING, a redelivered event is a harmless no-op, so a crash between “write row” and “save offset” cannot create a duplicate.

Side effects need the same care. A Slack alert is not naturally idempotent, so record a “sent” marker keyed by event_id before or atomically with the call. Exactly how to design those actions belongs to the automated actions article.

Ordering and Partition Keys

Order matters for sessions and funnels. A purchase_completed processed before its add_to_cart breaks any consumer that builds state per user. The broker only promises order within a partition, so choose the partition key to keep related events together.

For Acme Shop, key by anonymous_id. All events from one browser then land in one partition, in the order the collector accepted them. Keying by event_name is the classic bad choice, because it puts every page_view in one hot partition and scatters each user’s journey.

Note the limit: client clocks and retries mean occurred_at order and arrival order differ. Consumers that compute sessions should sort by occurred_at within a window, not trust arrival order.

A Small Working Slice: Log, Consumers and Replay

You do not need a cluster to learn the model. The code below is a complete, runnable Node.js example using only the standard library. It implements an append-only file as the log, stores each consumer’s offset in its own file, and shows replay by resetting an offset. It is a teaching model, not a production broker.

// log.js - a minimal append-only event log (teaching model)
import { appendFileSync, readFileSync, writeFileSync, existsSync } from 'node:fs';

const LOG = './events.log';

export function append(event) {
  appendFileSync(LOG, JSON.stringify(event) + '\n');
}

function readAll() {
  if (!existsSync(LOG)) return [];
  return readFileSync(LOG, 'utf8').split('\n').filter(Boolean).map((l) => JSON.parse(l));
}

export function getOffset(group) {
  const f = `./offset.${group}`;
  return existsSync(f) ? Number(readFileSync(f, 'utf8')) : 0;
}

export function commit(group, offset) {
  writeFileSync(`./offset.${group}`, String(offset));
}

export async function consume(group, handler) {
  const all = readAll();
  let offset = getOffset(group);
  for (; offset < all.length; offset++) {
    await handler(all[offset]);   // may throw: offset not advanced
    commit(group, offset + 1);    // commit after success: at-least-once
  }
}

Notice where the commit happens. Saving the offset after the handler succeeds gives at-least-once delivery. If the process dies between the two lines, the event is handled again on restart, which is why the handler must be idempotent.

// demo.js
import { append, consume, commit } from './log.js';
import { randomUUID } from 'node:crypto';
import { rmSync } from 'node:fs';

// Start from a clean slate so reruns give the same output
for (const f of ['./events.log', './offset.rollup', './offset.alerts']) rmSync(f, { force: true });

const names = ['page_view', 'add_to_cart', 'purchase_completed'];
for (const event_name of names) {
  append({
    event_id: randomUUID(),
    event_name,
    occurred_at: new Date().toISOString(),
    anonymous_id: 'anon-1',
    schema_version: 1,
    properties: {},
  });
}

// Consumer 1: counts events per name (stands in for a rollup)
const counts = {};
await consume('rollup', async (e) => {
  counts[e.event_name] = (counts[e.event_name] ?? 0) + 1;
});
console.log('rollup counts', counts);

// Consumer 2: alerts on purchases only (stands in for a webhook)
await consume('alerts', async (e) => {
  if (e.event_name === 'purchase_completed') console.log('ALERT purchase', e.event_id);
});

// Replay: reset the rollup offset and rebuild the counts from history
commit('rollup', 0);
const rebuilt = {};
await consume('rollup', async (e) => {
  rebuilt[e.event_name] = (rebuilt[e.event_name] ?? 0) + 1;
});
console.log('rebuilt counts', rebuilt);

Save both files in an empty folder with a package.json containing {"type": "module"}, then run node demo.js. The demo deletes its own log and offset files first, so every run prints the same result. Both consumers read the same three events, yet neither affects the other. The final step is the important one: resetting one offset to zero rebuilt a derived view from raw history without touching the producer.

Replay: The Feature That Changes How You Work

Replay is the event-driven analytics feature that turns bugs from permanent damage into a rerun. If a rollup consumer miscounted for two days, you fix the code, reset its offset to the start of the bad window and let it reprocess. The rerun must replace the old values, not add to them. A consumer that increments a counter per event is not idempotent, so write rollups as “recompute the affected days and overwrite” with an upsert.

Replay also lets you add consumers retroactively. Want a funnel table you did not plan for in January? Start a new consumer group in March at offset zero. The broker already holds the history, as long as retention covers it.

The cost is retention storage and the discipline to keep consumers deterministic. A consumer that calls Date.now() or an external API during processing gives different results on replay. Use the event’s occurred_at, never the wall clock.

Failure Modes You Must Plan For

  • Poison events. One malformed event makes a consumer throw forever and blocks its partition. Retry a bounded number of times, then move the event to a dead-letter store and continue.
  • Consumer lag. A slow consumer falls behind. Monitor the gap between the newest offset and the consumer’s offset. Lag is the primary health metric of this architecture.
  • Broker unavailable. Producers need a fallback, such as the bounded in-process buffer from the async article, with a clear policy on what to drop when it fills.
  • Retention shorter than recovery time. If a consumer is down for longer than the retention window, events expire unread. Size retention to your worst realistic outage.
  • Hot partitions. A bad partition key sends most traffic to one partition and caps your throughput at one consumer.

Choosing a Broker

The plan for this series is to start small and name the bigger options only when you outgrow the smaller ones. This table follows that order.

Option Strength Cost Fits when
In-process buffer No infrastructure Events lost if the process dies; one consumer One app, low volume, loss of a few events is acceptable
PostgreSQL as a queue Reuses your database; transactional Competes with analytics queries for the same resources; a claim-and-delete queue has no replay Small teams with one job to run; use FOR UPDATE SKIP LOCKED to claim rows
Redis Streams Consumer groups, acknowledgements, simple to run Memory-bound retention unless trimmed carefully Several consumers, moderate volume, short replay window
Apache Kafka Long retention, partitioned parallelism, mature ecosystem Real operational weight High volume, many consumer groups, long replay needs
Managed queues (SQS, Pub/Sub, Kinesis) No servers to run Vendor semantics; replay and retention vary by product You are already on that cloud

Check each product’s current documentation for retention limits and ordering guarantees, because they differ and they change. The Redis Streams documentation describes consumer groups and acknowledgement, and the PostgreSQL SELECT reference covers SKIP LOCKED in its locking clause section.

How Real Systems Do This

Snowplow documents a pipeline of collector, enrichment and loader stages connected by streams, and it supports more than one streaming backend, such as Kinesis. The shape matches the diagram above: a fast collector in front, independent processing behind. PostHog’s open-source repository includes a Kafka service in its Docker Compose setup, which reflects the same idea of a broker between capture and storage.

Segment follows the same pattern at the product level, with one ingestion API and many destinations. The lesson is not to copy their stacks. It is that all of them keep capture simple, durable and separate from everything that interprets events.

Adoption Path

  1. Start with the in-process buffer and idempotent inserts keyed by event_id.
  2. Add a schema_version field and a dead-letter table for rejected events.
  3. When a second consumer appears, introduce a shared log. Redis Streams gives you one with consumer groups. In PostgreSQL, keep events in a table with an identity column and give each consumer its own stored offset, because SKIP LOCKED only claims work and cannot fan out or replay.
  4. Split the writer, rollup and alert jobs into separate consumers with separate offsets.
  5. Move to Kafka or a managed service only when volume, retention or consumer count demands it.

Each step solves a pain you already feel. Skipping ahead buys operational burden without a matching benefit.

Decision Framework

  1. Do you have, or will you soon have, more than one consumer of the same events? If not, a direct write is fine.
  2. Can you tolerate a delay of seconds between capture and visibility? Event-driven systems are eventually consistent.
  3. Do you need replay to fix bugs or add views later? If yes, choose a broker with retention.
  4. What is your peak event rate, and how long must history stay replayable? Those two numbers pick the broker.
  5. Who will run it at 3 a.m.? Pick the option your team can operate, not the one with the best diagram.

When NOT to Use This

  • One app, one database, low volume. A collector inserting rows with idempotent keys is simpler and fully adequate. A broker adds a moving part with no payoff.
  • You need read-your-writes consistency. If a user must see their own action in a report instantly, asynchronous consumers add a lag you must design around.
  • You lack on-call capacity. Running Kafka without people who understand it is worse than a bigger Postgres. Consider a managed queue or a hosted analytics product, since buying often beats building here.

Common Mistakes

  • Committing the offset before the work finishes, which loses events when a consumer crashes mid-batch.
  • Skipping idempotency because “the broker is exactly-once,” then double-counting purchases after a redelivery.
  • Partitioning by event_name, which creates a hot partition and scrambles per-user order.
  • Changing a field’s meaning without a new schema version, which silently corrupts historical aggregates.
  • Letting a poison event block a consumer instead of routing it to a dead-letter store.
  • Using wall-clock time inside a consumer, so replays produce different numbers than the original run.

Key Takeaways

  • Separate capturing an event from using it, and let the collector only validate and append.
  • Choose at-least-once delivery and make every consumer idempotent with event_id.
  • Version your event schema, and treat it as an API between teams.
  • Partition by anonymous_id to keep each visitor’s events ordered.
  • Monitor consumer lag, and send poison events to a dead-letter store.
  • Use replay to repair bugs and to build new views from history.
  • Climb the ladder: buffer, database queue, Redis Streams, then Kafka or a managed service.

FAQ

What is event-driven architecture in analytics?

It is a design where the collector writes events to a durable log and separate consumers read that log to store rows, update rollups, send alerts or load a warehouse. Each consumer works independently, so failures stay isolated.

Do I need Kafka for event-driven analytics?

No. Many teams start with a PostgreSQL table, using SKIP LOCKED for a single job queue, or with Redis Streams. Kafka makes sense when volume, retention and consumer count are high enough to justify running it.

How do I prevent duplicate events in a consumer?

Use the client-generated event_id as an idempotency key. A unique constraint with ON CONFLICT DO NOTHING makes repeated delivery harmless for database writes.

What is the difference between at-least-once and exactly-once delivery?

At-least-once never loses an event but may deliver it twice. Exactly-once guarantees a single effect, but usually only inside one system. Idempotent consumers give you the practical equivalent everywhere.

Can I replay old events into a new consumer?

Yes, as long as the broker still retains them. Start a new consumer group at the earliest offset and let it process history, keeping the logic deterministic.

Conclusion

Event-driven analytics trades a little simplicity for isolation, fan-out and replay. Adopt it when a second consumer shows up, not before, and keep each consumer idempotent and single-purpose.

Rule of thumb: capture once, interpret many times, and never let interpretation slow down capture.

Last updated on 9 October 2026.

Share this article

Leave a Reply

Your email address will not be published. Required fields are marked *