IoT Data Pipelines Explained

Book 34 of 50 — AstolixGen Learning Series For researcher and publication students

Book cover


About This Book

Sensors produce data; decisions need information. The pipeline between them — ingestion, transport, stream processing, storage, and serving — is where most IoT projects succeed or fail. This book builds the full pipeline mental model: MQTT/CoAP ingestion, stream processors, time-series databases, batch vs stream trade-offs, and exactly-once semantics. You will learn to design, measure, and publish pipeline studies with honest end-to-end numbers.

Learning objectives: - Map the five stages of an IoT data pipeline - Choose ingestion protocols and brokers for scale - Compare stream vs batch processing with measured trade-offs - Select time-series storage (InfluxDB, TimescaleDB) by workload - Reason about delivery semantics: at-most/at-least/exactly-once - Handle late, missing, and out-of-order sensor data - Design a publishable pipeline evaluation


Chapter 1: The Pipeline Mental Model

Every IoT data pipeline has five stages. Ingest: devices publish readings (MQTT/CoAP/HTTP) to a broker or gateway. Transport/buffer: a durable log (Kafka, MQTT broker with persistence) absorbs bursts and decouples producers from consumers. Process: stream processors (Flink, Spark Streaming, or simple consumers) clean, aggregate, detect events, and enrich. Store: hot data in time-series DBs for dashboards, cold data in object storage for training. Serve: APIs, dashboards, alerts, and ML training jobs consume the data.

Failures concentrate at the joints: ingest without buffering loses data during outages; processing without idempotency double-counts on replay; storage without retention policy grows forever and bankrupts the project. Draw your pipeline as a diagram before writing any code — then label each arrow with data rate, format, and delivery guarantee. That diagram belongs in every IoT paper.

Example: A 500-sensor farm pipeline: sensors → MQTT (Mosquitto) → Telegraf → InfluxDB (hot, 30-day retention) → nightly export to Parquet on disk (cold) → Grafana dashboards + weekly training job. Each joint labeled: 500 msg/min, JSON, at-least-once.

For your research: The pipeline diagram with rates and guarantees is a required figure in systems papers. Reviewers check that every arrow is justified.

Key takeaway: Ingest → buffer → process → store → serve — label every joint with rate, format, and guarantee.


Chapter 2: Ingestion at Scale — Protocols and Brokers

