← All posts

Learning data streaming by building one

A student pipeline that ranked trending artists from live tweets, the ideas behind streaming, and how I'd run it on a real cluster.

In autumn 2019, a classmate and I built a live chart of the artists people were sharing on Twitter. Tweets containing a Spotify link went into Kafka. Spark Streaming looked up each track’s artist on the Spotify API, Cassandra stored the counts, and a Dash page redrew the top twenty every five seconds. It was my first streaming system.

It worked, and the recording looks convincing: bars grow, artists move up and down. Reading the code again a year and a half later, I can see that we had built the shape of a streaming pipeline without most of the things that make one trustworthy. This post starts with the ideas behind data streaming, then goes through what we built, the five things a stream has to get right, and how I’d deploy the same idea on a real cluster.

Data streaming, briefly

Most of the data I’d worked with until then sat still: a CSV, a table, a folder of files. Streaming starts from a different assumption. The data never stops arriving, and you want answers while it’s still arriving.

Batch or stream

A batch job processes a fixed chunk of data on a schedule: every night, every hour. A stream processor handles each event as it arrives and keeps its answer continuously up to date.

Batch versus stream processing The same events on two timelines. In batch, events pile up and a job produces a result at the end of each interval, so the first event waits for the next run. In a stream, each event produces an updated result right after it arrives. Batch runs every hour Stream always on time → event batch result stream result job runs first event waits for the next run
The same events, processed two ways. Batch answers on a schedule; a stream updates its answer as each event arrives, so results are seconds old instead of up to an hour.

The purpose is freshness. A batch job’s answer is as old as the time since its last run. A stream’s answer is seconds old, which matters when the question is “what’s happening now”: fraud, monitoring, or a live chart of what people are sharing. The cost is complexity. A batch job sees all its input at once and can simply be rerun if it fails; a stream must keep working, and keep being correct, while data keeps arriving.

The log in the middle

Most streaming systems put a log between the things that produce events and the things that process them. Kafka is the best-known example. A log is the simplest possible data structure: an ordered list of events that you can only append to. Each event gets a position, its offset.

A topic partition as a log Twelve numbered cells form an append-only log. A producer appends event 11 at the end. An archive job reads at offset 4 and a dashboard job reads at offset 9, each keeping its own position. ONE PARTITION OF A TOPIC: AN ORDERED, APPEND-ONLY LOG 0 1 2 3 4 5 6 7 8 9 10 11 Producer appends events oldest newest Archive jobnext offset: 4 Dashboard jobnext offset: 9 each reader keepsits own position
The log decouples writers from readers. Producers only append; each consumer reads at its own pace and remembers its offset. Rewinding the offset replays history, which a database table can't do.

That simple structure does a lot of work:

  • It decouples producers and consumers. The producer doesn’t know or care who reads. A slow consumer falls behind without slowing the producer down.
  • It absorbs bursts. When a thousand events arrive in one second, they wait in the log instead of being dropped.
  • It makes history replayable. Events stay in the log for a retention period. A consumer that crashed, or a new job that needs old data, just starts from an earlier offset.

Spreading the log across machines

One machine eventually runs out of disk, network or CPU, and one machine can die. So a topic is split into partitions, and the partitions are spread across several brokers. Each partition is also replicated: one broker holds the leader copy, and others keep followers that can take over if that broker fails.

Partitions, replication and a consumer group Three brokers each hold all three partitions: one as leader and two as copies. Each leader feeds one worker in a consumer group, so three workers read in parallel. leader: takes reads and writes copy: takes over if its broker dies Broker 1 partition 1 · copy partition 2 · copy partition 0 · leader Worker A Broker 2 partition 0 · copy partition 2 · copy partition 1 · leader Worker B Broker 3 partition 0 · copy partition 1 · copy partition 2 · leader Worker C consumer group: the three workers split the partitions
How the log scales out. A topic is split into partitions spread across brokers, and each partition is copied to other brokers so a machine can fail without losing data. Workers in a consumer group each take some partitions. Order is guaranteed within a partition, not across them.

