New to Kestra?
Use blueprints to kickstart your first workflows.
Realtime MQTT telemetry updates a Redis-backed digital twin, detects desired-vs-reported drift, and publishes reconcile commands only when state changes.
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.* }}.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.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.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.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.device_id, the full twin, drift, and whether a reconcile was published; errors posts failures to Slack with the device id.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.
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.
devices/<device_id>/telemetry, e.g. {"device_id": "dev-001", "power": "off", "mode": "manual", "report_interval_s": 60, "temp_c": 41.2}.desired_state input adjusted to the fields your devices report.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.disabled: false on telemetry_stream and save the flow.desired_state, e.g. mosquitto_pub -t devices/dev-001/telemetry -m '{"device_id": "dev-001", "power": "off", "mode": "manual", "report_interval_s": 60}'.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.unchanged and sends nothing.change_gate for out-of-range telemetry (for example temp_c > 80) that pages on-call alongside reconciliation.twin:* keys and alerts when last_seen is old.redis.string tasks.devices/<id>/config and let this flow confirm convergence.mqtt-realtime-alarm-reactor, mqtt-command-dispatch, mqtt-telemetry-batch-collector