TL;DR. Part 1 covered topics, partitions, offsets and consumer groups. This part is about running Kafka in production: scaling consumer groups horizontally, transforming and aggregating data in flight with Kafka Streams, shipping it to Elasticsearch, Neo4j and InfluxDB with Kafka Connect instead of hand-written consumers, and choosing between Lambda and Kappa architectures. The running example is a real-time social network analytics platform from a university data engineering course.
Picture a startup building a niche content platform: people share posts about places, recipes and local events, follow each other, like and comment. The product team has three questions and they want answers now, not in tomorrow’s batch report. Which hashtags are going viral? Who are the most influential users? Which people should we suggest in the “you might know” panel?
This is project P4 — Social Network Analytics Platform — from the Advanced Data Modeling course I teach at Politecnico di Torino. It combines five technologies: MongoDB for profiles and posts, Neo4j for the social graph, Elasticsearch for full-text search and trends, Grafana for KPIs, and Kafka as the nervous system tying it all together. In Part 1 we used Kafka to move sensor readings around. Here we use it for something more ambitious: letting several systems react to the same event at the same moment, without any of them knowing the others exist.

From request-response to event-driven
In a typical web application, a “like” works like this: the backend receives the request, writes to the database, then calls the notification service, updates a counter, perhaps pushes something to analytics. Every call is a potential failure point. If notifications are slow, the user waits. If analytics is down, the like may be lost or surface as an error.
In an event-driven design the backend does exactly one thing: it publishes an event to Kafka — {"event_type": "like", "user_id": "u42", "target_id": "p1337", "timestamp": "..."} — and returns immediately. Everything else happens downstream, asynchronously. Any service that cares about likes subscribes to the topic and does its job at its own pace.
Events as immutable facts
The shift is subtle but it matters. An event is not a command (“update the counter”); it is a fact (“user u42 liked post p1337 at 6:04 pm”). Facts are never edited or deleted — they accumulate. That is why Kafka’s durable log is the natural home for them, and why any consumer can rebuild its own state by replaying history from the beginning.
Three concrete benefits
- Decoupling. The service producing likes knows nothing about Neo4j or Elasticsearch. You can add a fifth consumer tomorrow without touching the producer.
- Replay. Got the trending logic wrong? Fix the code, rewind the consumer’s offset by seven days, and recompute everything.
- Fan-out. The same event feeds the social graph, the search index and the dashboard. One source of truth, many specialized views — the essence of polyglot persistence, which gets its own deep dive later in this series.
Consumer groups and horizontal scaling in production
Part 1 showed how a consumer group splits a topic’s partitions among its members. In production this is your main scaling lever. If social-events has 6 partitions and the Neo4j writer can’t keep up, start a second and third instance with the same group.id: Kafka redistributes the partitions and the load is shared. Beyond 6 instances you gain nothing — extra consumers sit idle. Partition count is the ceiling on parallelism, which is why you should size it thoughtfully when you create the topic.
Rebalancing and the new group protocol
Whenever a consumer joins or leaves, Kafka runs a rebalance and reassigns partitions. With the classic protocol the whole group paused for a few seconds — painful with dozens of consumers. Kafka 4.0 made the broker-side consumer group protocol (KIP-848) generally available, turning rebalances incremental: only partitions that actually need a new owner move. Kafka 4.2 (February 2026) brought the same approach to Kafka Streams applications (KIP-1071), GA with a limited feature set. If you’re starting a new cluster today, turn it on.
Delivery semantics: how many times does a message arrive?
Kafka gives you three guarantees, and which one you get depends on when you commit the offset:
- At-most-once: commit before processing. Crash mid-way and the message is gone. Fine for approximate metrics.
- At-least-once: commit after processing. Crash after writing but before committing and you’ll process it again. The most common default.
- Exactly-once: achievable with idempotent producers and Kafka transactions, as long as the destination is also Kafka (which is what Kafka Streams does).
In practice, the most robust pattern is at-least-once plus idempotent writes at the destination: if the same event arrives twice, nothing changes. In Neo4j that means MERGE instead of CREATE:
from kafka import KafkaConsumer
from neo4j import GraphDatabase
import json
consumer = KafkaConsumer(
"social-events",
bootstrap_servers="localhost:9092",
group_id="neo4j-writer",
enable_auto_commit=False,
value_deserializer=lambda v: json.loads(v.decode("utf-8")),
)
driver = GraphDatabase.driver("neo4j+s://xxxx.databases.neo4j.io", auth=("neo4j", "***"))
CYPHER = {
"follow": "MERGE (a:User {id:$uid}) MERGE (b:User {id:$tid}) MERGE (a)-[:FOLLOWS]->(b)",
"like": "MERGE (a:User {id:$uid}) MERGE (p:Post {id:$tid}) MERGE (a)-[:LIKES]->(p)",
}
for msg in consumer:
ev = msg.value
query = CYPHER.get(ev["event_type"])
if query:
with driver.session() as s:
s.run(query, uid=ev["user_id"], tid=ev["target_id"])
consumer.commit() # commit only after the write: at-least-onceBecause of MERGE, replaying the topic never creates duplicate relationships. The graph stays consistent even if you reprocess a full week of events.
Kafka Streams: transforming data in flight
So far our consumers just move data from Kafka into a database. Often, though, you want to process the stream before it lands: drop test events, enrich a like with the post’s category, count hashtags over the last ten minutes. Kafka Streams is the Java library (part of Apache Kafka itself) built for exactly this. It is not a separate cluster to operate like Spark or Flink — it’s a plain application that reads from topics, processes, and writes to other topics. To scale it, you run more instances; the mechanics are the same as consumer groups.
Stateless operations
These look at one event at a time with no memory: filter (drop bot traffic), map/mapValues (normalize timestamps), flatMap (turn one post with five hashtags into five events), split/branch (route likes, follows and comments to separate topics). They’re fast, easy to scale and need no local storage.
Stateful operations: aggregations, windows and joins
Things get interesting when you need memory. Counting hashtags means remembering previous counts; joining a like with post metadata means having that metadata at hand. Kafka Streams keeps this state in a local state store (RocksDB by default) and logs every change to an internal changelog topic. If an instance dies, another one rebuilds the state from the changelog and carries on.
Here’s a topology that computes trending hashtags over 10-minute windows:
StreamsBuilder builder = new StreamsBuilder();
builder.stream("posts", Consumed.with(Serdes.String(), postSerde))
.flatMapValues(post -> post.getHashtags()) // 1 post -> N hashtags
.selectKey((k, tag) -> tag.toLowerCase())
.groupByKey(Grouped.with(Serdes.String(), Serdes.String()))
.windowedBy(TimeWindows.ofSizeWithNoGrace(Duration.ofMinutes(10)))
.count(Materialized.as("hashtag-counts")) // state store
.toStream()
.map((wKey, count) -> KeyValue.pair(wKey.key(), count))
.to("trending-hashtags", Produced.with(Serdes.String(), Serdes.Long()));The output is a continuously updated trending-hashtags topic, ready to index in Elasticsearch or chart in Grafana.
KStream vs KTable: two ways to read the same log
Kafka Streams introduces a distinction worth getting right. A KStream treats the topic as a sequence of independent events: three likes from the same user are three separate facts. A KTable treats it as a changelog: for each key only the latest value matters, like a table row being updated. A user profile is a KTable (you want the current version); their likes are a KStream. Joining the two — “enrich every like with the user’s country” — is one of the most common production patterns.
Not a Java shop?
Kafka Streams runs on the JVM. If your team lives in Python, solid alternatives exist: Quix Streams and faust-streaming offer similar APIs with windows and state, and Flink SQL (also available on Confluent Cloud) lets you express the same aggregation in plain SQL. For a 20-hour student project, a Python consumer aggregating short windows in memory is often enough; in production, pick a tool that handles state and fault tolerance for you.
Kafka Connect: integrations without writing consumers
The Neo4j consumer above works, but multiply it by every source and every sink and you end up with dozens of small scripts to maintain, monitor and patch. Kafka Connect replaces them with a standard framework: instead of writing code, you configure a connector in JSON. Connect handles offsets, retries, parallelism and error handling.

