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.
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.
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.
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:
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
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 artistover 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.
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:
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:
- 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 BYevery 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
- Apache Kafka. Apache Kafka documentation. kafka.apache.org. https://kafka.apache.org/documentation/ — topics, partitions, replication and consumer groups.
- Martin Kleppmann. Designing Data-Intensive Applications. O’Reilly Media, 2017. https://dataintensive.net/ — logs, stream processing and delivery guarantees.
- Apache Spark. Spark Streaming Programming Guide. spark.apache.org. https://spark.apache.org/docs/latest/streaming-programming-guide.html
- Apache Cassandra. Apache Cassandra documentation. cassandra.apache.org. https://cassandra.apache.org/doc/latest/
- 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.
- Apache Flink. Timely Stream Processing. nightlies.apache.org. https://nightlies.apache.org/flink/flink-docs-stable/docs/concepts/time/
- 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.