RealtimeTrigger icon
Log icon
Get icon
Script icon
Set icon
If icon
Publish icon
SlackIncomingWebhook icon

IoT Edge Digital Twin with Drift Reconciliation

Realtime MQTT telemetry updates a Redis-backed digital twin, detects desired-vs-reported drift, and publishes reconcile commands only when state changes.

Categories
Infrastructure

How it works

  1. telemetry_stream (io.kestra.plugin.mqtt.RealtimeTrigger, shipped disabled) holds a live subscription on devices/+/telemetry and starts one execution per message, parsing JSON payloads so fields arrive as {{ trigger.payload.* }}.
  2. read_twin (io.kestra.plugin.redis.string.Get, failedOnMissing: false) loads the current shadow document from twin:<device_id>; a missing key simply means a new device.
  3. compute_delta (scripts.python.Script, Process runner) merges the message into the shadow: strips device_id into identity, diffs the reported map against the previous baseline to decide changed, diffs reported against the desired_state input field-by-field (sorted, so drift lists are deterministic), and builds both the persisted twin (last_seen, drift_fields, in_sync) and a ready-to-send reconcile command document. Everything returns through Kestra's ::json:: outputs protocol.
  4. write_twin (io.kestra.plugin.redis.string.Set) always persists the refreshed shadow so the next message has a baseline, even when nothing downstream happens.
  5. change_gate (core.flow.If) fires only on changed or is_new - identical telemetry updates last_seen silently instead of spamming the command topic. Inside it, drift_gate branches on reconcile_planned (present report AND new change AND missing desired fields): drift publishes io.kestra.plugin.mqtt.Publish to devices/<id>/commands with QoS 1 (flat JSON command shape from mqtt-command-dispatch) plus a Slack detail message; in-sync changes just log.
  6. Outputs expose device_id, the full twin, drift, and whether a reconcile was published; errors posts failures to Slack with the device id.

What you get

  • A queryable last-known state per device in Redis, refreshed by every telemetry message.
  • Field-level drift detection between desired configuration and what devices actually report.
  • Idempotent reconciliation: steady-state drift publishes once, then stays quiet until state changes again.
  • Per-message execution history - every twin update and reconcile command is auditable.

Who it's for

IoT platform teams maintaining shadow/ghost state for fleets, ops teams who want configuration drift caught automatically, and anyone replacing a bespoke MQTT-to-Redis bridge daemon.

Why orchestrate this with Kestra

A consumer loop that writes Redis and fires commands works until it silently dies or double-publishes under load. Kestra supervises the subscription, gives every message an execution with logs and retries, and turns the twin update plus reconciliation into one reviewable graph where adding a Slack digest, a ticket, or an anomaly check is a task away - not another daemon to babysit.

Prerequisites

  • An MQTT broker with devices publishing JSON telemetry to devices/<device_id>/telemetry, e.g. {"device_id": "dev-001", "power": "off", "mode": "manual", "report_interval_s": 60, "temp_c": 41.2}.
  • A Redis instance for shadow documents.
  • The desired_state input adjusted to the fields your devices report.

Secrets

  • MQTT_SERVER: broker URI, e.g. tcp://broker.example.com:1883.
  • REDIS_URL: Redis connection URL, e.g. redis://redis:6379/0.
  • SLACK_WEBHOOK_URL: Slack incoming webhook for drift alerts.

Quick start

  1. Add the three secrets to your Kestra namespace.
  2. Set disabled: false on telemetry_stream and save the flow.
  3. Publish telemetry that violates desired_state, e.g. mosquitto_pub -t devices/dev-001/telemetry -m '{"device_id": "dev-001", "power": "off", "mode": "manual", "report_interval_s": 60}'.
  4. Watch the execution: drift lists power, mode and report_interval_s, a reconcile command lands on devices/dev-001/commands, and Redis holds the shadow at twin:dev-001.
  5. Republish the same payload - the second execution logs unchanged and sends nothing.

How to extend

  • Add an anomaly branch inside change_gate for out-of-range telemetry (for example temp_c > 80) that pages on-call alongside reconciliation.
  • Age out stale twins with a scheduled flow that scans twin:* keys and alerts when last_seen is old.
  • Replace Redis with a dedicated device-shadow service by swapping the two redis.string tasks.
  • Push desired-state changes from a config flow: publish to devices/<id>/config and let this flow confirm convergence.

Links

See How

New to Kestra?

Use blueprints to kickstart your first workflows.