Adaptive QoS-aware Operator Placement for Distributed Processing of Sensor Streams

Fabian Paul · 2019

Over the past decade, the adoption of sensors and the Internet of Things (IoT) has increased rapidly. In order to analyze the high-volume data streams, efficient distributed stream processing is needed. To support complex analysis on these streams, Stream Processing Engines (SPEs) emerged. As of today, most of these SPEs run on a cluster in a central data center. However, many IoT use cases are highly distributed and the process ing must meet Quality of Service (QoS) goals; otherwise, the output quality degrades or the output becomes worthless. Centralized processing might not fit these use cases since all the data is transferred from a source to the data center through often congestd network links. On the other hand, sensor networks leverage compute capabilities while transferring the data but are often limited to simple operations on the data stream. In this work, we propose a system that bridges the gap between complex event processing and in-network processing while preserving QoS goals. We extend a state-of-the-art SPE to allow complex event processing in tree-based topologies. Thereby the system places the query operators on intermediate nodes between the data sources and the data center. The system also monitors the network link conditions and the operator outputs to avoid congested network links. In case the system detects an over-provisioned network link, it migrates the operators to a more suitable location in the topology. Thus, the system can adhere to QoS goals although the network links might not be stable. Our evaluation shows that our system offers near-constant throughput in case of changing network conditions, contrary to a centralized approach where the throughput suffers from the network regression. Furthermore, our system can reduce overall network congestion. Therefore transformations on the data stream are moved as close as possible to the data source. The system shows, in our experiments, an almost 50% reduce in received bytes at the root node in comparison to a centralized approach. \\ In den letzten zehn Jahren, hat die Verbreitung von Sensoren und dem Internet der Dinge (IoT) rapide zugenommen. Um die hochvolumigen Datenstrome zu analysieren, ist eine effiziente verteilte Verarbeitung der Datenstrome erforderlich. Zur Unterstutzung komplexer Analysen auf diesen Datenstromen, wurden Stream Processing Engines (SPEs) entwickelt. Bis heute werden die meisten dieser SPEs in einem Cluster innerhalb eines Rechenzentrums ausgefuhrt. Viele IoT-Anwendungsfalle sind jedoch stark verteilt, und die Verarbeitung muss die QoS-Ziele (Quality of Service) erfullen. Andernfalls verschlechtert sich die Ausgabequalitat oder die Ausgabe wird wertlos. Die zentralisierte Verarbeitung passt moglicherweise nicht zu diesen Anwendungsfallen, da alle Daten uber haufig uberlastete Netzwerkverbindungen von einer Datenquelle an das Rechenzentrum ubertragen werden. Andererseits nutzen Sensornetzwerke die Rechenfunktionen beim Ubertragen der Daten, beschranken sich jedoch haufig auf einfache Operationen auf dem Datenstrom. In dieser Arbeit schlagen wir ein System vor, das die Lucke zwischen komplexer Ereignisverarbeitung und Verarbeitung in Sensornetzwerken schliest und gleichzeitig die QoS-Ziele beibehalt. Wir erweitern eine hochmoderne SPE, um komplexe Ereignisverarbeitung in baumbasierten Topologien zu ermoglichen. Dabei platziert das System die Anfrageoperatoren auf Zwischenknoten zwischen den Datenquellen und dem Rechenzentrum. Das System uberwacht auch die Netzwerkverbindungsbedingungen und die Operatorenausgaben, um uberlastete Netzwerkverbindungen zu vermeiden. Wenn das System eine uberbelastete Netzwerkverbindung erkennt, migriert es die Operatoren an einen geeigneteren Ort in der Topologie. Somit kann das System die QoS-Ziele einhalten, obwohl die Netzwerkverbindungen moglicherweise nicht stabil sind. Unsere Auswertung zeigt, dass unser System bei sich andernden Netzwerkbedingungen einen nahezu konstanten Durchsatz bietet, im Gegensatz zu einem zentralisierten Ansatz, bei dem der Durchsatz unter der Netzwerkregression leidet. Daruber hinaus kann unser System die Gesamtuberlastung des Netzwerks reduzieren. Dafur werden Transformationen des Datenstroms so nah wie moglich an die Datenquellen verschoben. Das System zeigt im Vergleich zu einem zentralisierten Ansatz, in unseren Experimenten beinahe eine Halbierung der empfangenen Bytes am Wurzelknoten.

Read the paper · More papers on PaperTik