On the reading side, a consumer group shares the work: each partition goes to one worker in the group, so adding workers (up to the number of partitions) adds throughput. The price of this parallelism is ordering. Events are in order within a partition but not across partitions. That’s why events are usually partitioned by a key: all events for the same user or the same artist land in the same partition and stay in order.

This is where streaming becomes a distributed-systems problem. Machines fail, networks pause and workers restart, and the system has to keep going without losing or double-counting data.

Processing: stateless and stateful

Some processing steps look at one event at a time: filter out tweets without a link, pull out the track id, look up the artist. These are stateless, so they’re easy to parallelise and easy to restart.

Counting is different. “Shares per artist in the last minute” needs memory of earlier events: that’s state. And because a stream never ends, aggregates are computed over windows of time:

Tumbling and sliding windows Twenty events over three minutes. Tumbling one-minute windows count 7, 7 and 6 events. Sliding one-minute windows every 30 seconds overlap, so each event is counted in two windows. time → 0:00 0:30 1:00 1:30 2:00 2:30 3:00 Tumbling 1 min blocks 7 events7 7 events7 6 events6 Sliding 1 min long,every 30 s 7 events7 8 events8 7 events7 6 events6 6 events6
A stream never ends, so aggregates are computed over windows. Tumbling windows cut time into separate blocks; sliding windows overlap, giving a smoother "last minute" that updates more often.

Two more ideas matter once there’s state. First, which time? An event has the time it happened (event time) and the time it reached the processor (processing time). They differ whenever something is delayed, and counting by the wrong one puts events in the wrong window. Second, state must survive a crash, so processors periodically checkpoint their state and their offsets to durable storage.

Engines differ in how they process: Spark Streaming, which we used, collects events into small micro-batches every few seconds, while engines like Flink process each record as it arrives. Both give you windows, state and checkpoints.

When things fail

Failures are where streaming gets hard. Suppose a worker processes an event, updates a count, and crashes before it records that it has finished:

sequenceDiagram
  participant K as Kafka
  participant W as Worker
  participant DB as Database
  K->>W: event at offset 41
  W->>DB: count + 1
  Note over W: crashes before saving "done up to 41"
  W->>K: restarted: resume from offset 41
  K->>W: event at offset 41, again
  W->>DB: count + 1 (counted twice)

Depending on when a consumer saves its position, you get one of three guarantees:

Guarantee What can go wrong How you get it
At most once events can be lost save the offset before processing
At least once events can be counted twice save the offset after processing
Exactly once, in effect nothing, if done right at least once, plus writes that are safe to repeat

In practice, most pipelines choose at least once and make their writes idempotent, so that repeating one does no harm. Writing “the count for this artist in this window is 12” is safe to repeat; “add 1 to the count” is not.

Our project had all of these pieces in miniature. Looking back, that’s what makes it useful: I can see which of these ideas it got right and which it skipped.

What we built

The pipeline as we built it Twitter stream API to a Tweepy producer to Kafka (one broker, one partition, replication 1) to Spark Streaming on one machine in 30-second batches. Spark calls the Spotify API once per tweet and writes to a single Cassandra node keyed by artist and second. Dash sums the last 60 minutes every 5 seconds. Twitter stream API keyword: spotify com language: en Producer (Tweepy) created_at + first URL only Kafka single broker 1 partition, RF 1 Spark Streaming 30 s micro-batches one machine Spotify API 1 call per tweet errors dropped Cassandra key (artist, second) single node, RF 1 Dash sum of last 60 min re-query every 5 s single point offailure or silent loss
As built. Every component ran as a single copy on one laptop; the orange labels are where data could be lost without anyone noticing.

Each stage had one job:

  • The producer listened to Twitter’s filtered stream for tweets containing “spotify com”, kept only the timestamp and the first link, and published a small JSON event to Kafka.
  • Kafka buffered those events, so a slow consumer didn’t lose tweets while it caught up.
  • Spark Streaming read from Kafka in 30-second micro-batches, pulled the track id out of each link, asked the Spotify API for the artist, and wrote one row per tweet to Cassandra.
  • Cassandra stored the rows, and Dash re-ran a SUM ... GROUP BY artist over the last hour every five seconds to draw the chart.

