Skip to content

Examples/IoT

Device telemetry platform (example)

Sensors publish readings over MQTT into a Kafka log. A batch writer loads a time-series store, a rules engine pages on alerts, and Postgres holds devices, credentials and firmware. Replace any part with how your own system works.

Scale: 20k readings/s sustained, 13 months retention

123456789101112CLIENTSensorsEDGEMQTT brokerSERVICEIngestLOG / STREAMKafkaWORKERBatch writerDATABASETime-seriesstoreWORKERRules engineDATABASEPostgresEXTERNALPaging and SMSSERVICEQuery APICLIENTDashboardOBJECT STOREObject storageWORKERCompactor

Select a component to see what it is responsible for and which state it owns.

  1. 1Sensors → MQTT broker: Publish reading
  2. 2MQTT broker → Ingest: Forward messages
  3. 3Ingest → Kafka: Append to log
  4. 4Batch writer → Kafka: Consume readings
  5. 5Batch writer → Time-series store: Batch insert
  6. 6Rules engine → Kafka: Consume for rules
  7. 7Rules engine → Paging and SMS: Send alert
  8. 8MQTT broker → Postgres: Authenticate device
  9. 9Dashboard → Query API: Load charts
  10. 10Query API → Time-series store: Query series
  11. 11Compactor → Time-series store: Roll up old data
  12. 12Compactor → Object storage: Write Parquet
  • Asynchronous
  • Request / response
  • Bulk data

Select a part to trace its flows.

  • Asynchronous
  • Request / response
  • Bulk data

Parts 13

Flows 12

  1. 1Sensors → MQTT brokerPublish readingAsynchronous
  2. 2MQTT broker → IngestForward messagesAsynchronous
  3. 3Ingest → KafkaAppend to logAsynchronous
  4. 4Batch writer → KafkaConsume readingsRequest / response
  5. 5Batch writer → Time-series storeBatch insertBulk data
  6. 6Rules engine → KafkaConsume for rulesRequest / response
  7. 7Rules engine → Paging and SMSSend alertRequest / response
  8. 8MQTT broker → PostgresAuthenticate deviceRequest / response
  9. 9Dashboard → Query APILoad chartsRequest / response
  10. 10Query API → Time-series storeQuery seriesRequest / response
  11. 11Compactor → Time-series storeRoll up old dataRequest / response
  12. 12Compactor → Object storageWrite ParquetBulk data

Invariants 3

  • A device's readings are processed in order

    Kafka key = device id, so one partition per device; readings carry a device sequence number.

    Kept by

  • Redelivered readings are stored once

    Unique (device_id, ts, seq) with ON CONFLICT DO NOTHING; offsets committed after the insert.

    Kept by

  • An incident pages once, not every window

    Per (rule, device) state row in the registry; the ok → firing transition is a conditional UPDATE WHERE state = 'ok', so only one window pages. Resolving waits for hysteresis (several clean windows).

    Kept by