0Pricing
AI Agents · Lektion

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

  1. MQTT-Protokoll zur Agentenintegration
  2. Zeitreihendaten in Agenten verarbeiten
  3. Automatisierte Reaktion auf Sensorereignisse
  4. Bereitstellung leichtgewichtiger Agenten am Edge
← Zurück zu AI Agents