On paper that’s the textbook architecture: a durable log, a stream processor, a store and a view. In practice every component was a single copy. Kafka had one broker, one partition and replication factor (RF) 1; Spark ran as local[*] on one machine; Cassandra was a single node with replication factor 1. That’s fine for a demo. It also means none of the properties we chose these tools for were actually switched on.

Five things a stream has to get right

1. Know where you resume after a restart

A batch job that crashes can simply run again. A stream that crashes has to know where it stopped. Our Spark job read Kafka directly but never saved its offsets, so after a restart it started from the latest message. Anything that arrived while it was down was silently skipped. Kafka still held those tweets; we had just lost track of where we were in them.

The fix is to checkpoint offsets together with the processing state, and restart from the checkpoint rather than from “now”.

2. Make writes idempotent

The Cassandra table looked like this:

CREATE TABLE artistshare (
  created_at timestamp, artist text, count int,
  PRIMARY KEY (artist, created_at)
);

Each tweet was written as a row with count = 1, keyed by artist and by time to the second. In Cassandra, writing a row with an existing primary key overwrites it. If two people shared the same artist within the same second, the second write replaced the first.

Two shares in the same second become one row Tweet A and tweet B both share artist X at 18:50:00. Both write the key (artist X, 18:50:00) with count 1. The second write overwrites the first, so Cassandra stores one row. Tweet A18:50:00 · artist X Tweet B18:50:00 · artist X (artist X, 18:50:00)count = 1 write overwrites Stored in Cassandraartist X · 18:50:00 · 1 2 shares, 1 row
The primary key is (artist, second), so two shares of one artist in the same second land on the same row.

So the undercount was largest at exactly the moment the chart exists to show: a new release, when many people share the same artist at once. The ranking we displayed was a flattened version of the real one.

The general lesson is that a streaming sink must produce the same result no matter how many times a record arrives or how records collide. Streams deliver at least once after any failure, so writes need to be idempotent. Here that means aggregating in the stream and writing one count per (window, artist), which is safe to overwrite, or keying raw rows by tweet id.

3. Count by event time, not arrival time

We did store each tweet’s own timestamp, which was the right call. But the window was computed by the dashboard, from the dashboard’s clock:

The window behind the chart A time axis ending at now. The chart sums the 60 minutes ending 30 seconds before now. The last 30 seconds are excluded because Spark writes a batch every 30 seconds. Dash re-queries every 5 seconds. Summed on the chart: the 60 minutes ending 30 s ago summed on the chart: 60 minutes 30 s now now − 60 min − 30 s last batch still in flight, excluded Spark writes a batch every 30 s · Dash re-queries every 5 s · not to scale
The dashboard's window was anchored to its own clock. A tweet that arrived late, after its minute had already scrolled out, was never counted.

Tweets don’t arrive in order. A slow producer, a Kafka backlog or a Spark restart can deliver a tweet minutes late, and by then its minute may already have scrolled out of the window. A stream processor can handle this properly. It groups by the event’s own time, and a watermark says how long to wait for late events before closing a window. That turns “the last 60 minutes” into a definition the pipeline enforces, instead of whatever happened to arrive by the time the dashboard asked.

4. Keep slow calls off the hot path

For every tweet, Spark made a blocking HTTP call to the Spotify API. That puts an external service’s latency and rate limits directly on the stream. If Spotify slowed down, the whole batch slowed down; if a call failed, the code returned None and the record was filtered out without being counted.

The same few thousand tracks get shared again and again, so a cache in front of the API would absorb most lookups. The remaining calls can run asynchronously with a rate limiter. Failures shouldn’t vanish: they belong on a dead-letter topic, where they can be counted, inspected and replayed once the cause is fixed.

5. Decide what one event is

Every filter in a stream quietly decides what counts:

  • Language: only English tweets, as detected by Twitter.
  • First link, first artist: a tweet sharing two tracks counted once, and a collaboration counted only for whoever Spotify listed first.
  • Tracks only: albums and playlists were dropped.
  • No tweet id: a retweet was indistinguishable from an original share, and there was no way to remove duplicates.

None of these is wrong on its own, but each changes the chart, and none was visible on the page. In a stream, the definition of an event is the schema of everything downstream. It deserves a written decision, and each filter deserves a counter.