Source and sink connectors
There are two kinds. Source connectors read from an external system and write to Kafka. The most important is Debezium, which tails the transaction log of PostgreSQL, MySQL or MongoDB and publishes every INSERT, UPDATE and DELETE as an event. That technique — Change Data Capture (CDC) — is the de facto standard for keeping different databases in sync without dual writes in application code. Sink connectors go the other way: they read a topic and write to a destination.
Here’s an Elasticsearch sink that indexes P4’s posts:
{
"name": "posts-to-elasticsearch",
"config": {
"connector.class": "io.confluent.connect.elasticsearch.ElasticsearchSinkConnector",
"topics": "posts",
"connection.url": "http://elasticsearch:9200",
"key.ignore": "false",
"schema.ignore": "true",
"tasks.max": "3",
"errors.tolerance": "all",
"errors.deadletterqueue.topic.name": "dlq-posts",
"errors.deadletterqueue.context.headers.enable": "true"
}
}Look at the last three lines. With errors.tolerance and a dead letter queue, a malformed document doesn’t stall the whole pipeline — it’s parked in a separate topic for later inspection. Details like this are what separate a prototype from a production system.
Neo4j and InfluxDB as destinations
For the graph there’s the Neo4j Connector for Kafka, which lets you attach a Cypher statement to each topic — essentially our Python MERGE logic, declared as configuration. For time series you have two options: an InfluxDB sink connector, or Telegraf, InfluxData’s agent, which ships a native kafka_consumer input plugin. For sensor pipelines it’s usually the simplest way to get batched writes into InfluxDB.
Case study: real-time social network analytics
Let’s put it together. The project’s simulator emits an event every 0.5–2 seconds with a uniform JSON payload. Each event type has one or more destinations, each chosen for the kind of question it has to answer:
| Event | Topic | Consumed by | Question it answers |
|---|---|---|---|
| New post | posts | MongoDB + Elasticsearch (Connect) | “Show me posts about wine bars near downtown” |
| Follow | social-events | Neo4j (idempotent consumer) | “Who should I follow?” (2-hop friend-of-friend) |
| Like / comment | social-events | Neo4j + Kafka Streams | “Which posts are growing fastest in their first 2 hours?” |
| Aggregated hashtags | trending-hashtags | Elasticsearch / Grafana | “What’s trending in the last hour?” |
| Engagement KPIs | engagement-metrics | InfluxDB → Grafana | “Posts per day, like ratio, follower growth” |
One detail that surprises people: Grafana has no native Kafka data source. To chart KPIs in real time you first write them to something Grafana can query — InfluxDB or PostgreSQL. That’s not a limitation; it’s consistent with Kafka’s philosophy: the log transports and retains, specialized databases serve queries. For a refresher on building panels, see our Grafana guide for technical and IoT data; graph modeling is covered in the Neo4j article.
The query that makes the graph pay off
Once follows land in Neo4j in real time, “people you might know” is a few lines of Cypher:
MATCH (me:User {id: $uid})-[:FOLLOWS]->(friend)-[:FOLLOWS]->(suggestion)
WHERE NOT (me)-[:FOLLOWS]->(suggestion) AND suggestion <> me
RETURN suggestion.id, count(friend) AS mutual_friends
ORDER BY mutual_friends DESC
LIMIT 10Kafka’s contribution isn’t the query itself — it’s that the result reflects follows from a few seconds ago, not last night’s snapshot.
Lambda vs Kappa: where Kafka fits
Once data arrives as a stream, you eventually ask: how do I also run analytics over the full history? The two classic answers have names.

