EKI05 - Complex Event Processing
Complex Event Processing (CEP) verarbeitet einen kontinuierlichen Strom von Ereignissen in Echtzeit, erkennt darin Muster und leitet daraus hoehere Schluesse ab - typischerweise mit Regeln und Sliding Windows.
Überblick
Complex Event Processing (CEP) gehoert in der KI-Landkarte (nach B. Humm) zum Bereich Reasoning innerhalb der Symbolic AI / Knowledge-based AI. Dort steht es neben Logic programming und Probabilistic reasoning. Der Kernunterschied zu klassischer Datenverarbeitung: Nicht gespeicherte, ruhende Daten werden abgefragt, sondern ein fliessender Strom von Ereignissen wird laufend ausgewertet.
Definition laut Folie (5.10):
Complex Event Processing (CEP): Das Verarbeiten eines Stroms von Ereignissen und das Ableiten von Schluessen daraus.
Der Foliensatz gliedert sich in:
- Überblick
- CEP Anwendungen
- CEP Konzepte und Technologien
- Beispiel (COVID-19 Monitoring)
- Mini-Test
CEP Anwendungen
Alle Anwendungsbeispiele folgen demselben Grundmuster: Ein Ereignisstrom fliesst in eine CEP-Komponente, diese prueft auf Bedingungen/Muster und erzeugt daraus Ausgaben (Warnungen, Hinweise, Entscheidungen).
flowchart LR
A["Strom von Eingabe-Ereignissen"] --> B["CEP (Regel-Auswertung)"]
B --> C["Ausgabe: Warnung / Hinweis / Entscheidung"]
| Anwendung | Eingabestrom | CEP prueft auf | Ausgabe |
|---|---|---|---|
| Betrugserkennung (fraud detection) | Kreditkarten-Transaktionen | ungewoehnlich hohe Transaktionen in kurzer Zeit | Betrugswarnung: menschlicher Eingriff notwendig |
| Predictive Maintenance | Sensordaten einer Maschine in einer Fabrik | Muster, die auf Verschleiss hinweisen (z.B. Vibrationen) | Hinweis: Teil sollte bald ersetzt werden |
| Logistik | RFID-Signale | Bedingungen, die bestimmte Aktionen benoetigen | Logistikpartner/Kunde informieren; Warnung z.B. bei unterbrochener Kuehlkette |
| Wertpapierhandel | Boersen- und andere Marktdaten | Klassifikation der Situation fuer Kauf-/Verkauf-Entscheidungen | Kauf / Verkauf |
| COVID-19 Monitoring | Benachrichtigungen aus Testzentren und Krankenhaeusern | Berechnung von Inzidenz und Kritikalitaet | Massnahmen |
CEP Konzepte und Technologien
Was ist ein Ereignis?
Ereignis (event): Etwas besonderes, was passiert ist.
Beispiele:
- Eine Finanztransaktion
- Ein landendes Flugzeug
- Ein Sensor, der einen Messwert erfasst hat
- Eine Aenderung in einer Datenbank, ein Zustandswechsel
- Ein Tastendruck
- Ein historisches Ereignis, z.B. die franzoesische Revolution
Ereignisobjekt
Ereignisobjekt (event object): Ein Datensatz, welcher ein Ereignis repraesentiert.
Beispiele:
- Eine Kaufbestaetigung
- Eine Nachricht mit einem Kurswechsel
- Ein RFID-Datensatz mit Sensordaten
Hinweis der Folie: Der Begriff Ereignis wird oftmals stellvertretend fuer das Ereignisobjekt benutzt.
Ereignistyp
Ereignistyp (event type, event class): Der Ereignistyp spezifiziert die Struktur von Ereignisobjekten, d.h. ihrer Attribute und Datentypen.
Beispiel - die Klasse Kaufereignis mit den Attributen:
- Zeitstempel
- Kaeufer
- Produkt
- Preis
Die drei Begriffe verhalten sich damit zueinander wie in objektorientierter Programmierung: Der Ereignistyp ist die Klasse (Schema), das Ereignisobjekt eine Instanz davon, und Ereignis ist der umgangssprachliche Oberbegriff.
DBMS vs. CEP
Der zentrale Gegensatz (Folie 5.14): In einem DBMS sind die Daten persistent gespeichert und eine Query wird einmalig (volatil) darueber ausgefuehrt. Bei CEP ist umgekehrt die Regel persistent (sie laeuft dauerhaft mit), waehrend die Daten als fliessender Ereignisstrom an ihr vorbeilaufen.
| Aspekt | DBMS | CEP |
|---|---|---|
| Daten / Eingabe | Persistent (gespeichert) | Strom von Ereignissen: Fliessend |
| Query bzw. Regel | Query-Ausfuehrung: Volatil (einmalig) | Regel: Persistent (dauerhaft laufend) |
| Verarbeitungsmodell | Daten kommen zur Query | Query/Regel liegt am Datenstrom |
flowchart TB
subgraph DBMS
Q["Query (volatil)"] --> DB[("Persistente Daten")]
end
subgraph CEP
R["Regel (persistent)"] --> S["Ereignisstrom (fliessend)"]
end
Iteratives CEP
CEP wertet einen Strom von Eingabeereignissen (z.B. von Sensoren) durch Ausfuehrung von Regeln aus und erzeugt daraus einen Strom von Ausgabeereignissen (z.B. zur Entscheidungs-Unterstuetzung / decision support).
Der Prozess ist iterativ: Ausgabeereignisse koennen erneut als Eingabe dienen, sodass aus low-level Ereignissen schrittweise high-level Ereignisse abgeleitet werden.
flowchart LR
IN["Strom von Eingabe-Ereignissen (z.B. Sensoren)"] --> CEP["CEP: Ausfuehrung von Regeln"]
CEP --> OUT["Strom von Ausgabe-Ereignissen (decision support)"]
OUT -. "iterativ: low-level zu high-level" .-> CEP
Sliding "Rolling" Window
Ein Sliding Window (Gleitfenster, auch "rolling window") betrachtet immer nur einen Ausschnitt des Ereignisstroms in der juengsten Vergangenheit. Das Fenster wandert mit fortschreitender Zeit weiter (von "jetzt" ueber "frueher" bis "noch frueher").
Zwei Parameter:
- n: Dauer (duration) - die Laenge/Breite des Fensters
- m: Haeufigkeit (windowing frequency) - wie oft bzw. um welchen Schritt das Fenster weiterrueckt
Das aktuelle Fenster (die juengsten Ereignisse, hier t0, t1, t2, ...) kann zur Verarbeitung genutzt werden; davor liegen die vorherigen Fenster.
flowchart LR
A["... noch frueher ..."] --> B["Vorherige Fenster"]
B --> C["Aktuelles Fenster (n = Dauer)"]
C --> D["jetzt (t0)"]
Schematische Zeitachse laut Folie (juengste Ereignisse rechts):
... ... ... t(n+2) t(n+1) t(n) | t(n-1) t(n-2) ... t2 t1 t0 | ...
noch frueher frueher jetzt -> t
n spannt die Fensterbreite auf, m ist der Versatz zwischen aufeinanderfolgenden Fensterpositionen.
CEP Technologien
Der Foliensatz nennt folgende CEP-/Stream-Processing-Technologien:
- Apache Kafka
- Apache Flink
- Apache Spark
- Drools Fusion
- MS Azure Stream Analytics
- Oracle Stream Analytics
- SAG Apama
- SAP ESP
- SAS ESP
- TIBCO
- IBM WebSphere Business Events
CEP Technologien sind komplex (Folie 5.18):
- Komplex aufgrund der inhaerenten Verteilung
- Schwer zu installieren und handzuhaben
- Daher: Im Kurs wird mit Python das Konzept der Sliding Windows ausprobiert, ohne eine vollstaendige CEP-Technologie zu benutzen.
Hinweis: Der Foliensatz zeigt keine dedizierte CEP-Abfrage-/Regelsprache (z.B. EPL - Event Processing Language). Die einzigen konkreten Code-/Regel-Beispiele sind Python/pandas-Code und mathematische Definitionen (siehe Beispiel unten). Es werden also bewusst keine EPL-Queries behandelt.
Beispiel: COVID-19 Monitoring
Anwendung des CEP-Grundmusters auf die Pandemie-Ueberwachung:
flowchart LR
A["Strom von Benachrichtigungen aus Testzentren und Krankenhaeusern"] --> B["CEP: Berechne Inzidenz und Kritikalitaet"]
B --> C["Massnahmen"]
Datensatz
COVID-19 Singapur (data.world, hxchua/covid-19-singapore): Zeitreihendaten seit dem COVID-19-Ausbruch (Jan. 2020 bis Jan. 2022) mit vielen Features - u.a. Inzidenz (Daily Confirmed), Krankenhausdaten (Still Hospitalised, Intensive Care Unit), Todesrate (Daily/Cumulative Deaths) und weitere Spalten. Der Datensatz umfasst laut Folie 1 Datei mit 36 Spalten.
Daten laden und vorbereiten
Der Python-Code konvertiert die Datum-Strings nach datetime, um sie als DataFrame-Index zu verwenden:
import pandas as pd
from datetime import datetime
df = pd.read_csv('Covid-19_SG.csv')
df.index = df['Date'].apply(lambda s: datetime.strptime(s, '%Y-%m-%d'))
7-Tage Inzidenz (Sliding Window)
Die 7-Tage-Inzidenz wird ueber ein 7-Tage Sliding Window berechnet - die Summe der taeglichen Inzidenzen der letzten 7 Tage, normiert auf 100.000 Einwohner:
7
7 inc = SUM inc_i * 100,000 / population
i=1
- Population Singapur: 5,686 Millionen (deutsche Schreibweise, d.h. rund 5,686 Mio. Einwohner)
Hinweis: Im Original wird die Bevoelkerung als "5,686 Millionen" angegeben. Das Komma ist hier das deutsche Dezimaltrennzeichen; gemeint sind ca. 5,686 Millionen Einwohner (nicht 5686 Millionen).
Der resultierende Verlauf (Folie 5.24) zeigt die 7-Tage-Inzidenz von Jan. 2020 bis Jan. 2022 mit Spitzen bis knapp 500.
Pandas DataFrame Windowing
Zur Umsetzung des Sliding Window in Python verweist der Foliensatz auf die pandas Windowing operations (rolling). Das gezeigte Doku-Beispiel:
In [1]: s = pd.Series(range(5))
In [2]: s.rolling(window=2).sum()
Out[2]:
0 NaN
1 1.0
2 3.0
3 5.0
4 7.0
dtype: float64
Die Fenster entstehen, indem von der aktuellen Beobachtung um die Fensterlaenge zurueckgeblickt wird. Zur Visualisierung wird Matplotlib genutzt.
Iteratives CEP: Kritikalitaet aus der 7-Tage Inzidenz
Aus der 7-Tage-Inzidenz wird ein hoeheres Ereignis abgeleitet - die Kritikalitaet als Grundlage fuer Regierungsmassnahmen. Das ist das iterative Prinzip (low-level Inzidenz -> high-level Kritikalitaet).
Beispiel-Definition der Kritikalitaet:
4 , if 7 inc > 400
3 , else if 7 inc > 200
crit(7 inc) = 2 , else if 7 inc > 100
1 , else if 7 inc > 50
0 , else
7-Tage Kritikalitaets-Score (ueber ein 7-Tage Sliding Window):
7
7 crit = min ( crit(7 inc_i) )
i=1
Hinweis: Der 7-Tage-Kritikalitaets-Score verwendet das Minimum ueber die Kritikalitaetswerte der letzten 7 Tage (nicht das Maximum). Das ist so im Foliensatz definiert: Der Score erreicht eine Stufe nur, wenn die Inzidenz an allen 7 Tagen mindestens diese Stufe hielt (glaettende, "vorsichtige" Bewertung). Fuer eine Warn-Metrik ist ein Minimum untypisch - dies ist bewusst als Beispiel-Definition ("Beispiel-Definition der Kritikalitaet") formuliert und keine allgemeingueltige Vorschrift.
Der Kritikalitaets-Verlauf (Folie 5.27) zeigt Stufen von 0 bis 4, mit dem hoechsten Wert (4) im Zeitraum um Ende 2021.
Prüfungsrelevanz
- Definition CEP wortgetreu koennen: Verarbeiten eines Stroms von Ereignissen und Ableiten von Schluessen daraus.
- Drei Begriffe klar abgrenzen: Ereignis (etwas Passiertes), Ereignisobjekt (Datensatz, der es repraesentiert), Ereignistyp (Struktur/Schema: Attribute + Datentypen).
- DBMS vs. CEP: DBMS = persistente Daten + volatile Query; CEP = persistente Regel + fliessender Ereignisstrom. Dieser Umkehr-Gegensatz ist ein klassischer Pruefungspunkt.
- Iteratives CEP: Regel-Auswertung erzeugt aus Eingabeereignissen Ausgabeereignisse; Rueckkopplung fuehrt von low-level zu high-level Ereignissen.
- Sliding "Rolling" Window mit den Parametern n (Dauer/duration) und m (Haeufigkeit/windowing frequency) erklaeren koennen.
- Anwendungen aufzaehlen koennen (Betrugserkennung, Predictive Maintenance, Logistik, Wertpapierhandel, COVID-19 Monitoring).
- CEP-Technologien nennen koennen (Kafka, Flink, Spark, Drools Fusion, Azure/Oracle Stream Analytics, Apama, SAP/SAS ESP, TIBCO, IBM WebSphere Business Events).
- COVID-Beispiel: 7-Tage-Inzidenz und Kritikalitaet als Sliding-Window-Berechnung; pandas
rolling(window=...)als praktische Umsetzung.
Mini-Test
Mini-Test "Complex Event Processing":
- Was ist CEP?
- Nennen Sie CEP Anwendungen
- Was ist ein Ereignis, Ereignisobjekt und Ereignistyp?
- Wie funktioniert iteratives CEP?
- Was sind sliding ("rolling") windows?
- Nennen Sie CEP Technologien