How I’d deploy it for real

Nothing about this product needs a big cluster; the stream fit comfortably on a laptop. But if it had to run continuously and survive failures, the shape changes like this:

How I would deploy it on a real cluster Two ingest services keep tweet ids and write to a three-broker Kafka cluster with six partitions and replication 3. A stream processor with several workers deduplicates, extracts every track and artist, enriches through a cache before calling Spotify, windows by event time with a watermark, and counts per window and artist. Failed lookups go to a dead-letter topic. State and offsets are checkpointed to object storage, and metrics track lag and drops. Counts go to a three-node Cassandra cluster keyed by window and artist, behind a query API that pushes updates to the dashboard. Twitter stream filtered stream API Ingest service ×2 keeps the tweet id reconnect, backoff Kafka ×3 brokers 6 partitions RF 3 Dead-letter topic failed lookups replay after a fix Stream processor · N workers Flink or Spark Structured Streaming 1 drop duplicate tweet ids 2 extract every track, every artist 3 enrich: cache first, then Spotify 4 window by event time + watermark 5 count per (window, artist) Artist cache track → artists Spotify API rate-limited, async Checkpoints state + offsets object storage Metrics consumer lag drops by reason Cassandra ×3 RF 3 key (window, artist) Query API cached top 20 Dashboard pushed updates on a miss
The same product on a real cluster. Every stateful piece is replicated, every drop is counted, and a restart resumes from a checkpoint instead of from "now". Blue marks what changed.
  • Ingest: two small, stateless services hold the Twitter connection, reconnect with backoff, and keep the tweet id. The id is what makes deduplication and idempotent writes possible later.
  • Kafka: three brokers with replication factor 3, so losing a machine loses no data. Six partitions let several stream workers read in parallel.
  • Stream processor: Flink or Spark Structured Streaming, running several workers. This is where the five fixes live: deduplicate by tweet id, extract every track and every artist, enrich through the cache, window by event time with a watermark, and count per window and artist.
  • Checkpoints: processing state and Kafka offsets are saved together to object storage, so a restart resumes exactly where it stopped.
  • Dead-letter topic: failed lookups are kept, counted and replayable instead of dropped.
  • Cassandra: three nodes with replication factor 3, storing one row per (window, artist). Rewriting a count is harmless, so replays can’t corrupt it.
  • Serving: a small query API caches the current top 20 and pushes updates to the dashboard, instead of every open browser running its own GROUP BY every five seconds.
  • Metrics: consumer lag says whether the stream is keeping up; drop counters by reason say how much of the stream the chart actually represents.

The biggest differences from what we built aren’t the extra machines. They’re the things that make results repeatable: saved offsets, idempotent writes, event time and counted drops. Those would have been worth doing even on one laptop.

The code and the report are on GitHub.

References

  1. Apache Kafka. Apache Kafka documentation. kafka.apache.org. https://kafka.apache.org/documentation/ — topics, partitions, replication and consumer groups.
  2. Martin Kleppmann. Designing Data-Intensive Applications. O’Reilly Media, 2017. https://dataintensive.net/ — logs, stream processing and delivery guarantees.
  3. Apache Spark. Spark Streaming Programming Guide. spark.apache.org. https://spark.apache.org/docs/latest/streaming-programming-guide.html
  4. Apache Cassandra. Apache Cassandra documentation. cassandra.apache.org. https://cassandra.apache.org/doc/latest/
  5. Tyler Akidau et al. The Dataflow Model: A Practical Approach to Balancing Correctness, Latency, and Cost in Massive-Scale, Unbounded, Out-of-Order Data Processing. Proceedings of the VLDB Endowment 8(12), 2015. https://www.vldb.org/pvldb/vol8/p1792-Akidau.pdf — event time, windows and watermarks.
  6. Apache Flink. Timely Stream Processing. nightlies.apache.org. https://nightlies.apache.org/flink/flink-docs-stable/docs/concepts/time/
  7. Apache Spark. Structured Streaming Programming Guide. spark.apache.org. https://spark.apache.org/docs/latest/structured-streaming-programming-guide.html — checkpointing, event-time windows and watermarks.