Hands-on: building a pipeline to collect and store IoT sensor data in a client's factory
Factory sensors send data every second. The hard part is making a few decisions before the first packet crosses the network: what the data looks like, which protocol carries it and which key it is stored under.
In brief
- Kafka's documentation names factory sensor data collection as a use case. According to Kai Waehner, Kafka complements MQTT and OPC UA rather than competing with them.
- Cassandra and Bigtable both handle large volumes. What decides performance is the row key, because Bigtable sorts every row by it.
- IoT devices are easy to attack because they are always connected, so access to the factory network must be agreed before any work begins.
- 1Sensor on the machineDetects changes in its environment; used to monitor machine performance
- 2MQTT or OPC UA at the edgeKeep the protocol the client already uses; MQTT suits unreliable networks
- 3Normalise the eventMap to one schema: plant, line, sensor, metric, value, time
- 4Publish to KafkaWrite to a topic; consumers subscribe to read the event stream
- 5Write to Cassandra/BigtableBuild the row key around the questions the client will query
Every block involves its own decisions, but the row key in the last block determines whether the client's queries are fast or slow.
Graphic: FDE Times
Picture a production line with 200 sensors, each sending one reading per second. Multiply by the 86,400 seconds in a day and you have 17,280,000 records. At that volume, a poor choice of storage key is no longer a minor mistake, because it is repeated tens of millions of times a day.
Apache Kafka’s introductory documentation names this use case directly: continuously capturing and analysing sensor data from IoT devices, for example in factories or wind farms. If you want to work as an FDE for manufacturing clients, sooner or later you will have to build a pipeline like this inside their plant.
What follows sketches that pipeline in five steps, from sensor to edge protocol to Kafka to storage. The aim is not to install each component but to settle the right design decisions before you write the first line of code on the client’s site.
What will you build, and what do you need?
The end result is a five-block diagram, an event schema, pseudocode for writing to Kafka and reading back out, and a row key design with a clear rationale. All code in this article is simplified pseudocode, not tied to any library’s API.
You should understand publish/subscribe, be able to read JSON and have a basic grasp of key/value NoSQL. Keep a sheet of scrap paper handy: the load calculation is the part newcomers most often skip.
Step 1: Know which sensors the plant has before writing code
IBM defines a sensor as a device that detects changes in its environment. In manufacturing, industrial IoT is used to monitor machine performance. So the first question on site is not “which version of Kafka?” but “which machines need watching, and which sensors measure them?”
Draw up a sensor inventory with these columns: id, the machine it is attached to, the quantity measured, unit, send frequency and protocol. This is the first thing you sit down to review with the client’s operations engineers.
Check: every row in the table must have a send frequency. Without that column you cannot calculate the load in step 5.
Step 2: Settle the shape of an event
Kafka defines event streaming as capturing data in real time from sources such as databases and sensors. Every reading becomes an event, so normalise events from the start:
{
"plant": "plant-a",
"line": "line-3",
"sensor_id": "temp-042",
"metric": "temperature_c",
"value": 71.4,
"ts": "2026-10-09T08:00:01Z"
}
The unit is written into the metric name, and time uses one consistent format. This is a design choice made for this article, not a mandatory standard. It does, however, prevent a classic error: two production lines sending temperature in different units without anyone noticing.
Check: looking at any single event, you can tell immediately which plant, line and sensor it came from and what time it was measured.
Step 3: Let MQTT and OPC UA do their work at the edge
Kai Waehner, who writes extensively on data streaming, argues that MQTT is well suited to unreliable networks with limited bandwidth. In his view, OPC UA suits industrial automation, and Kafka complements rather than competes with both protocols.
In practice, you keep whatever protocol the client’s equipment already uses and add a bridge layer from MQTT or OPC UA into Kafka. Respecting existing infrastructure wins the operations team’s trust faster than any architecture slide deck.
Check: the protocol column in the step 1 table now has a value for every sensor.
Step 4: Publish to Kafka, subscribe to write to storage
Kafka lets you publish (write) and subscribe to (read) streams of events. The writing side looks like this in pseudocode:
# Pseudocode, simplified
for message in mqtt_or_opcua_source:
event = normalize(message) # matches the schema from step 2
publish(topic="sensor-readings", value=event)
The reading side subscribes to the same topic and writes to storage:
# Pseudocode, simplified
subscribe(topic="sensor-readings")
for event in stream:
row_key = make_row_key(event) # designed in step 5
write(table="readings", key=row_key, value=event)
The normalize (normalise) function is where units, time zones and missing fields are handled. Put all cleaning logic there rather than scattering it across individual consumers.
Step 5: The row key is the most important decision
Bigtable is a distributed NoSQL database. Each table is a sorted key/value map, and each row is indexed by exactly one row key. Bigtable is used for both time-series and IoT data. Because rows are stored in key order, the way you compose the row key determines which queries are fast.
Key starting with time
- 2026-10-09T08:00:01#temp-042
- Every new write lands at the end of the key range
- Suits questions about all sensors at one moment
Key starting with device
- nha-may-a#line-3#temp-042#2026-10-09T08:00:01
- One sensor's history sits in a single contiguous range
- Suits questions like this machine's temperature over the past 24 hours
The reason lies in the table being sorted by key. If the key starts with time, every new record has a key larger than all previous keys, so every write piles up at the same end of the table.
The consequence follows directly: in a distributed system, that end of the table lives on one part of the cluster, so a single node carries almost all 17,280,000 writes a day while the others sit nearly idle. This is a write hotspot. Adding machines does not fix it, because new writes still converge on the same place.
When the key starts with the device instead, writes from 200 sensors spread across 200 different key ranges. Each sensor’s history is also stored contiguously, and the question “how did machine X do last week?” requires reading exactly one key range.
The row key analysis above rests on Bigtable’s sorted-key model. If the client uses Cassandra, the project’s documentation states that read and write throughput increase linearly as machines are added. Its masterless architecture lets Cassandra survive the loss of an entire data centre without losing data, and failed nodes can be replaced with no downtime.
Those properties take care of scale, but the key is still yours to design, following the same logic. One option (simplified) is a partition key made of sensor_id and date, for example temp-042 + 2026-10-09, with rows inside each partition ordered by ts.
Each partition then holds exactly 86,400 readings from one sensor for one day, writes are spread evenly by sensor, and the question “machine 3 yesterday” touches a single partition.
Check: write down the three queries the client cares about most, then show how your row key serves each one.
Mistakes that break the pipeline in the first week
The most common mistake is choosing the key by arrival order. Go back to the hypothetical plant at the start: one week produces 7 × 17,280,000 = 120,960,000 records, of which the sensor on machine 3 accounts for 604,800. The plant manager asks: “How has machine 3 done this week?”
With a time-first key, those 604,800 records are interleaved with data from the other 199 sensors. The system has to scan the full range of 120,960,000 rows to filter out what it needs. With a device-first key, it is a single contiguous range.
The next mistake is skipping the load calculation. Multiply the number of sensors by the send frequency and the seconds in a day before sizing the cluster, rather than waiting until the system is already slow.
Nor should security be left for later. Because they are always connected, IoT devices are vulnerable to intrusion and raise privacy concerns. Every time you pull data out of the plant network, you open another way in.
So before your first day on site, ask the client’s IT/OT team: which networks you may access, in which direction data is allowed to flow, and who approves changes. Your architecture must follow those answers, not the other way round.
How this skill shows up on client sites and in your CV
Clients ask why machine 3 keeps stopping and want to see its history. A good FDE turns that question into a row key, a topic and a sensor inventory, then takes those back to the operations engineers for discussion.
When reading job descriptions, look for the keywords IIoT, OPC UA, MQTT, time-series and streaming when they appear alongside “customer-facing”. On your CV, do not just write “used Kafka”. Describe the problem: how many sensors, how often they sent data, which row key you chose and why.
If you have no real project yet, build one yourself with simulated sensor data, including the load calculation and the reasoning behind your key design, and bring it to the interview. By the second week in the plant, the client needs someone who understands the machine producing the data, not another person who only knows how to install software.
Was this article useful?
Thanks for the feedback!
6 sources
- Introduction | Apache Kafka · 2026-05-22
- Apache Cassandra | Apache Cassandra Documentation
- What is the Internet of Things (IoT)?
- Internet of things - Wikipedia
- What is Google Bigtable? | Definition from TechTarget · 2024-01-25
- OPC UA, MQTT, and Apache Kafka – The Trinity of Data Streaming in IoT · 2022-02-11