Ingestion is the narrow waist: everything flows through it. For under ~10k msg/s, a single Mosquitto or EMQX node suffices; beyond that, you need clustered brokers (EMQX, HiveMQ) or a Kafka fronted by MQTT (via Kafka's MQTT proxy or a bridge). Protocol choice at ingest follows Book 31/33: MQTT for pub/sub fan-out, CoAP for UDP-constrained links, HTTP for simplicity.

The key ingestion metrics: throughput (msg/s sustained), fan-in (concurrent publishers), p99 ingest latency, and behavior under backpressure (what happens when consumers lag — drop? buffer? block?). Test backpressure explicitly: kill the consumer for 60 seconds and observe. Systems that silently drop data fail open in dangerous ways.

Example: A campus deployment: 2,000 sensors × 1 msg/10 s = 200 msg/s — trivially handled by one Mosquitto instance at <5% CPU. The same team later added 50 cameras' metadata at 30 msg/s each and hit Mosquitto's single-node ceiling — migrating to EMQX clustered took a weekend because topics were unchanged.

For your research: Ingestion benchmarks (broker X vs Y at N publishers) are clean, reproducible papers — provided you report hardware, message sizes, QoS, and the backpressure policy.

Key takeaway: Size ingestion for 10× your current load; test backpressure; keep topics stable across broker migrations.


Chapter 3: Stream vs Batch — When Each Wins

Stream processing handles each event as it arrives (milliseconds to seconds of latency): alerting, live dashboards, online control. Batch processing handles bounded datasets on a schedule (minutes to hours): daily reports, model retraining, billing. Micro-batch (Spark Streaming's model) is a middle ground with seconds of latency.

Choose by the decision's timescale: a gas-leak alert cannot wait for the nightly batch; a monthly yield report gains nothing from streaming. Many systems do both from the same log (the "Kappa" idea): the stream serves real-time needs, and replays of the log serve batch needs — one ingestion path, two consumers.

The research-relevant trade-off: streaming systems are harder to make correct (out-of-order data, exactly-once) and their evaluation must include late-data handling. Batch is simpler and reproducible. A paper should justify the choice against the application's latency requirement, stated numerically.

Example: Cold-chain monitoring: stream path — temperature > 8°C for 10 min → SMS alert within 60 s. Batch path — nightly: compute per-truck compliance scores, retrain the spoilage-risk model weekly. Same MQTT topics feed both.

For your research: If you propose a streaming algorithm, evaluate it against a batch baseline on the same data and report the latency/accuracy trade-off. "Streaming" alone is not a contribution.

Key takeaway: Match processing model to decision latency; justify streaming with a numeric latency requirement.


Chapter 4: Time-Series Storage — InfluxDB, TimescaleDB, and Tiers

Sensor data is time-series: timestamped, append-only, queried by time range and tag. General databases handle it poorly at scale. InfluxDB (tag-set model, built-in downsampling/retention) excels at high-ingest monitoring. TimescaleDB (PostgreSQL extension, hypertables) excels when you need SQL joins with relational data. Tiering is the cost control: hot tier (SSD, days–weeks) for dashboards, warm tier (compressed, months), cold tier (Parquet on object storage, years) for training and audit.

Design decisions: schema (one measurement per metric vs wide rows — the former scales better), retention policies (drop raw after 30 days, keep 1-min rollups for a year), and downsampling tasks (continuous aggregates). Get retention wrong and storage costs grow unboundedly — the most common IoT project failure after security.

Example: 1,000 sensors × 1 reading/10 s × 50 bytes = 432 MB/day raw. Hot: 14 days raw (6 GB). Warm: 1-min averages for 1 year (105 MB). Cold: Parquet yearly archives. Total manageable; without downsampling, 5-year raw retention would be 788 GB.

For your research: Storage benchmarks (ingest rate, query latency for typical dashboard queries, compression ratio) across InfluxDB/TimescaleDB on identical workloads are useful, citable systems work.

Key takeaway: Use time-series DBs, design retention + downsampling on day one, tier by access pattern.


Chapter 5: Delivery Semantics — The Hardest Concept

At-most-once: messages may be lost, never duplicated (MQTT QoS 0). At-least-once: messages never lost, may be duplicated (QoS 1, Kafka default) — consumers must be idempotent (processing a duplicate changes nothing). Exactly-once: each message processed once (QoS 2 per-hop, Kafka transactions + idempotent producer) — expensive and often unnecessary.

The pragmatic rule: use at-least-once everywhere and make consumers idempotent (keyed writes, deduplication on message IDs). Reserve exactly-once for money and safety. Most IoT analytics (averages, counts over windows) tolerate duplicates if you dedupe by (device_id, timestamp) — a 10-line change that removes the need for expensive exactly-once machinery.

Example: A rainfall counter double-counting due to QoS 1 redelivery overstated irrigation by 8%. Fix: consumer keyed on (sensor_id, reading_ts) with "insert if not exists" — duplicates became harmless, no protocol change needed.

For your research: Semantics papers should demonstrate the failure: show duplicate rates under redelivery, quantify the error without dedup, then show the fix. Measured failure + fix is a complete story.

Key takeaway: At-least-once + idempotent consumers beats exactly-once machinery for most IoT analytics.


Chapter 6: Late, Missing, and Out-of-Order Data

Real sensor data is messy. Late data arrives after its window closed (network delay). Out-of-order data arrives with earlier timestamps after later ones. Missing data never arrives (dead sensor, lost packets). Your pipeline must define policies: watermarks (how long to wait for late data), allowed lateness, and gap-filling (forward-fill, interpolate, or mark missing — never silently fabricate).

For ML, data quality is a first-class concern: a model trained on silently-interpolated gaps learns fiction. Log data-quality flags alongside readings (was_interpolated, gap_seconds) so training pipelines can exclude or weight suspect data.

Example: A vibration monitor's 2% packet loss created gaps that a naive pipeline forward-filled — the anomaly detector then "detected" the flat lines as anomalies. Fix: propagate a data-quality flag; the detector learned to ignore flagged windows. False alarms dropped 60%.

For your research: Data-quality handling is an under-published topic. A paper characterizing gap patterns in a real deployment and evaluating mitigation strategies fills a genuine gap.

Key takeaway: Define late/out-of-order/missing policies explicitly; propagate quality flags to ML.


Chapter 7: Schema, Formats, and Evolution

Choose a wire format deliberately. JSON is human-readable and debuggable but verbose (bad for LoRa). CBOR/MessagePack are binary JSON — 40–60% smaller, still schemaless. Protobuf/Avro need schemas but give the smallest payloads and safe evolution (additive field changes). Rule: JSON for gateways and debugging, binary for constrained links, schema-based when multiple teams consume the data.

Schema evolution is where projects rot: a firmware update adds a field and breaks the dashboard. Use a schema registry (or at minimum versioned topics like v2/...), enforce backward-compatible changes (only add optional fields), and test old consumers against new producers in CI.

Example: Adding battery_mv to a sensor payload: with Protobuf, old consumers ignore the new field — zero downtime. With ad-hoc CSV, the dashboard parsed the wrong column for a week. The paper's lesson section was titled "use schemas."

For your research: Format comparisons (bytes per message, parse time on MCU, energy) across JSON/CBOR/Protobuf on your hardware are compact, useful measurements.

Key takeaway: Binary + schemas for constrained links; version everything; test evolution in CI.


Chapter 8: Observability — Monitoring the Pipeline Itself

Monitor the pipeline, not just the sensors. Golden signals: throughput (msg/s per stage — a drop means blockage), lag (consumer offset behind the log — growing lag means under-provisioning), error rate (parse failures, auth failures), duplication rate (are idempotency keys working?), and end-to-end latency (sensor timestamp → dashboard render).

Alert on pipeline health, not just sensor values: "ingest lag > 5 min" pages the engineer before the dashboard goes stale. Keep a dead-letter queue for unparseable messages — every one is a bug report from the field.

Example: A farm dashboard froze every night at 2 AM. Pipeline metrics showed the batch export locking the hot DB — the fix (read replica for exports) was obvious once lag and query-latency were graphed. The postmortem became a conference talk.

For your research: Observability studies — what to monitor, alert thresholds derived from data — are practical contributions that industry tracks value.

Key takeaway: Monitor throughput, lag, errors, duplication, and end-to-end latency; alert on pipeline health.


Chapter 9: Edge Processing in the Pipeline

Push processing to the edge where it pays: filtering (send only interesting events), aggregation (1-min averages instead of 1-s raw), compression (delta encoding), and event detection (thresholds, tiny models). The pipeline then handles orders of magnitude less data, and the cloud bill collapses.

The design question is what the edge may discard: raw data needed for future model training must still be sampled (store 1% raw or event-centered windows). "Filter everything" destroys your future training set — a common regret.

Example: Acoustic pest trap: edge classifier scores every 10 s of audio; only scores > 0.7 upload a 5-s clip; 1% random raw samples upload for retraining. Uplink: 2.4 GB/day → 31 MB/day. Retraining data preserved via sampling.

For your research: Edge-filtering papers must report the preservation metric: what fraction of true events survived filtering (recall at the filter stage), not just the data reduction.

Key takeaway: Filter and aggregate at the edge, but sample raw data for future training; report filter recall.


Chapter 10: Security Across the Pipeline

Each stage is an attack surface. Ingest: TLS, per-device credentials, topic ACLs (Book 31). Transport: encrypt the log, authenticate consumers. Processing: validate and sanitize inputs (a malicious sensor can inject payloads that break parsers — fuzz your parsers). Storage: encrypt at rest, least-privilege DB roles, audit logs. Serve: authenticated APIs, rate limiting.

The pipeline-specific risk is data integrity: an attacker who can inject plausible sensor readings can manipulate decisions (fake "all clear" from a safety sensor). Sign readings at the device (HMAC with a device key) and verify at ingestion — integrity, not just confidentiality.

Example: A building-management pipeline accepted unsigned JSON; a compromised thermostat injected "CO2 normal" during a real ventilation failure. Post-incident: HMAC-signed CBOR at the device, verification at the broker bridge, alerts on signature failures.

For your research: Pipeline security analyses — threat model per stage, measured cost of HMAC/TLS on the device — are strong applied-security contributions.

Key takeaway: Secure every stage; sign sensor data at the source for integrity, not just TLS for confidentiality.


Chapter 11: Cost Engineering the Pipeline

Pipeline costs hide in: ingest (per-message broker/cloud fees), storage (hot tier $/GB-month), processing (always-on stream workers), and egress (dashboards pulling data out of the cloud). Model each: monthly = ingest_msgs × fee + hot_GB × rate + worker_hours × rate + egress_GB × rate. Then attack the biggest term: downsample sooner, filter at edge, cache dashboards, move cold data to cheap object storage.

Example: A 5,000-sensor pipeline cost $1,900/month: 55% hot storage, 25% stream workers, 20% ingest. Cutting hot retention 30→7 days and adding edge aggregation saved $1,100/month with no dashboard change — the cost teardown was the paper's most-read section.

For your research: Cost teardowns of real pipelines are rare and valuable. Anonymize the deployment, show the model, show the optimization — industry citations follow.

Key takeaway: Model cost per stage, attack the biggest term, publish the teardown.


Chapter 12: Your Pipeline Study — Design to Publication

Template:

  1. Problem: A pipeline pain (e.g., "dashboard staleness during cellular outages on a 2,000-sensor farm").
  2. Hypothesis: Falsifiable, e.g., "edge buffering + store-and-forward cuts data loss from 8% to <0.5% during 4-hour outages at 1.3× device energy."
  3. Design: Five-stage diagram with rates, formats, guarantees; components named with versions.
  4. Method: Fault injection (kill links, delay, duplicate), metrics (loss, duplication, lag, end-to-end latency, energy), runs with variance.
  5. Baselines: The naive pipeline (direct MQTT→DB, no buffer) and/or a cloud-only alternative.
  6. Results: Tables, lag graphs, cost model.
  7. Limitations & threats: Emulation fidelity, hardware specificity.
  8. Artifact: Docker Compose of the whole pipeline, configs, workload generator, raw logs.

Write in IEEE format; the architecture diagram is Figure 1; evaluation leads with the end-to-end numbers.

For your research: This is your A1→A4 pipeline again: problem + lit review, method, results, presentation — now for data infrastructure.

Key takeaway: Pain → hypothesis → five-stage design → fault-injected evaluation → artifact = a publishable pipeline paper.


Learning Dashboard

# Chapter Core idea Research use
1 Pipeline model 5 stages, labeled joints Architecture figure
2 Ingestion Brokers, backpressure Broker benchmarks
3 Stream vs batch Decision-latency rule Justify with numbers
4 TS storage InfluxDB/TimescaleDB, tiers Storage benchmarks
5 Semantics At-least-once + idempotency Failure demos + fixes
6 Messy data Late/missing/out-of-order Gap-pattern studies
7 Formats JSON→CBOR→Protobuf Format measurements
8 Observability Lag, duplication, e2e latency Monitoring studies
9 Edge processing Filter + sample raw Filter-recall reporting
10 Security Per-stage threats, signed data Threat models + costs
11 Cost Per-stage model, teardown Published cost analyses
12 Study template Pain→hypothesis→artifact A1–A4 mapping

References

[1] L. Atzori, A. Iera, and G. Morabito, "The Internet of Things: A survey," Computer Networks, vol. 54, no. 15, pp. 2787–2805, 2010. [2] A. Banks and R. Gupta, "MQTT Version 3.1.1," OASIS Standard, Oct. 2014. [3] W. Shi, J. Cao, Q. Zhang, Y. Li, and L. Xu, "Edge Computing: Vision and Challenges," IEEE Internet of Things Journal, vol. 3, no. 5, pp. 637–646, 2016. [4] M. Satyanarayanan, "The Emergence of Edge Computing," Computer, vol. 50, no. 1, pp. 30–39, 2017. [5] J. Kreps, N. Narkhede, and J. Rao, "Kafka: A distributed messaging system for log processing," in Proc. NetDB Workshop, 2011. [6] T. Akidau et al., "The Dataflow Model: A practical approach to balancing correctness, latency, and cost in massive-scale, unbounded, out-of-order data processing," Proc. VLDB Endowment, vol. 8, no. 12, 2015. [7] S. Wolfert, L. Ge, C. Verdouw, and M. J. Bogaardt, "Big Data in Smart Farming — A review," Agricultural Systems, vol. 153, pp. 69–80, 2017. [8] R. Roman, P. Najera, and J. Lopez, "Securing the Internet of Things," Computer, vol. 44, no. 9, pp. 51–58, 2011. [9] ISO/IEC 20922:2016, "Information technology — MQTT v3.1.1," 2016. [10] H. Kopetz, Real-Time Systems: Design Principles for Distributed Embedded Applications, 2nd ed. Springer, 2011. (book)


Glossary

  • Ingestion — entry of device data into the pipeline
  • Backpressure — what happens when consumers lag producers
  • Stream processing — per-event computation with low latency
  • Batch processing — scheduled computation over bounded data
  • Time-series DB — database optimized for timestamped data
  • Retention policy — how long raw/aggregated data is kept
  • Downsampling — reducing resolution of old data (1 s → 1 min)
  • Idempotent — safe to apply more than once
  • Watermark — stream-processing bound on how late data may be
  • Dead-letter queue — holding area for unprocessable messages
  • Schema registry — versioned store of data formats
  • Egress — data leaving the cloud (often billed)

Practice Exercises

  1. Draw the five-stage pipeline for a 200-sensor deployment you imagine. Label each arrow with rate, format, guarantee.
  2. Your consumer lags during a 1-hour outage. Compare three backpressure policies and their consequences.
  3. When is streaming unjustified? Give a numeric latency argument.
  4. Size hot/warm/cold tiers for 5,000 sensors at 1 reading/30 s, 60 bytes. Show your math.
  5. Explain why at-least-once + idempotency beats exactly-once for dashboard aggregates.
  6. Design watermark and allowed-lateness policies for a sensor network with 5% late data.
  7. Compare JSON, CBOR, and Protobuf for a LoRa payload: bytes, parse cost, evolution safety.
  8. List five pipeline golden signals and what each failure mode looks like.
  9. Your edge filter drops 99% of data. What metric proves you didn't lose real events?
  10. Write a cost model for a pipeline you design; identify the biggest term and one optimization.

End of Book 34. Next: Book 35 — Smart Agriculture with IoT.