Automatisierte Reaktion auf Sensorereignisse
Wenn Temperatur > Schwellenwert → Alarm → Aktor auslösen: agentengesteuerte IoT-Regelkreise.
Automatisierte Reaktion auf Sensorereignisse ist eine kostenlose AI Agents-Lektion auf CoddyKit. Dies ist Lektion 3 von 4. Du kannst die komplette Lektion unten kostenlos lesen – dann übst du sie direkt im Browser mit einem integrierten Code-Editor und einem KI-Tutor rund um die Uhr. Sie ist Teil des AI Agents-Lernpfads, und dein Fortschritt wird über Web und CoddyKit-App synchronisiert. Der AI Agents-Kurs umfasst insgesamt 4 Lektionen.
Automatisierte ereignisgesteuerte Agentenantworten
Wenn ein Sensor einen Schwellenwert überschreitet, muss der Agent automatisch und ohne menschliches Eingreifen reagieren. Die zentralen Herausforderungen bestehen darin, zu entscheiden, welche Aktion ausgeführt werden soll, sicherzustellen, dass dasselbe Ereignis keine doppelten Aktionen auslöst, und einen Cooldown-Zeitraum einzuhalten, damit der Agent Aktoren nicht mit Befehlen überflutet.
Aktionsrichtlinien definieren
Eine Aktionsrichtlinie ordnet Sensorbedingungen den Antworten des Agenten zu. Definieren Sie Richtlinien deklarativ, damit sie leicht zu lesen und zu ändern sind, ohne den Logikcode anzupassen. Jede Richtlinie enthält eine Bedingung, eine Priorität und eine oder mehrere Aktionen.
ACTION_POLICIES = [
{
'name': 'HIGH_TEMP_ALERT',
'topic': 'sensors/temperature',
'condition': lambda v: v > 38,
'priority': 'critical',
'actions': ['TURN_ON_COOLING', 'ALERT_MAINTENANCE', 'LOG_EVENT']
},
{
'name': 'HIGH_TEMP_WARNING',
'topic': 'sensors/temperature',
'condition': lambda v: 35 < v <= 38,
'priority': 'warning',
'actions': ['ALERT_MAINTENANCE', 'LOG_EVENT']
},
{
'name': 'LOW_HUMIDITY',
'topic': 'sensors/humidity',
'condition': lambda v: v < 30,
'priority': 'warning',
'actions': ['TURN_ON_HUMIDIFIER', 'LOG_EVENT']
}
]
def match_policies(topic: str, value: float) -> list:
return [
p for p in ACTION_POLICIES
if p['topic'] == topic and p['condition'](value)
]
if __name__ == '__main__':
matches = match_policies('sensors/temperature', 39)
print('Matched policies for temperature=39:')
for p in matches:
print(f" {p['name']} ({p['priority']}): {p['actions']}")
Aktionswarteschlange
Eine Aktionswarteschlange entkoppelt die Ereigniserkennung von der Aktionsausführung. Ereignisse werden in die Warteschlange eingefügt; ein Worker entnimmt sie und führt sie aus. Dadurch wird verhindert, dass die MQTT-Empfangsschleife blockiert, und bei einem Fehler einer Aktion sind Wiederholungsversuche möglich.
import queue
import threading
from datetime import datetime
action_queue: queue.Queue = queue.Queue(maxsize=500)
def enqueue_action(action_name: str, context: dict, priority: str = 'normal'):
item = {
'action': action_name,
'context': context,
'priority': priority,
'enqueued_at': datetime.utcnow().isoformat()
}
try:
action_queue.put_nowait(item)
print(f'Enqueued: {action_name}')
except queue.Full:
print(f'WARNING: Action queue full, dropping {action_name}')
def action_worker(executor_fn):
"""Run in a background thread, executing actions from the queue."""
while True:
item = action_queue.get()
try:
executor_fn(item['action'], item['context'])
except Exception as e:
print(f'Action failed: {item["action"]} — {e}')
finally:
action_queue.task_done()
# Start worker thread:
# worker_thread = threading.Thread(target=action_worker, args=(execute_action,), daemon=True)
# worker_thread.start()
if __name__ == '__main__':
enqueue_action('TURN_ON_COOLING', {'zone': 'server-room'}, priority='critical')
enqueue_action('LOG_EVENT', {'msg': 'temperature nominal'})
print('Queue size:', action_queue.qsize())
Ereignisdeduplizierung
Ohne Deduplizierung erzeugt eine Temperatur, die bei 1 Messwert pro Sekunde 10 Minuten lang über 38 °C bleibt, 600 identische Ereignisse. Die Deduplizierung stellt sicher, dass dieselbe Kombination aus (topic, condition, action) pro Ereignis nur einmal ausgelöst wird. Sie wird zurückgesetzt, sobald die Bedingung nicht mehr erfüllt ist.
class EventDeduplicator:
def __init__(self):
# active_events: (topic, policy_name) -> event_start_time
self._active: dict = {}
def is_new_event(self, topic: str, policy_name: str) -> bool:
key = (topic, policy_name)
return key not in self._active
def mark_active(self, topic: str, policy_name: str):
self._active[(topic, policy_name)] = datetime.utcnow()
def clear_event(self, topic: str, policy_name: str):
key = (topic, policy_name)
if key in self._active:
duration = (datetime.utcnow() - self._active.pop(key)).seconds
print(f'Event cleared: {policy_name} (lasted {duration}s)')
def clear_topic_if_normal(
self, topic: str, value: float, normal_fn
):
if normal_fn(value):
keys = [k for k in self._active if k[0] == topic]
for k in keys:
self.clear_event(k[0], k[1])
dedup = EventDeduplicator()
dedup.mark_active('sensors/temperature', 'HIGH_TEMP_ALERT')
print('New event?', dedup.is_new_event('sensors/temperature', 'HIGH_TEMP_ALERT'))Cooldown-Zeitraum
Selbst wenn ein Ereignis beendet und erneut ausgelöst wird, verhindert ein Cooldown-Zeitraum ein zu schnelles erneutes Auslösen. Kombinieren Sie die Deduplizierung (einmaliges Auslösen, solange die Bedingung erfüllt ist) mit dem Cooldown (N Minuten nach dem Ende der Bedingung warten, bevor derselbe Alarm erneut ausgelöst werden darf).
from datetime import datetime, timedelta
class CoolDownManager:
def __init__(self, cool_down_minutes: int = 15):
self.cool_down = timedelta(minutes=cool_down_minutes)
self._cleared_at: dict = {} # (topic, policy) -> cleared_datetime
def is_in_cool_down(self, topic: str, policy_name: str) -> bool:
key = (topic, policy_name)
cleared_at = self._cleared_at.get(key)
if cleared_at is None:
return False
return datetime.utcnow() - cleared_at < self.cool_down
def record_clear(self, topic: str, policy_name: str):
self._cleared_at[(topic, policy_name)] = datetime.utcnow()
def time_remaining(self, topic: str, policy_name: str) -> int:
key = (topic, policy_name)
cleared_at = self._cleared_at.get(key)
if cleared_at is None:
return 0
elapsed = datetime.utcnow() - cleared_at
remaining = self.cool_down - elapsed
return max(0, int(remaining.total_seconds()))
cooldown = CoolDownManager(cool_down_minutes=15)
cooldown.record_clear('sensors/temperature', 'HIGH_TEMP_ALERT')
print('In cool-down?', cooldown.is_in_cool_down('sensors/temperature', 'HIGH_TEMP_ALERT'))Aktionsentscheidung mit Unterstützung eines LLM
Bei komplexen Situationen — mehreren gleichzeitigen Alarmen, widersprüchlichen Richtlinien oder ungewöhnlichen Kombinationen von Messwerten — übertragen Sie die Entscheidung an das LLM. Das LLM erhält den vollständigen Sensorkontext und empfiehlt einen priorisierten Aktionsplan.
import anthropic
import json
def llm_decide_actions(
sensor_readings: dict,
active_policies: list
) -> list:
client = anthropic.Anthropic(api_key='YOUR_API_KEY')
context = json.dumps({
'readings': sensor_readings,
'triggered_policies': [p['name'] for p in active_policies]
}, indent=2)
prompt = (
f'Current sensor state:\n{context}\n\n'
'Multiple alert policies are active. '
'Recommend an ordered list of actions to take. '
'Consider conflicting effects (e.g., humidifier and cooling may conflict).\n'
'Return JSON: {"recommended_actions": [str], "reasoning": str}'
)
response = client.messages.create(
model='claude-opus-4-5', max_tokens=512,
messages=[{'role': 'user', 'content': prompt}]
)
return json.loads(response.content[0].text)Aktionen über MQTT ausführen
Aktionen werden ausgeführt, indem Befehlsnachrichten an gerätespezifische MQTT-Topics veröffentlicht werden. Der Payload des Befehls folgt einem Standardschema: Aktionsname, Parameter, Request-ID zur Bestätigung und TTL (der Befehl läuft ab, wenn das Gerät zu lange offline ist).
import json
import uuid
from datetime import datetime, timedelta
ACTION_TOPICS = {
'TURN_ON_COOLING': 'devices/hvac/commands',
'TURN_OFF_COOLING': 'devices/hvac/commands',
'TURN_ON_HUMIDIFIER': 'devices/humidifier/commands',
'ALERT_MAINTENANCE': 'notifications/maintenance',
'LOG_EVENT': 'logs/agent_events'
}
def execute_action(action_name: str, context: dict, mqtt_client) -> str:
topic = ACTION_TOPICS.get(action_name)
if not topic:
print(f'No topic defined for action: {action_name}')
return 'unknown_action'
request_id = str(uuid.uuid4())[:8]
ttl = (datetime.utcnow() + timedelta(minutes=5)).isoformat()
payload = json.dumps({
'action': action_name,
'request_id': request_id,
'context': context,
'ttl': ttl
})
mqtt_client.publish(topic, payload, qos=1)
print(f'Executed {action_name} -> {topic} (req={request_id})')
return request_id
if __name__ == '__main__':
class FakeMQTT:
def publish(self, topic, payload, qos=1):
pass
execute_action('TURN_ON_COOLING', {'zone': 'server-room'}, FakeMQTT())
Bestätigung von Aktionen
Geräte sollten empfangene Befehle bestätigen, indem sie eine Nachricht an ein Bestätigungs-Topic veröffentlichen. Der Agent abonniert Bestätigungs-Topics und kann den Befehl wiederholen, wenn innerhalb eines Zeitlimits keine Bestätigung empfangen wird.
import threading
from collections import defaultdict
class AckTracker:
def __init__(self, timeout_seconds: int = 30):
self.timeout = timeout_seconds
self._pending: dict = {} # request_id -> {'action', 'send_time', 'ack_event'}
def register(self, request_id: str, action_name: str):
event = threading.Event()
self._pending[request_id] = {
'action': action_name,
'send_time': datetime.utcnow(),
'ack_event': event
}
# Schedule timeout check
t = threading.Timer(self.timeout, self._on_timeout, args=[request_id])
t.daemon = True
t.start()
def acknowledge(self, request_id: str):
entry = self._pending.pop(request_id, None)
if entry:
entry['ack_event'].set()
print(f'Ack received for {entry["action"]} (req={request_id})')
def _on_timeout(self, request_id: str):
if request_id in self._pending:
action = self._pending.pop(request_id)['action']
print(f'TIMEOUT: No ack for {action} (req={request_id}) — retry?')Vollständige Pipeline für Sensoreignisse
Alle Komponenten zusammengefügt: MQTT-Empfang → Richtlinienabgleich → Prüfung auf Deduplizierung und Cooldown → Aktionen in die Warteschlange einreihen → Worker führt sie über eine MQTT-Veröffentlichung aus → Bestätigungen verfolgen. Diese Architektur verarbeitet Tausende von Sensoreignissen pro Minute, ohne zu blockieren.
class IoTAgentPipeline:
def __init__(self, mqtt_client):
self.mqtt = mqtt_client
self.dedup = EventDeduplicator()
self.cooldown = CoolDownManager(cool_down_minutes=15)
self.ack_tracker = AckTracker(timeout_seconds=30)
def on_sensor_message(self, topic: str, value: float):
policies = match_policies(topic, value)
self.dedup.clear_topic_if_normal(
topic, value,
normal_fn=lambda v: v <= 35 # below warning threshold
)
for policy in policies:
name = policy['name']
if not self.dedup.is_new_event(topic, name):
continue # already active, skip
if self.cooldown.is_in_cool_down(topic, name):
print(f'In cool-down: {name}')
continue
self.dedup.mark_active(topic, name)
ctx = {'topic': topic, 'value': value, 'policy': name}
for action in policy['actions']:
enqueue_action(action, ctx, policy['priority'])Eskalationspfad
Einige Situationen erfordern eine Eskalation an Menschen: wiederholtes Ausbleiben der Bestätigung eines Befehls, anhaltende kritische Zustände oder mehrere widersprüchliche Richtlinien. Definieren Sie einen Eskalationspfad, der eine Push-Benachrichtigung sendet oder ein Ticket erstellt.
import requests
def escalate_to_human(
reason: str,
sensor_data: dict,
webhook_url: str = 'https://hooks.slack.com/services/YOUR/SLACK/WEBHOOK'
):
message = {
'text': (
f'*IoT Agent Escalation* \n'
f'Reason: {reason}\n'
f'Sensor data: {sensor_data}\n'
f'Time: {datetime.utcnow().isoformat()}'
)
}
try:
response = requests.post(webhook_url, json=message, timeout=5)
response.raise_for_status()
print(f'Escalation sent: {reason}')
except requests.RequestException as e:
print(f'Escalation failed: {e}')
# Fall back: log to file
with open('escalations.log', 'a') as f:
import json
f.write(json.dumps({'reason': reason, 'data': sensor_data}) + '\n')Die Ereignispipeline testen
Schreiben Sie vor der Bereitstellung Ihrer Ereignispipeline in der Produktion automatisierte Tests, die Sensoreignisse simulieren und überprüfen, ob die richtigen Aktionen in die Warteschlange eingereiht werden. Testen Sie jede Richtlinie einzeln sowie das Verhalten der Deduplizierung und den Ablauf des Cooldowns.
import time
def test_high_temp_policy_fires_once():
dedup = EventDeduplicator()
cooldown = CoolDownManager(cool_down_minutes=0) # disable cooldown for test
pipeline = IoTAgentPipeline(None)
pipeline.dedup = dedup
pipeline.cooldown = cooldown
actions_fired = []
action_queue.queue.clear()
# Fire same event 5 times in a row
for _ in range(5):
pipeline.on_sensor_message('sensors/temperature', 40.0)
# Only 1 set of actions should have been enqueued
actions = list(action_queue.queue)
assert len(actions) > 0, 'At least one action should fire'
print(f'Actions enqueued: {len(actions)} (expected: just 1 event worth)')
return True
result = test_high_temp_policy_fires_once()
print('Test passed:', result)Wissenscheck
Was ist der wichtigste Zweck der Ereignisdeduplizierung in einer Pipeline für Sensoreignisse?
Zusammenfassung: Automatisierte Reaktion auf Sensoreignisse
Sehr gut! Das haben Sie gelernt:
- Aktionsrichtlinien: deklarative Zuordnung von Bedingungen zu Aktionen mit Prioritäten
- Aktionswarteschlange: Erkennung und Ausführung entkoppeln; ein Worker verarbeitet die Aktionen
- Deduplizierung: einmal pro Ereignis auslösen, nicht einmal pro Messwert
- Cooldown: erneutes Auslösen unmittelbar nach dem Ende einer Bedingung verhindern
- LLM-Eskalation: komplexe Situationen mit mehreren Richtlinien an die Schlussfolgerungen des LLM delegieren
- Bestätigungsverfolgung: nicht bestätigte Befehle erkennen und wiederholen
Als Nächstes: die Bereitstellung kompakter Agenten auf Edge-Geräten wie dem Raspberry Pi.
Häufig gestellte Fragen
Ist die Lektion „Automatisierte Reaktion auf Sensorereignisse“ kostenlos?
Ja — der vollständige Text von „Automatisierte Reaktion auf Sensorereignisse“ ist hier im Web kostenlos zu lesen. Um sie interaktiv zu üben (integrierter Code-Editor und 24/7 KI-Tutor) und den Rest des AI Agents-Kurses freizuschalten, upgrade auf CoddyKit PRO. Der AI Agents-Kurs umfasst insgesamt 4 Lektionen.
Was lerne ich in „Automatisierte Reaktion auf Sensorereignisse“?
Wenn Temperatur > Schwellenwert → Alarm → Aktor auslösen: agentengesteuerte IoT-Regelkreise. Du übst AI Agents mit praktischem Code, den du direkt im Browser ausführst, und ein 24/7 KI-Tutor beantwortet deine Fragen während du die Lektion bearbeitest.
Brauche ich Erfahrung, um AI Agents zu starten?
Keine Vorkenntnisse erforderlich. AI Agents auf CoddyKit ist für Anfänger bis fortgeschrittene Lernende strukturiert, sodass du hier starten oder von Anfang an beginnen und in deinem eigenen Tempo voranschreiten kannst. Dies ist Lektion 3 von 4.
Wie lange dauert die Lektion „Automatisierte Reaktion auf Sensorereignisse“?
Die meisten CoddyKit-Lektionen dauern etwa 5–10 Minuten. Jede ist kompakt und interaktiv, sodass du stetig Fortschritte machst und genau dort weitermachst, wo du aufgehört hast – im Web und in der App.
Kann ich in dieser AI Agents-Lektion Code schreiben und ausführen?
Ja. Jede AI Agents-Lektion enthält einen integrierten Code-Editor, sodass du echten Code direkt in deinem Browser schreibst und ausführst und sofort KI-Feedback erhältst — ohne lokale Einrichtung erforderlich.
Alle Lektionen in diesem Kurs
- MQTT-Protokoll zur Agentenintegration
- Zeitreihendaten in Agenten verarbeiten
- Automatisierte Reaktion auf Sensorereignisse
- Bereitstellung leichtgewichtiger Agenten am Edge