Lambda Architecture
Proposed by Nathan Marz in 2011, Lambda has three layers: a batch layer that periodically recomputes the entire history (the ground truth), a speed layer that processes recent events in real time to cover the batch latency, and a serving layer that merges both. The upside is reliability — batch always corrects streaming mistakes. The downside is that you write and maintain the same logic twice, in two different technologies.
Kappa Architecture
In 2014 Jay Kreps — one of Kafka’s creators at LinkedIn — proposed dropping the batch layer. Everything flows through streaming; to recompute history, you replay the log from the start with a new version of the code. One codebase, simpler operations. The cost: you need mature streaming infrastructure and long enough retention in Kafka.
| Criterion | Lambda | Kappa |
|---|---|---|
| Historical reprocessing | Natural (batch jobs) | Needs long Kafka retention |
| Code complexity | Two pipelines | One pipeline |
| Correctness guarantee | Batch is ground truth | Reprocess to correct |
| Team skills | Batch + streaming | Streaming only |
| Choose it when | Strict audit/compliance | High velocity, agile teams |
My advice to students is pragmatic: start with Kappa if your team is comfortable with streaming, and move to Lambda only if a regulatory constraint demands certified batch recomputation. For P4, Kappa is the obvious fit: everything starts life as an event and the Kafka log is already the source of truth.
Common mistakes (and how to avoid them)
After mentoring dozens of project teams, these are the pitfalls I see most:
- Too few partitions. A single-partition topic can’t scale past one consumer. Fine for a lab, not for production. You can add partitions later but never remove them.
- CREATE instead of MERGE. The first replay fills the graph with duplicates. Make every write idempotent.
- Wrong message key. Ordering is only guaranteed within a partition. If one user’s events must stay in order, key by
user_id. - No dead letter queue. One malformed JSON can block a connector for hours.
- Ignoring consumer lag. If consumers fall behind, your “real-time” data is minutes old. Monitor lag and alert on it — it’s the single most important Kafka metric.
Where to start
To experiment without installing anything, use Confluent Cloud’s free tier as we did in Part 1. If you prefer local, a docker-compose file with a single-node Kafka image in KRaft mode (no ZooKeeper — it was removed for good in Kafka 4.0) gives you a full environment in minutes. Suggested path: a Python producer and consumer first, then an Elasticsearch sink connector, then a small Kafka Streams topology with a time window.
That wraps up our Kafka chapter. Next in the Digital Stack 2026 series we go one level down and look at where sensor readings actually live: InfluxDB, the database built for time series. If you’re weighing an event-driven architecture for your organization and want to know whether you really need Kafka or something simpler will do, request a data and AI analysis — we start from your actual data flows, not from whatever’s trending.
