Analytics: OLAP, stream and batch
Row vs columnar storage, batch vs streaming, lambda and kappa, windows and watermarks, late data, and what "exactly once" really covers.
Analytics design is how you answer big questions without making the product database do reporting work. This page gets you ready to choose OLTP, OLAP, batch, stream, pre-aggregation and lakehouse answers honestly.
Read this if your last attempt…
- You used the product database for dashboards and could not explain the blast radius
- You said "use Kafka and Flink" without naming windows, lateness or sink semantics
- You confuse row stores with columnar stores
- You promise exactly-once analytics without scoping the guarantee
The concept
Analytics systems answer aggregate questions over many events: conversion, revenue, fraud and error rates. OLTP systems answer transactional questions for one user or one entity. Keep those jobs separate.
OLTP versus OLAP
OLTP is the product path: short reads and writes, invariants, and user-facing latency. A checkout update should not wait behind a dashboard scan. Row-oriented SQL fits because a row keeps one entity together.
The product database serves the user path. Events and changes flow to stream or batch processing, then land in an OLAP store for dashboards and investigation.
Analytics choices: pick by freshness, mutability and query shape.
| Choice | Use it when | Watch out for |
|---|---|---|
| OLTP row store | Transactions, point reads, invariants and small updates on the product path. | Dashboard scans steal cache, locks or I/O from users. |
| Columnar OLAP store | Large scans over a few columns, GROUP BY queries and historical investigation. | Point updates, deletes and tiny trickle writes need deltas, merges or compaction. |
| Batch pipeline | Freshness can wait and replayability is more important than seconds-old results. | Dashboards lag behind reality and incident alerts arrive too late. |
| Stream pipeline | Live counts, alerts, fraud signals, sessions or fresh materialised views. | State, watermarks, late data and sink idempotency become design work. |
| Lambda architecture | You need a real-time layer and a separate batch layer is already the reliable path. | Two code paths can drift and query-time merging is a source of bugs. |
| Kappa architecture | One retained log can feed both live results and reprocessing by replay. | Retention, replay speed and schema evolution must support your backfill horizon. |
| Pre-aggregation | Hot dashboards need cheap predictable reads over known dimensions. | Corrections and new dimensions require backfills or new aggregate tables. |
| Query-time aggregation | Analysts need flexible slicing and the scan cost is acceptable. | Hot product surfaces can become slow or expensive. |
How interviewers grade this
- You keep user-facing OLTP traffic away from broad analytical scans.
- You choose row or columnar storage from the read and write pattern, not from brand names.
- You state freshness: batch, stream, or a hybrid, and what stale means to users.
- You name window type, event time, watermarks and late-data policy for stream metrics.
- You scope exactly-once to source, processor, state and sink rather than promising it globally.
- You separate pre-aggregated hot metrics from flexible query-time analysis.
Variants
Batch warehouse
Load events into an OLAP store on a schedule, then query with SQL.
The simplest correct analytics shape. Events land in object storage or a staging table, a scheduled job cleans and models them, and analysts query the warehouse. Use it for reporting, finance, experimentation and backfills where minutes or hours are fine. The cost is freshness. A dashboard may show yesterday or a recent partition, not the last click.
Pros
- +Simple to reason about and replay.
- +Works well with SQL, governance and analyst tools.
- +Cheap for large historical jobs when scheduled carefully.
Cons
- −Not suitable for live alerting.
- −Late data needs partition repair or restatement.
- −A bad model can require a full backfill.
Choose this variant when
- Daily reporting.
- Analyst-driven ad hoc queries.
- Historical backfills and model changes.
Streaming analytics
Compute continuously from an event log into fresh tables or serving views.
Use a stream processor when freshness is part of the product. Kafka carries the events; Flink or a similar engine windows, joins, enriches and writes results. The design work is the time model: event time, watermark strategy, window type, allowed lateness and sink semantics. Keep raw events so you can repair results when code or lateness policy changes.
Pros
- +Fresh results for dashboards, alerts and fraud.
- +Can update materialised views incrementally.
- +Keeps compute close to the event stream.
Cons
- −Stateful jobs need checkpointing and operational ownership.
- −Late data creates corrections or side outputs.
- −External sinks still need idempotent or transactional writes.
Choose this variant when
- Live dashboards.
- Fraud or anomaly detection.
- Sessionisation or near-real-time feature generation.
Lambda architecture
Batch layer for complete history, speed layer for fresh results.
Lambda is honest when the batch path is the source of truth and the speed path covers the freshness gap. It is also expensive because the transformation logic and operations often exist twice. Use it when the batch and stream systems are deliberately different, such as complete nightly billing plus a provisional live counter.
Pros
- +Strong historical backfill story.
- +Fresh-enough speed layer for product surfaces.
- +Can keep mature batch tooling.
Cons
- −Two implementations can drift.
- −Query serving must merge batch and speed results.
- −Operational burden is high.
Choose this variant when
- Existing batch platform is reliable.
- Live result can be provisional.
- Historical recompute is too large for the stream path alone.
Kappa architecture
One stream processing path, replayed from a retained log when needed.
Kappa removes the duplicate batch implementation. Retain the immutable input log, deploy a new version of the stream job, replay from the beginning of the retained horizon into a new output, validate, then switch readers. This works when retention is long enough, replay is fast enough and the stream processor can handle historical throughput.
Pros
- +One code path for live and reprocessed results.
- +Replay uses the same semantics as live processing.
- +Simpler serving layer than lambda.
Cons
- −Retention can be expensive.
- −Large historical recomputes may take too long.
- −Schema evolution in the log must be disciplined.
Choose this variant when
- You already retain the event log.
- Same computation works live and historical.
- Operational simplicity matters more than separate batch optimisation.
Warehouse versus lakehouse
Managed SQL service versus open table formats on object storage.
A warehouse is the easiest interview default for analytics: load tables, govern access, run SQL and let the service own much of the tuning. A lakehouse keeps data in open files and table formats on object storage, then lets several engines query or train on it. The lakehouse buys openness and ML or batch flexibility, but you inherit table maintenance, compaction, catalog and access-control work.
Pros
- +Warehouse: fast managed path for SQL analytics.
- +Lakehouse: open storage and multiple compute engines.
- +Both can separate analytics from OLTP.
Cons
- −Warehouse can be costly or lock you into one service.
- −Lakehouse operations are easy to understate.
- −Neither removes modelling, quality and governance work.
Choose this variant when
- Warehouse for standard BI and SQL-first teams.
- Lakehouse for open data, ML and multi-engine processing.
Worked example
Numbers in this section are illustrative.
Scenario: design analytics for a ticket platform. Numbers in this section are illustrative.
Setup. The product path sells tickets and records payments. The analytics path needs live sale velocity, a fraud signal within a few seconds, and daily finance reports. At launch, assume 10k checkout events per second during a hot sale and 200M historical events kept for analysis.
Separate OLTP and OLAP. Checkout writes orders, payments and seat holds to the OLTP database. It also writes an outbox event in the same transaction. A CDC or outbox relay publishes those events to Kafka. The user path does not run GROUP BY over orders.
Live stream. Kafka feeds a Flink job keyed by event_id and seller_id. It computes a tumbling 1-minute sale-rate window for dashboards and a sliding 5-minute failure-rate window for fraud. The job uses event time because mobile and payment events can arrive late. Watermarks allow, for example, 2 minutes of lateness. Events later than that go to a side output for audit and a daily correction job.
Serving views. Hot dashboards read a pre-aggregated table keyed by event_id + minute, written with idempotent upserts. Analysts query raw facts in the warehouse. Daily finance uses a batch job from the same event log and payment tables, then reconciles against the ledger.
Backfill. If the fraud rule changes, replay the retained log into a new aggregate table and validate counts for an illustrative day before switching readers. If retention is too short, use the lake or warehouse copy as the historical source.
Answer to say. "OLTP owns seat and payment correctness. Kafka carries committed events. Flink computes fresh aggregates with event-time windows and a late-data policy. OLAP keeps raw facts plus hot aggregates. Exactly-once is scoped: checkpoints and offsets cover the stream job, while the dashboard sink is safe because writes are idempotent upserts by aggregate key."
Good vs bad answer
Numbers in this section are illustrative.
Interviewer probe
“Your product database has 500M order events and the CEO wants live revenue dashboards. What do you build?”
Weak answer
"I will add read replicas and run SQL aggregates on the orders table every few seconds. If that is slow, I will cache the dashboard."
Strong answer
"I keep OLTP for checkout and move analytics off the user path. Each committed order writes an outbox event to Kafka. For live revenue, a stream job keys by merchant and minute, uses event time with watermarks, and writes idempotent upserts into an OLAP serving table. The warehouse keeps raw facts for ad hoc analysis and daily reconciliation. I would not claim exactly-once to every sink: Kafka/Flink cover source offsets and state, and the external table is safe because the write key is merchant + minute + metric."
Why it wins: It protects OLTP, names the eventing path, separates pre-aggregation from ad hoc analysis, uses event time and watermarks, and scopes exactly-once to the parts that actually provide it.
When it comes up
- The prompt asks for dashboards, metrics, event analytics, BI or fraud signals.
- The interviewer asks whether the product database can serve reports.
- You need to choose batch, streaming or both.
- Late and out-of-order events would change the answer.
Order of reveal
- 11. Split OLTP and OLAP. Keep the product database for transactions and move analytical scans to an OLAP path.
- 22. Choose row or columnar by access pattern. Rows are good for point updates and entity reads; columnar is good for scanning selected columns across many rows.
- 33. State freshness. If minutes or hours are fine, batch is simpler. If the product needs seconds-old answers, use a stream job.
- 44. Window live metrics. For streaming metrics I name the window, the time model, watermark and late-data policy.
- 55. Scope correctness. Flink checkpoints and Kafka transactions have real but bounded guarantees; external sinks still need idempotency or transactions.
Signature phrases
- “OLTP protects invariants; OLAP scans facts.” — Separates the two workloads in one sentence.
- “Freshness is a product requirement, not a technology choice.” — Forces batch versus stream from user value.
- “Event time for when it happened; processing time for when we saw it.” — The cleanest windowing distinction.
- “Pre-aggregate the hot questions and keep raw facts for the questions we have not met yet.” — Balances serving cost with analytical flexibility.
Likely follow-ups
?“How do you handle late events after a dashboard already showed a number?”Reveal
Pick a product policy. For a live dashboard, update the same aggregate row when the late event is within allowed lateness and show that recent windows may still settle. For finance, keep an append-only correction and restate the report in batch. Events beyond the bound go to a side output or audit topic so they are visible rather than silently dropped.
?“When do you pre-aggregate?”Reveal
Pre-aggregate when a metric is hot, stable in shape and expensive to compute at read time, such as revenue per merchant per minute. Keep query-time aggregation for exploratory dimensions where analysts may ask new questions. The mistake is pre-aggregating every possible slice before you know which are used.
?“How would you reprocess after a bug in the stream job?”Reveal
Keep raw events. Deploy a fixed job that writes to a new output table from the retained log or from historical lake files. Validate totals against known partitions, switch readers, then retire the bad table. If you cannot replay the source, the analytics design is fragile.
Code examples
CREATE TABLE revenue_by_event_minute (
event_id BIGINT NOT NULL,
bucket_minute TIMESTAMPTZ NOT NULL,
currency TEXT NOT NULL,
gross_cents BIGINT NOT NULL,
order_count BIGINT NOT NULL,
updated_at TIMESTAMPTZ NOT NULL,
PRIMARY KEY (event_id, bucket_minute, currency)
);
-- Stream job writes idempotent upserts keyed by the aggregate identity.
INSERT INTO revenue_by_event_minute
(event_id, bucket_minute, currency, gross_cents, order_count, updated_at)
VALUES ($1, $2, $3, $4, $5, now())
ON CONFLICT (event_id, bucket_minute, currency)
DO UPDATE SET
gross_cents = EXCLUDED.gross_cents,
order_count = EXCLUDED.order_count,
updated_at = now();WatermarkStrategy<OrderEvent> wm =
WatermarkStrategy
.<OrderEvent>forBoundedOutOfOrderness(Duration.ofSeconds(30))
.withTimestampAssigner((event, ts) -> event.eventTimeMillis());
DataStream<Revenue> revenue =
orders.assignTimestampsAndWatermarks(wm)
.keyBy(event -> event.eventId())
.window(TumblingEventTimeWindows.of(Time.minutes(1)))
.allowedLateness(Time.minutes(2))
.sideOutputLateData(tooLateTag)
.aggregate(new RevenueAggregate());Common mistakes
A report that scans months of rows can steal I/O and cache from checkout or login. Read replicas help, but the better shape is a separate OLAP path with its own freshness and backfill contract.
Columnar layout is built for scans and compression. If the core workload is small point updates, returns, status flips and transactional constraints, keep the canonical record in OLTP and project facts into OLAP.
Arrival time is not the same as when a purchase, click or payment happened. Processing-time windows mis-bucket delayed events. Use event time for analytics unless the question is explicitly about arrival.
A watermark is not a guarantee that no late data exists. State the allowed lateness, whether late records update results, and where very late events go.
Lambda can be useful, but two implementations need golden tests and reconciliation. If the batch and speed layers disagree, users trust neither number.
A checkpointed stream job can still replay a write to an external system. Make the sink an idempotent upsert, a transactional sink, or a repairable append with dedupe.
Practice drills
Numbers in this section are illustrative.
Why is a row store a poor default for an analytics dashboard over months of events?Reveal
A row store reads whole rows and is optimised for transactional access. A dashboard usually scans a few columns across many rows and groups them. A columnar OLAP store can read only the needed columns, compress repeated values and keep the scan away from the product database.
A mobile app sends events late after being offline. Which time model do you use?Reveal
Use event time if the metric is about when the user action happened. Processing time would put the event in the bucket where the server received it. Add a watermark and allowed-lateness policy so the system can trade result latency against completeness.
When would lambda beat kappa?Reveal
Lambda can be a good fit when the historical recompute is huge, the batch path is already mature, or the live path intentionally serves a smaller provisional result. Kappa wins when a retained log and one stream job can support both live processing and replay.
Does Flink exactly-once mean the dashboard database cannot see duplicate writes?Reveal
No. Flink can recover state and source positions from checkpoints, but an external sink must cooperate. Use an idempotent upsert by aggregate key, or a transactional sink tied to the checkpoint, or accept append-and-dedupe repair.
Deep dives
Late data policy in depth
Late data is where many streaming answers become vague. A watermark is a practical statement that event time has progressed, not a promise that no older event can ever appear. The policy should match the business value of the metric. A live operations dashboard may update a recent window for a short allowed-lateness interval and then send very late records to an audit stream. A billing report may keep corrections forever and restate closed periods through batch reconciliation. A fraud signal may ignore very late events for online action, but still persist them for model training and investigation.
Say the policy out loud: "We allow an illustrative two minutes of lateness for live dashboards, emit correction updates while the window is retained, and write later events to a side-output topic for daily reconciliation." That one sentence covers latency, correctness and operational repair. It also prevents the common mistake of treating dropped late events as if they never existed.
Pre-aggregation versus query-time aggregation
Pre-aggregation is a serving decision. If a product surface asks the same expensive question all day, store the answer at the grain the product needs. Query-time aggregation is a flexibility decision. If analysts are exploring new slices, keep raw or lightly modelled facts and pay the scan cost when they ask. Most mature systems use both: raw facts for truth and backfills, modelled tables for common analysis, and small aggregate tables for hot product views.
The trade-off is change. If you pre-aggregate by country and device, a new request for campaign and browser needs a new table or a backfill. If you aggregate at query time, the same question is easy to ask but may be too slow or too expensive for every page view. In interviews, pick the known hot metrics for pre-aggregation and leave exploratory analysis in the warehouse or lakehouse.
State, checkpoints and savepoints
A stream job is a long-running state machine. The late data policy decides what happens when event time is messy; checkpoints decide whether the job can survive failure without losing its remembered state. This agrees with the page's exactly-once section: Flink can recover its own state and source positions, but the sink still has to be idempotent or transactional.
Keyed state and windows
Flink's state docs say keyed state is scoped to the key of the current input element and requires a keyed stream. Jobs partition by key so all events for the same user, ad or auction reach the same logical state, while different keys run in parallel. Session windows are gap-based: Flink's window docs describe session assigners as creating windows separated by inactivity gaps, and note that session windows merge when events bridge the gap.
Checkpoints and replay
Flink's checkpoint docs say checkpoints recover state and corresponding stream positions, giving the application the same semantics as failure-free execution. The mechanism comes from asynchronous barrier snapshotting: Carbone, Fóra, Ewen, Haridi and Tzoumas's 2015 ABS paper describes checkpoint barriers moving through a distributed dataflow so operators capture a consistent snapshot without stopping the whole job. During aligned checkpoints, Flink's large-state tuning docs say channels that already received a barrier are blocked until the remaining channels catch up. Flink's checkpointing-under-backpressure docs say unaligned checkpoints include in-flight buffered data so barriers can overtake buffers; that can reduce checkpoint time under backpressure, but adds state-storage I/O.
After failure, the job restores from the last completed checkpoint and replays input from the recorded source position. That is why the source must be replayable: Kafka's consumer-position design says a consumer position is the offset of the next message in each partition and can be checkpointed; it also says consumers can rewind to an old offset and re-consume data.
Savepoints and state backends
Checkpoints are the automatic failure-recovery log. Flink's checkpoints-versus-savepoints docs say checkpoints are created, owned and released by Flink for unexpected failures and stored in backend-specific native format. Savepoints are the operator-controlled handoff: the same docs say savepoints are created, owned and deleted by the user, use a backend-independent canonical format by default, and support planned upgrades and rescaling. Flink's savepoint docs also describe savepoints as consistent images used to stop and resume, fork or update jobs, with a native format available when you choose that trade-off.
State backend choice sets the cost shape. Flink's state backend docs say HashMap state keeps objects on the Java heap, while EmbeddedRocksDB stores serialized bytes on local RocksDB and, among those two, supports incremental checkpoints. Shorter checkpoint intervals reduce replay after failure, but add more checkpoint I/O. Longer intervals reduce steady overhead, but increase recovery time because more input must replay.
Cheat sheet
- •OLTP = small transactional reads and writes. OLAP = large scans and aggregates.
- •Rows keep entities together; columns keep attributes together.
- •Columnar scans and compresses well, but point updates need deltas or rewrites.
- •Batch when freshness can wait; stream when the product needs live answers.
- •Windows: tumbling, sliding and session.
- •Event time answers when it happened; processing time answers when we saw it.
- •Watermarks close windows, not reality. State allowed lateness and side-output policy.
- •Lambda = batch plus speed layers. Kappa = one stream path plus replay.
- •Pre-aggregate hot known metrics; query raw facts for exploration.
- •Exactly-once claims stop at the source, state and sink contract.
Practice this skill
These problems exercise Analytics: OLAP, stream and batch. Try one now to apply what you just learned.
Read this if