Full text
University of Minho School of Engineering Pedro Miguel de Leal Meireles Pereira MulletBench: Multi-layer Edge Time Series Database Benchmark october 2023
University of Minho School of Engineering Pedro Miguel de Leal Meireles Pereira MulletBench: Multi-layer Edge Time Series Database Benchmark Master’s Dissertation Master in Informatics Engineering Dissertation supervised by Fábio André Castanheira Luís Coelho (advisor) João Tiago Medeiros Paulo (co-advisor) Luís Manuel Meruje Ferreira (co-advisor in Research Center) october 2023
Copyright and Terms of Use for Third Party Work This dissertation reports on academic work that can be used by third parties as long as the internationally accepted standards and good practices are respected concerning copyright and related rights. This work can thereafter be used under the terms established in the license below. Readers needing authorization conditions not provided for in the indicated licensing should contact the author through the RepositóriUM of the University of Minho. License granted to users of this work: CC BY-NC https://creativecommons.org/licenses/by-nc/4.0/ i
Acknowledgements This work was only possible due to the support of many people, to whom I would like to express my gratitude. First and foremost, I would like to thank my supervisors, Professor Fábio Coelho and Professor João Paulo, for their guidance and insight throughout the development of this dissertation. I would also like to express my gratitude to Luís Meruje for his continuous support and help on all aspects of this work. Without their support, it would not have been possible. To all my colleagues at High-Assurance Software Laboratory (HASLab) who provided much-needed breaks, laughter, and encouragement. To my family for their support and understanding throughout this academic journey. To my girlfriend for always being by my side regardless of distance, and for always believing in me. Finally, I would like to thank INESC TEC that co-funded this work through Component 5 - Capitalization and Business Innovation, integrated in the Resilience Dimension of the Recovery and Resilience Plan within the scope of the Recovery and Resilience Mechanism (MRR) of the European Union (EU), framed in the Next Generation EU, for the period 2021 - 2026, within project ATE, with reference 56 (10113/BII-E_B4/2023) ii
Statement of Integrity I hereby declare having conducted this academic work with integrity. I confirm that I have not used plagiarism or any form of undue use of information or falsification of results along the process leading to its elaboration. I further declare that I have fully acknowledged the Code of Ethical Conduct of the University of Minho. University of Minho, Braga, october 2023 Pedro Miguel de Leal Meireles Pereira iii Assinado por: Pedro Miguel de Leal Meireles Pereira Num. de Identificação: 14673529 Data: 2023.10.31 17:02:35+00'00'
Abstract Internet of Things (IoT) systems generate massive amounts of time series data that need to be stored for historical analysis. As a result, Database Management Systems (DBMSs) for these scenarios have particular requirements in their ability to ingest large amounts of data and to optimise aggregation, filtering and time-ranged queries over this data, which are essential for historical analysis. Through the use of Fog Computing, combining both Edge and Cloud layers, it is possible to achieve reduced latency and increased scalability, privacy and connectivity through the Edge, while still benefiting from the enhanced computing and storage power of the Cloud. This has led to the development of Fog DBMSs. Database benchmarking allows standardising performance assessment and comparison of different solutions. However, current time series database benchmarking tools are not designed for multi-layer architectures, such as the ones used by Edge-Cloud hybrid DBMSs. This thesis proposes MulletBench, a benchmarking tool that is able to evaluate the internal loadbalancing capabilities of a multi-layer Time Series Database Management System (TSDBMS). This is achieved by integrating automated deployment features, per-node and per-layer performance and system resource metrics, allowing for a more detailed analysis of the SUTs’ performance than previously possible. The performance of InfluxDB and IoTDB is evaluated using the developed tool, comparing their performance in multiple workloads and deployment scenarios. Results show that the Edge layer can be used to improve performance by distributing the workload over multiple layers and performing downsampling at the Edge layer, increasing overall throughput and reducing latency at the Cloud. These conclusions are enabled by MulletBench’s novel features, and would not have been possible with previously existing solutions. Keywords Benchmark, Internet of Things, Fog, Time Series, Databases iv
Resumo Os Sistemas de Internet das Coisas (IdC) geram enormes quantidades de dados de séries temporais que precisam de ser armazenados para análise histórica. Consequentemente, os sistemas de gestão de bases de dados (SGBD) para estes cenários têm requisitos particulares não só para a sua capacidade de ingerir grandes quantidades de dados mas também para a sua capacidade de otimizar queries de agregação ou filtragem e em intervalos de tempo sobre estes dados, que são essenciais para análise histórica. Através do uso de computação em Fog , combinando as camadas de Edge e Cloud , é possivel conseguir latência reduzida e maior escalabilidade, privacidade e conectividade através da Edge , beneficiando simultaneamente do superior poder de computação e armazenamento da Cloud . Isto levou ao desenvolvimento de sistemas de gestão de bases de dados Fog . O Benchmarking de bases de dados permite a estandardização do método de avaliação de desempenho e a comparabilidadde de resultados de diferentes soluções. Contudo, as atuais ferramentas de benchmarking para bases de dados de séries temporais não estão desenhadas para arquiteturas multicamada, tais como as utilizadas por SGBDs híbridos Edge-Cloud . A tese propõe o MulletBench, uma ferramenta de benchmarking capaz de avaliar as capacidades de balanceamento interno de carga de um SGBD de séries temporais multi-camada. Esta ferramenta alcançao através da integração de funcionalidades de automação de deployment e metricas de desempenho e de recursos de sistema por nó e por camada, permitindo uma análise mais detalhada do desempenho dos sistemas testados do que era previamente possível. O desempenho das bases de dados InfluxDB e IoTDB é avaliado usando a ferramenta desenvolvida, comparando o seu desempenho com multiplas cargas e em multiplos cenários. Os resultados demonstram que a camada Edge pode ser usada para melhorar o desempenho destas, através da distribuição da carga por múltiplas camadas e realização de downsampling na camada Edge , aumentando o débito total e reduzindo a latência na Cloud . Estas conclusões são possíveis graças às funcionalidades inovadoras do MulletBench, e não seriam possíveis com soluções existentes anteriormente. Palavras-chave Benchmark , Internet das Coisas, Fog, Séries Temporais, Bases de Dados v
Contents 1 Introduction 1 1.1 Problem Statement .................................. 2 1.2 Objectives and Contributions ............................. 3 1.3 Results ........................................ 3 1.4 Document Structure .................................. 4 2 Background and Related Work 5 2.1 IoT and Edge Computing ............................... 5 2.2 Time Series Databases ................................ 6 2.2.1 Cloud-first Time Series Databases ...................... 7 2.2.2 Multi-layer Time Series Databases ...................... 8 2.3 Benchmarking .................................... 10 2.4 Survey on TSDB Benchmarking Systems ........................ 13 2.4.1 Workloads .................................. 14 2.4.2 Data ..................................... 16 2.4.3 Metrics ................................... 18 2.4.4 Orchestration ................................ 20 2.4.5 Discussion .................................. 21 3 MulletBench 23 3.1 Design Principles ................................... 23 3.2 Architecture ...................................... 26 3.2.1 Orchestrator ................................. 26 3.2.2 Client .................................... 28 3.3 Implementation .................................... 31 3.3.1 Orchestrator ................................. 31 vi
25 Client IoTDB options ................................. 96 26 Client configuration for insertion workload targeting InfluxDB in 2 Edge + Cloud scenario 97 27 Client configuration for insertion workload targeting IoTDB in 2 Edge + Cloud scenario . 97 28 Client configuration for query workload Cloud-only scenario ............... 98 29 Pre-population Client configuration for mixed workload Cloud-only scenario ....... 98 30 Insert Client configuration for mixed workload Cloud-only scenario ........... 98 31 Query Client configuration for mixed workload Cloud-only scenario ........... 99 32 Client configuration for insertion workload targeting InfluxDB in 2 Edge + Cloud scenario 99 33 Client configuration for insertion workload targeting IoTDB in 2 Edge + Cloud scenario . 99 34 Query Client configuration for query workload in 2 Edge + Cloud scenario ........100 35 Pre-population Client configuration for mixed workload 2 Edge + Cloud scenario . . . . 100 36 Insert Client configuration for mixed workload in 2 Edge + Cloud scenario ........100 37 Query Client configuration for mixed workload in 2 Edge + Cloud scenario . . . . . . . 101 38 Cloud query Client configuration for mixed workload with downsampling in 2 Edge + Cloud scenario ....................................101 39 Client configuration for insertion workload targeting InfluxDB in 4 Edge + Cloud scenario 101 40 Client configuration for insertion workload targeting IoTDB in 4 Edge + Cloud scenario . 102 41 Query Client configuration for query workload in 4 Edge + Cloud scenario ........102 42 Pre-population Client configuration for mixed workload 4 Edge + Cloud scenario . . . . 102 43 Insert Client configuration for mixed workload in 4 Edge + Cloud scenario ........103 44 Query Client configuration for mixed workload in 4 Edge + Cloud scenario . . . . . . . 103 xiii
Acronyms DBMS Database Management System. GAN Generative Adversarial Network. IoT Internet of Things. LSM Log-Structured Merge-Tree. PCP Performance Co-Pilot. PSoTSDBs Performance Study of Time Series Databases. SUT system under test. TSDBMS Time Series Database Management System. TSM Time-Structured Merge Tree. WAL Write Ahead Log. YCSB Yahoo! Cloud Serving Benchmark. xiv
Chapter 1 Introduction The Internet of Things (IoT) concept consists of a network of devices referred to as Things , that is, any device that is either a sensor or an actuator in an IoT application. Sensors are commonly connected through a sensor network , allowing for their collaboration and enabling the collection, processing, analysis and dissemination of the data they produce. IoT systems have contributed to an exponential increase in the number of connected devices and generated data, resulting in a projection that by 2025 there will be 175 Zettabytes of data, of which 79.4 Zettabytes are to be created by IoT devices (Wang and Wang [2020]). The need for storing such a high volume of data, as well as performing historical time-ranged queries, among other reasons, has led to a rise in popularity of Time Series Database Management Systems (TSDBMSs) (dbe [2022b]), which are optimised for high data ingestion rates and timestamped data, suiting IoT scenarios in which data is generated by sensors. A possible solution for storing data generated by IoT devices is doing it at the Edge, which consists of computing and storage resources located at the Internet’s edge, close to mobile devices and sensors (Satyanarayanan [2017]). This leads to benefits such as lower latency, reduced network strain and improved data privacy, but this solution also suffers from drawbacks, with limited hardware and heterogeneous resource capabilities. Combining both Edge and Cloud allows for the balance of data storage and processing between both layers. This solution also brings additional benefits, enabling the choice of whether data should be kept temporarily or permanently at each layer, and where to perform data processing tasks. This makes use of a concept called Fog or Edge-Cloud hybrid Computing, leveraging both the Edge and Cloud layers, to achieve reduced latency and increased scalability, privacy, and connectivity through the Edge, while still benefiting from the enhanced computing and storage capabilities of the Cloud. The most common sensor network architecture follows the sink node approach, in which sensorgenerated data is offloaded to a local, more powerful node called the sink node, which in turn communicates with the Cloud layer, achieving some of the aforementioned benefits (Perera et al. [2013]). There 1
are several database solutions that make use of both the Edge and Cloud layers, each using different mechanisms to achieve this and consequently having different performance advantages and drawbacks. Since different database solutions have differing implementations and focus on different aspects of their performance, they show significant performance differences. This results in a need for evaluation tools that provide metrics which allow for informed decision-making, either by database researchers, vendors, or customers. As a result, the establishment of industry-standard benchmarks can offer a means for standardised performance assessment and comparison of competing ideas and implementations (Seltzer et al. [1999]). While there are benchmarks specifically targetting TSDBMSs, there are no current solutions that target multi-layer TSDBMSs for sensor networks which distribute storage between the Edge and Cloud layers. 1.1 Problem Statement There are no current benchmark implementations capable of adequately evaluating hybrid Edge-Cloud database solutions for time series data generated in sink node sensor networks. Benchmarking databases distributed over multiple layers is a complex task, due to the high flexibility of architectures, as well as the need for metrics that can accurately represent not only the database’s performance as a whole but also the performance of each layer or even node. While there are benchmarks that evaluate time series DBMSs, none of them are sensitive to important aspects of Edge-Cloud distributed DBMSs, such as how the latency of operations is affected by which node replies to the request and how these databases scale with different numbers of both Edge and Cloud nodes, as well as with different distributions of nodes between layers. Metrics that represent these aspects, such as per-node ingestion throughput and query latency, are important to understand how an Edgeenabled sensor network Database Management System (DBMS) behaves in real-world scenarios and with real-world workloads, as they enable the analysis of the database’s internal load balancing and, consequently, a more in-depth understanding of its behaviour and performance. Furthermore, features that automate the deployment of database nodes, improving not only the ease of use of the benchmark, but also the reproducibility of its results are uncommon, with no solution offering it for multiple layers. Given that the deployment of a database cluster is a complex task, requiring the installation and configuration of multiple nodes, automation is very useful and eases the evaluation of 2
different database node configurations. 1.2 Objectives and Contributions Since a lack of existing benchmarks for this domain has been identified, this thesis presents a benchmarking system capable of providing metrics that accurately measure the performance of a time series DBMS deployed in an Edge-Cloud distributed architecture. It focuses on addressing the highlighted features that have not been observed in any existing solution, such as per-node metrics and orchestration features. To achieve this, this thesis provides the following contributions: • An in-depth study and survey of currently existing TSDBMS benchmarks. This study includes a new taxonomy of the existing benchmarks, analysing their workloads, data, metrics and orchestration features and concludes that existing solutions lack metrics and orchestration features suitable for the evaluation of Edge-Cloud hybrid TSDBMSs. • A new TSDBMS benchmarking tool, named MulletBench, with novel features such as per-node and per-layer metrics, automatic multi-layer database deployment and orchestration of distributed benchmark clients. • A validation of MulletBench and evaluation of InfluxDB InfluxData [2013] and Apache IoTDB Wang et al. [2020]. This evaluation includes a comparison of the targeted databases in multiple workloads and deployment scenarios, and its results show that IoTDB generally presents superior performance to InfluxDB, and that Edge-Cloud hybrid deployments enable significant performance improvements when compared to single Cloud node deployments. Additionally, a flaw in both systems’ replication mechanisms is identified when these are subjected to high throughput workloads. 1.3 Results The development of this dissertation resulted in an open-source TSDBMS benchmark providing multiple novel features 1. Issues were created in the respective GitHub repositories of the targeted databases documenting the 1MulletBench: https://github.com/pedropereira98/MulletBench 3
identified replication flaws 2 3 4. 1.4 Document Structure This document is organised into four chapters. Chapter 2 gives a more in-depth look at the subjects of IoT,Time Series Database Management System (TSDBMS) and DBMS benchmarking, as well as a survey looking into the existing TSDBMS benchmarking solutions. Chapter 3 presents MulletBench, including a specification of its overall features, an architectural overview of the developed system and implementation details. Chapter 4 goes into the evaluation and results, presenting the used methodology, the results of the performed feature validation and an evaluation of multiple databases using several workloads and scenarios. This allowed for comparisons between not only these databases but also their different deployment options. Finally, Chapter 5 presents the final conclusions that could be drawn from the developed work, as well as how it could be further expanded. 2InfluxDB issue: https://github.com/influxdata/influxdb/issues/24434 3IoTDB replication issue: https://github.com/apache/iotdb/issues/11428 4IoTDB unsequenced issue: https://github.com/apache/iotdb/issues/11441 4
Chapter 2 Background and Related Work This chapter introduces relevant background concerning IoT and Edge Computing, going into its characteristics as well as the data handled by it, TSDBMSs, what they achieve and how, and DBMS benchmarking. Afterwards, a survey analysing the existing TSDBMS benchmarking solutions is presented, from which came the conclusion that there is an opportunity for contributions. 2.1 IoT and Edge Computing IoT systems can have several applications with some of the most common being smart systems, such as Smart Agriculture, Smart Cities, Smart Healthcare and Smart Industry (also known as Industrial IoT) (Kamienski et al. [2019]; Wollschlaeger et al. [2017]). The IoT attempts to increase the automation processes present in these domains and ultimately achieve higher efficiency. At the core of the IoT are its sensors, which have limited computing power and may have limited energy. These sensors are potentially heterogeneous and communicate among them using wired or wireless technologies, making up sensor networks (Perera et al. [2013]). Although the concept of a sensor network has existed since before the emergence of the IoT, it was previously only used to achieve specific purposes in limited domains. The IoT can therefore be seen as a general-purpose extension of sensor networks, focusing on the interactions between the users and the physical environment, rather than on sensing and reporting low-level data (Firner et al. [2011]). As IoT devices have limited computing power, they often need to offload tasks to more powerful Cloud nodes, allowing for benefits such as increased computing and storage capacity. New concepts such as Edge and Fog computing attempt to provide better latency and connectivity when compared to systems that offload all data to the Cloud. Edge computing also enhances the management, storage, and processing power of data generated by connected devices by using small compute and storage nodes at the edge of the network, and as close 5
as one hop to the IoT devices (Reale [2019]), providing better latency than Cloud offloading systems. Fog or Edge-Cloud hybrid computing, on the other hand, consists of deploying nodes located closer to end devices as a way to complement the Cloud, aggregating, processing and/or pre-processing data, essentially extending the Cloud’s capabilities and allowing for a balance between processing power and latency (Yousefpour et al. [2019]). The most common architecture used in sensor networks for Fog computing is the sink node sensor network, in which the network’s sensor nodes offload their collected data to a local, more powerful node, called the sink node or gateway (Perera et al. [2013]). This central point of the sensor network is often responsible for the pre-processing of the data it receives and stores, before offloading it to a Cloud TSDBMS (Buevich et al. [2013]). This architecture is one possible way of achieving a multi-layer data storage and processing solution, making use of both the Edge and Cloud layers, and will for that reason be the one on which this work focuses. Another possible approach is in-network processing , in which each sensor node stores its own data locally and queries are disseminated through the network and processed by the sensor nodes. The queried data is sent back to parent nodes where it is combined and once more sent to their parent nodes until the final data reaches the gateway (Diallo et al. [2013]). 2.2 Time Series Databases The data being generated and used in IoT applications needs to be properly understood for it to be correctly stored. It usually consists of time series data, which has a format associating a timestamp for a recorded instant in time with one or multiple values (also sometimes called fields), as well as optionally having tags which give additional information about the measurement, also known as metadata, and are used to more easily categorise the data. IoT data can often be categorised as Big Data (Cai et al. [2016]), which is generally considered to have the following characteristics (Han et al. [2017]): •Volume:IoT devices generate massive amounts of data, estimated to reach 79.4 Zettabytes by 2025 (Wang and Wang [2020]). •Velocity:IoT data is generated at a high rate, with a single device being able to generate hundreds of data points every 20ms, and big deployments reaching the billions of data points generated daily (Hall [2020]). In time series data, this specifically relates to the timestamp interval between samples and the sampling frequency of the sensors. •Variety: As IoT devices can be very varied, ranging from environmental sensors to GPSs or cars, 6
the data and data formats they generate can also vary widely (Sasaki [2021]). •Veracity: Since IoT data is often generated by sensors and these are not 100% accurate, this data becomes inherently imprecise, containing noise or errors. In addition, IoT data also has some unique characteristics (Cai et al. [2016]): •Multisource:IoT applications acquire data from different sensors with different characteristics and in parallel. •Low-Level:IoT data from sensors is low-level and has weak semantics, meaning processing is needed for valuable information to be extracted from it. All of these characteristics present in IoT data introduce challenges for systems that aim to store, process or analyse it. DBMSs for these applications need to be able to store massive amounts of data at a very high rate. On top of being able to withstand the ingestion of this data, these DBMSs also need to support and, ideally, optimise certain querying operations, as IoT applications often require either real-time or historical analysis. These applications make heavy use of aggregation operations over specific time ranges and hardly ever require data granularity at the individual data point level, allowing for varying query time resolutions or downsampling of data. Monitoring online sensors, incident analysis or long-term trend analysis are examples of common query scenarios with different data freshness requirements and ranging very different time spans (Karlstetter et al. [2022]). 2.2.1 Cloud-first Time Series Databases The need for these data storage characteristics gave way to the rise in popularity of TSDBMSs. The simultaneous emergence of both IoT and Cloud Computing resulted in the majority of TSDBMS solutions being meant for the Cloud layer, storing all the generated data. These are currently the most popular solutions on the market, with the top 5 TSDBMSs, according to DB-Engines Rankings (dbe [2022a]), all being cloud-first databases. InfluxDB (InfluxData [2013]) is currently ranked as the number one most popular TSDBMS and features a storage engine that uses a Write Ahead Log (WAL), cache, a Time-Structured Merge Tree (TSM), as well as a time series index to ensure no data loss, a high ingestion rate and efficient data retrieval. To improve efficiency, the TSM only stores deltas between values in a series, while the time series index ensures queries remain fast as data cardinality grows. 7
Another cloud TSDBMS solution, BTrDB Andersen and Culler [2016], achieves relatively fixed response times for analytical queries, regardless of the underlying data size, by using a time-partitioning copy-on-write version-annotated k-ary tree. Its base data points are stored in the leaves, with the depth of the tree defining the interval between data points. Since each internal node holds summaries of the subtrees below it, consisting of statistical aggregates that are computed as nodes are updated, the tree only needs to be traversed to the depth corresponding to the desired resolution for the retrieval of these statistical records. A warm query throughput test (with a preheated cache), which simulates typical tailinganalytics workloads, reached 119 million readings per second. BTrDB also supports out-of-order insertions, which tests have shown to not significantly affect the measured throughput of 53 million inserted values per second, using a four-node cluster. 2.2.2 Multi-layer Time Series Databases Storage and more complex data analysis have mostly been reserved for the Cloud layer, with Edge devices being, up to this point, mainly used for data processing, such as compression, filtering and aggregation, and no data being stored at this layer. However, as the needs of new IoT applications have grown, a demand for increased privacy-sensitivity, lower latencies and enhanced ability to mask connection losses that is not possible to achieve with thus far existing architectures has led to a growth in the need for storage at the Edge layer (Yousefpour et al. [2019]). As a result, multi-layer TSDBMSs, supporting both Edge and Cloud deployment, have also become more common. There are 4 main reasons to store data either completely or partially at the edge level (Satyanarayanan [2017]): •Latency: since Edge nodes are located closer to the users, there is a reduction in latency when compared to Cloud databases, allowing them to be able to fulfil the ever-increasing latency requirements of newer applications. •Scalability: by storing data and applying analytics at the Edge level, the amount of data sent to the Cloud can be reduced, allowing for a reduction of its bandwidth usage, and consequently a potential increase in scalability. •Privacy: the use of Edge nodes allows for the storing of sensitive data in a more geographically restricted area, when compared to Cloud storage. Additionally, the processing of data at the Edge enables the denaturing of sensitive data, such as the blurring of faces or coarse aggregation of sensor readings, before it is sent to the Cloud. 8
TSDBMSs (Andersen and Culler [2016], Paparrizos et al. [2021], Wang et al. [2020]). For the analysed benchmarks to be considered to have realistic workloads, at most one of the listed desired features can be missing. None of the analysed benchmarks support the use of workload traces, possibly due to these not being publicly available. Only IoTDB-Benchmark Liu and Yuan [2019] presents all the desired features for arealistic synthetic workload. Benchmarks that were only missing one of the listed features were also considered to be realistic, even while being inferior to IoTDB-Benchmark. These were SmartBench Kamienski et al. [2019], TSM-Bench Khelifati et al. [2023], and TSBS Timescale [2022], the latter of which was missing a mixed workload, while the first two were missing out-of-order insertions. The remaining benchmarks presented unrealistic synthetic workloads. SciTS Mostafa et al. [2022], PSoTSDBs Shah et al. [2022] and TS-Benchmark Hao et al. [2021] were missing mixed workloads. Out-of-order insertion was the most absent feature with only the aforementioned IoTDB-Benchmark and TSBS supporting it. The other considered ingestion feature, batch ingestion, was missing from PSoTSDBs, IoTDataBench Zhu et al. [2021], TSDBBench Bader [2016] and TPCx-IoT [TPC]. PSoTSDBs and TPCx-IoT were also missing downsampling and out-of-range or filtering queries from their workloads, with TSDBBench also missing the latter. The operation types supported by the supported workloads were: ingestion,querying, or a mix, in which ingestion and querying operations are used concurrently. Both data loading and appending workloads were categorised as ingestion workloads. The vast majority of the surveyed benchmarks support workloads consisting of either only ingestion or only querying operations. TPCx-IoT [TPC] and IoTDataBench were the only benchmarks missing these pre-defined workloads, opting only for mixed workloads. IoTDataBench, SmartBench, IoTDBBenchmark, TSDBBench, TPCx-IoT, and TSM-Bench all support mixed workloads, concurrently performing ingestion and querying operations, and achieving a more demanding workload. The client rate category analyses the way workloads scale by varying each client’s throughput rate. This can fall into one of three categories: maximum, meaning the workload tries to increase the client’s throughput rate until the database is stressed; static, meaning the workload keeps a constant throughput rate for each client; and dynamic, in which the workload tries to vary the throughput rate of each client between a minimum and maximum value. Most surveyed benchmarks maximise the client’s throughput. IoTDataBench, TSDBBench and TPCxIoT use a static client rate since they are all built on top of YCSB Cooper et al. [2010] which throttles the operation throughput, while TSBS supports both a maximum throughput mode and static throughput. 15
TSM-Bench maximises throughput when running its insertion workload and uses a static throughput rate when running its mixed workload. Moreover, none of the benchmarks use a dynamic client rate strategy. 2.4.2 Data Table 2: Data characteristics presented by surveyed benchmarks Origin Datasets Real Realistic Unrealistic Multiple Scalable SciTS Mostafa et al. [2022] PSoTSDBs Shah et al. [2022] TS-Benchmark Hao et al. [2021] IoTDataBench Zhu et al. [2021]**** SmartBench Gupta et al. [2020] IoTDB-Benchmark Liu and Yuan [2019] TSDBBench Bader [2016] TSBS Timescale [2022] TPCx-IoT [TPC] TSM-Bench Khelifati et al. [2023] * supports user-provided Table 2presents the observed data characteristics sub-categorised into the data’s origin and data scaling options For the origin of the datasets used for benchmarking, three options were observed. Real data, meaning data extracted from real-world scenarios; synthetic data, consisting of either purely random data, or random data generated by using probabilistic distributions which can be more realistic but were still categorised as purely synthetic; and data based on real data, which can be generated by a tool that receives real data and uses it to generate data that approximates the input. Only two of the analysed solutions, PSoTSDBs Shah et al. [2022] and IoTDataBench Zhu et al. [2021] support the use of real datasets, with the former featuring them for testing, while the latter supports user-provided real datasets. TS-Benchmark Hao et al. [2021], TSM-Bench Khelifati et al. [2023], and SmartBench Kamienski et al. [2019] both provide realistic synthetic data by using data generation tools that were trained 16
using real data. The first two make use of trained Generative Adversarial Network (GAN) models, while the latter uses temporal, speed, device, and user scaling to generate realistic data. In addition, IoTDataBench also supports the training of a data generation model using a user-provided dataset, allowing for customisable realistic synthetic data generation. Overall, the vast majority of the surveyed benchmarks allow using purely synthetic data, with five out of the ten surveyed solutions supporting only the use of synthetic data, and supporting no other data sources. The scaling capabilities of the analysed data solutions were evaluated by considering whether multiple pre-made datasets were provided or if the data generation tools allowed for the volume of generated data to be scaled. PSoTSDBs, TS-Benchmark and TSM-Bench provide multiple dataset volume options for the scaling of tests, while the remaining provide a way of scaling the data for testing, usually by having the user be able to decide the volume of data to be generated. TSM-Bench provides multiple datasets with different volumes as well as providing a data-generation tool, allowing for scaling of the generated dataset. 17
2.4.3 Metrics Table 3: Metrics presented by surveyed benchmarks Performance Resource monitoring Throughput Latency Sampling Target Readings Continuous Average Inserts Queries Average Percentiles Min/Max Standard Deviation CPU Memory Storage (I/O) Network (I/O) Storage Quota SciTS Mostafa et al. [2022] PSoTSDBs Shah et al. [2022]** TS-Benchmark Hao et al. [2021]**** IoTDataBench Zhu et al. [2021] SmartBench Gupta et al. [2020]** IoTDB-Benchmark Liu and Yuan [2019] TSDBBench Bader [2016] TSBS Timescale [2022] TPCx-IoT [TPC] TSM-Bench Khelifati et al. [2023] Table 3organises the observed performance metrics in throughput and latency sub-categories. Throughput metrics represent a very important aspect of a TSDBMS’s performance since these solutions are subjected to a high volume of data ingestion. Overall, the analysed benchmarks mostly measure the throughput of inserts, with TSDBBench Bader [2016] being an exception as it does not measure throughput of any kind, and TPCx-IoT [TPC] and IoTDataBench Zhu et al. [2021] being the remaining exceptions by measuring throughput for all operations. Both data loading measurements and data appending measurements were considered to be inserts. Apart from PSoTSDBs Shah et al. [2022], which presents the time taken to load a certain amount of data in seconds, all other benchmarks present throughput in operations per second (also mentioned as points or rows per second), with TS-Benchmark Hao et al. [2021] also producing throughput for loading performance in data volume per second. Finally, the throughput values produced by these benchmarks are either presented or logged in a continuous way throughout its execution or simply calculated as an average result and presented at the 18
end of the execution. Benchmarks such as TPCx-IoT, continuously present a rolling overall average, while others, such as SciTS Mostafa et al. [2022], continuously present the calculated average observed over regular intervals; both of these were considered to be continuous. Two sub-categories were considered for latency metrics: the type of operation for which latency is measured, as well as what type of readings are presented. Latency metrics can be measured for both insert and query operations, with queries being considered any data fetching operation ranging from simple reads to aggregation and ranged queries, which are very common operations for the analysis of IoT sensor data. Both TPCx-IoT and IoTDataBench Zhu et al. [2021] do not present latency metrics, being the only analysed benchmarks to not do so, while SciTS and IoTDB-Benchmark Liu and Yuan [2019] were the only surveyed benchmarks to present latency metrics for inserts. SciTS used insert latency for different batch sizes to study how these affected performance. All benchmarks that measured latency did so for queries, either displaying latency metrics for individual queries or all queries. While benchmarks such as PSoTSDBs, TS-Benchmark, SmartBench, and TSM-Bench Khelifati et al. [2023] present only average values for latency, most surveyed benchmarks present more useful statistics, which allow for more in-depth analysis of the performance of the SUT. These included several percentile values, minimum and maximum values and standard deviation. All these statistics are produced by SciTS and TSBS Timescale [2022], with IoTDB-Benchmark and TSDBBench presenting all but the standard deviation values. Table 3also presents the resource monitoring metrics measured by each benchmark. Resource monitoring was considered as it allows for a better understanding of a SUT’s behaviour and performance. Most surveyed benchmarks include little to no resource monitoring, with PSoTSDBs Shah et al. [2022] and SmartBench Kamienski et al. [2019] including some monitoring only during testing, without it being integrated into the benchmark. TS-Benchmark Hao et al. [2021] also uses CPU and Memory monitoring in testing, and storage monitoring integrated into the benchmark. While both SciTS Mostafa et al. [2022] and IoTDB-Benchmark Liu and Yuan [2019] feature most considered resource monitoring (with SciTS missing storage quota monitoring), SciTS does so by integrating an external tool (Glances), while IoTDB-Benchmark implements its system monitoring. Storage quota monitoring usually consists of simply measuring the final storage used by the database, and is mostly used to evaluate compression, by comparing the measured value to the expected storage usage by considering how much data was inserted. 19
2.4.4 Orchestration Table 4: Orchestration presented by surveyed benchmarks Multi-client Local Remote Database node scaling* SciTS Mostafa et al. [2022] PSoTSDBs Shah et al. [2022] TS-Benchmark Hao et al. [2021] IoTDataBench Zhu et al. [2021] SmartBench Gupta et al. [2020] IoTDB-Benchmark Liu and Yuan [2019] TSDBBench Bader [2016] TSBS Timescale [2022] TPCx-IoT [TPC] TSM-Bench Khelifati et al. [2023] Automatic Manual * Single-layer only Table 4details whether the surveyed benchmarks support orchestration, i.e., the automatic deployment of multiple clients and database nodes. Support for the deployment of multiple clients was categorised in two different ways. First, whether the analysed benchmarks support multiple local clients, multiple remote clients or do not support multiple clients at all. Second, the benchmarks that support multiple clients were categorised by whether the deployment of additional clients is done automatically or manually through a set number of clients. Ideally, a benchmark should support the automatic deployment and usage of multiple remote clients, for a more realistic simulation of a distributed workload as well as ease of use. SciTS Mostafa et al. [2022] supports automatically testing over a set of numbers of clients defined by the user through a parameter, and such was considered to be automatic, while TS-Benchmark Hao et al. [2021] automatically increases the number of threads in its continuous injection workload to simulate a growing workload. In addition, IoTDB-Benchmark Liu and Yuan [2019], TSDBBench Bader [2016] and TSBS Timescale [2022] support multiple local clients with the use of parameters. IotDataBench Zhu et al. [2021] and TPCx-IoT [TPC] are the only surveyed systems that support usage of remote clients, 20
with the former doing so automatically to increase the load on the SUT, and the latter allowing the user to define which remote machines to use for a certain benchmark run. Analysed benchmarks only feature single-layer database node scaling, since none of them consider multi-layer TSDBMSs. IoTDataBench and TSDBBench are the only benchmarks that feature database node scaling. IoTDataBench uses database commands to increase the number of database nodes while simultaneously adding corresponding clients for load generation during execution of its scalability test. However, it does not seem to take into account any form of connection between database nodes. On the other hand, TSDBBench uses exploits Python and Vagrant (HashiCorp [2010]) scripts in its Elastic Infrastructure to enable the automatic setup and deployment of database nodes in virtual machines in a range of both public and private cloud solutions. 2.4.5 Discussion Some of the analysed benchmarks were capable of producing realistic workloads, without which it is impossible to accurately assess a SUT’s performance when deployed in a real system. Through the use of concurrent insert and querying operations, multiple solutions achieve realistic workloads that simulate demanding real-world scenarios in which data is being inserted while historical queries are being performed. Several of the existing solutions provide either real or realistic data, essential for the proper evaluation of the targeted databases, since these systems may feature mechanisms that make use of patterns present in real time series data, such as InfluxDB InfluxData [2013] which stores deltas between sequential values. By using unrealistic data, a benchmark may not be able to accurately access the impact of these mechanisms in a particular system’s performance. In addition, ways to scale the volume of data were also present. These are needed for the benchmarking tools to stay relevant as either the SUTs or the hardware they run on evolves to be capable of handling more data. For throughput metrics, benchmarks should ideally continuously track both operation and volume rates for inserts, with one achieving that and several other existing solutions only missing the volume rates, meaning there are existing solutions that are satisfactory in this regard. By providing continuous metrics, benchmarks enable the analysis of a system’s continuous behaviour, such as how the insertion throughput or query latency values vary throughout a test’s execution, e.g., identifying a severe performance degradation halfway through a test, and which would not be obvious by looking at a single value. The analysed benchmarks focus mostly on having latency metrics for querying operations, and several solutions present statistical latency results, which can be very useful to assess query performance. Often times, mean latency values do not tell the whole story about a system’s performance. Two systems 21
with identical mean latencies may in practice perform very differently if one has significantly higher 90th percentile latencies, which can affect certain applications’ performance. Resource monitoring metrics can help understand a SUT’s behaviour but were mostly neglected by existing benchmarking solutions, with only a couple featuring extensive monitoring. This kind of metric can be used to identify system bottlenecks, which can allow the benchmark users to build more efficient deployments, such as in an Edge-Cloud deployment in a situation where the edge nodes are the system’s bottleneck and in which case adding a cloud node would not improve its overall performance. None of the surveyed benchmarks present metrics at the node level, measuring only the database’s overall performance. This makes the existing solutions suboptimal for the assessment of the multi-layer DBMSs as their query performance can vary depending on not only in which node but also in which layer the necessary data is stored. Orchestration features are not commonly found, with none of the surveyed benchmarks supporting multi-layer database node scaling, which is essential for the proper evaluation of Edge-cloud sink node TSDBMSs as it allows the assessment of how the SUT’s performance varies depending on how many nodes are present on each layer or how they are configured. In addition, the automatic deployment of both multiple benchmark clients and database nodes can greatly increase a benchmark’s ease of use by not only allowing the user to easily deploy a database cluster several nodes but also to more easily scale the benchmark’s generated load. In summary, while there are good existing solutions for benchmarking TSDBMSs, all of them lack features that are essential for the evaluation of multi-layer Edge systems. 22
Chapter 3 MulletBench This chapter presents an overview of the MulletBench benchmarking tool, as well as its architecture, features and general characteristics, including its usage and configurations. MulletBench operates through two distinct phases: deployment and execution. Its two components, the Orchestrator and the Client, are responsible for the entire process. During deployment, the Orchestrator configures database nodes, simulates resource and network limitations, and deploys Clients, enabling the testing and simulation of a multi-layer TSDBMS deployment. In the execution phase, Clients are triggered to run in sequential stages, logging all operations and resource usage, and the Orchestrator aggregates all results. These stages can be configured to accommodate various workloads, such as including prepopulation stages, or having several consecutive stages of varying load. By providing aggregated metrics per node, layer, Client or stage, as well as overall results, MulletBench allows for a fine-grained analysis of the SUT’s performance. 3.1 Design Principles MulletBench consists of two main components: the Orchestrator which is responsible for the installation of all tools necessary for the entire benchmarking process, running Client instances, orchestrating the Clients’ execution and finally retrieving and aggregating all generated metrics from each Client; and the Client, responsible for submitting the specified load to the database and registering all submitted operations. Each Client can perform either insertion or query operations, and multiple instances of each can be run simultaneously to generate mixed workloads. Overall, the benchmark features the following characteristics: Workload Workloads can be defined by their type, the operations composing these and what rate is used to submit these operations. MulletBench includes: 23
• Synthetic workloads: Through the composition of both insert and query Clients, different kinds of workloads can be defined. The following are examples of workloads that were defined and used our experiments. –Ingestion workload: Simple insertion of a set amount of data into an empty database. This workload allows for a synthetic analysis of a DBMS’s ability to handle large amounts of data, an important aspect for TSDBMSs; –Query workload: Performing a limited rate of queries over data previously loaded into the database. Measures query latency as well as query throughput, and allows for the measuring of a database’s isolated query performance; –Pre-populated mixed workload: This workload constitutes a realistic workload as per the parameters defined in subsection 2.4.1. An initial stage of insertion is performed, followed by the insertion of a set amount of data while simultaneously performing a limited rate of queries. Measures both insert and query latency and throughput, including for the initial insertion stage. This subjects the SUT to a realistic workload, allowing also for the analysis of how a database performs when it contains higher volumes of data. • Operations: To compose the desired workloads, Clients should be configured to perform either insertion or query operations. –Inserts: perform the insertion of data using batches with configurable size. This allows for the evaluation of a database’s ability to ingest data, including how it performs with different batch sizes; –Queries: perform time-ranged aggregations, downsampling, and filtering queries, which are the most common query types for historical analysis. These can be used to create a query workload that can be representative of the real usage. • Client Rate: The rate at which operations are submitted can be explicitly defined for both insertions and queries, and remains constant for the whole execution of a Client instance. A single test can be defined to have an increased rate of operations through the use of multiple stages. This allows for the simulation of a scenario where a constant amount of analysis is done on an increasing amount of data. Data The characteristics of the data used in a benchmark can have a significant impact on the results. Using a real dataset allows MulletBench to benefit from the advantages of using real data, containing real 24
x=σz +µ(3.1) This value is used as a minimum value for a filter query with the purpose of finding maximum values. A value for finding minimum values can easily be calculated by using negative zin the previous equation. For every successfully performed query, the Client registers both the timestamp in which the query was submitted, the timestamp in which the result was received, the size of the received result (how many entries it contains) and the type of query performed, allowing for query-type specific results. 3.3 Implementation 3.3.1 Orchestrator The deployment component of the Orchestrator was implemented through the use of Ansible1, an opensource tool that simplifies all kinds of automation and deployment. Through the implementation of an Ansible playbook, composed of several plays specifying which tasks to perform, a single configuration file (commonly called inventory file) allows for the definition of all hosts and variables needed for the execution of a wide range of tests. This allows for the installation of all tools needed on all machines, along with the running of any program needed for the system. Docker2and Docker Compose are both installed by this component on each machine as both the execution component of the Orchestrator and the Clients are run in Docker containers, with each of these components having their own image. Additionally, supported databases should also run in Docker containers. For the automatic installation of database nodes as well as setting up replication, database-specific plays were implemented. Support for other databases can be added by implementing additional plays for the following: starting and stopping the instance, starting and stopping replication, starting downsampling (if supported). This component is also responsible for setting limited resources on database instance containers by altering the used Docker Compose files through templating. Docker limits disk IO usage through the use of blkio, when using the blkio_config option in a Compose file. In practice, the container’s use of the storage device is limited, which does not necessarily exactly match the application’s use of disk IO, due to the usage of buffered IO. CPU and RAM usage are limited in the Compose files through the deploy-resources1Ansible: https://www.ansible.com/ 2Docker: https://www.docker.com/ 31
Table 5: Raspberry Pi 3 B+ resource limiting profile Variable Value limited_resources_cpu 4 limited_resources_mem 1024M limited_resources_read_bps 22M limited_resources_write_bps 16M limited_resources_read_iops 2700 limited_resources_write_iops 1200 Table 6: Raspberry Pi 4 resource limiting profile Variable Value limited_resources_cpu 4 limited_resources_mem 4096M limited_resources_read_bps 44M limited_resources_write_bps 40M limited_resources_read_iops 2700 limited_resources_write_iops 1200 limits keys cpus and memory , respectively. These resource limits can be defined with each individual parameter, but some pre-defined profiles that emulate common edge devices such as the Raspberry Pi 4 have already been defined and can be used directly. The CPU and memory values were based on each device’s specifications and the IO limits were based on benchmarks performed on SD cards found on PiBenchmarks3. The resource limiting profiles for the Raspberry Pi 3 B+ and Raspberry Pi 4 can be found in Table 5and 6, respectively. The Orchestrator also achieves network limitations through the use of tcconfig4, an open-source wrapper for tc commands. Performance Co-Pilot (PCP)5is also installed and run if monitoring is desired. The execution component of the Orchestrator component was implemented using Java, with Maven for build automation. This component is provided all relevant information about the test configuration through a YAML file generated in the deployment phase. Since the orchestrator needs to perform resource monitoring of both database and Client nodes through PCP, it runs multiple pmrep processes and, when a test’s execution is completed, kills them. 3PiBenchmarks: https://pibenchmarks.com/ 4tcconfig: https://github.com/thombashi/tcconfig 5Performance Co-Pilot: https://pcp.io/ 32
3.3.2 Client data Client DatabaseConnector Stats Collector SUT Query Generator Database Driver Orchestrator Workers Insertion Worker Query worker config Time Controller Query Builder Dataset Reader Dataset Options Option specific Database specific Dataset specific Implementation Data flow Control flow Figure 5: Benchmark Client structure The Client component was also implemented using Java. Diagram 5presents the main classes used in the Client component, with options, database, and dataset specific classes being omitted for simplicity, showing only respective interfaces or abstract classes. Since Clients should be abstracted from which database they are evaluating, database connectors must be implemented for each supported database, using a common interface that can be used by the workers. Any database-specific classes, such as Measurement classes for InfluxDB or Schema classes for Apache IoTBD, should be implemented and used exclusively by that database connector implementation. These will need to implement each dataset’s parsing, so that the original string data can be correctly parsed into each database’s preferred representation. Dataset classes are alternative ways data can be stored and used, implementing the Dataset interface to standardize access to any dataset representation’s data. In cases where the dataset used for an insertion workload can fit entirely into memory, it can be loaded into memory only once and subsequently accessed by all workers simultaneously using the SharedDataset . Otherwise, IteratorDataset is used, avoiding loading the whole dataset into memory, but resulting in higher runtime overhead since it 33
the dataset needs to be read from disk during execution. On the other hand, DatasetReader classes implement parsing for each specific dataset to be used as well as provide a standardized way for workers to access data from a used dataset, and make use of Dataset classes to access that data in any representation. An implementation of Dataset was created for a Passive Vehicular Sensors Dataset, in particular the gps_mpu dataset (Menegazzo and von Wangenheim [2020]). This dataset contains inertial sensor data combined with GPS data in a comma-separated values (CSV) format, containing 32 columns including the timestamp and 144,036 lines or lines, totalling 4,609,152 values, and with a frequency of 100 samples per second. QueryGenerator classes also need to be implemented for each supported dataset since these need to know the name of each column and how to process the values in the dataset. Any dataset containing only float values can simply extend FloatsQueryGenerator and needs only to implement the process method, meaning the way queries are generated can be consistent for different datasets. Each supported database should also implement a QueryBuilder , allowing the QueryGenerators to be abstracted from which database they are generating queries for. The TimeController classes allow for there to be multiple timestamp options. CurrentTimeController allows for datasets to me manipulated by injecting the current time, while DatasetTimeController allows for datasets to be inserted with their original timestamp, and for their original time sequence to be extended when the dataset is being looped over. A mechanism for rate limiting is used in cases where a lot of submitted operations are pending a response. Since operations are submitted through a thread pool, this mechanism does not impact how many operations are actually submitted to the database since that is already limited by the size of the pool, instead limiting only how many operations are in memory ready to be submitted. This avoids scenarios were very slow SUTs result in increasingly high memory usage on the Client through an increasing number of pending operations. 3.3.3 Configurations For the configuration and executions of tests, the user needs to specify an Ansible inventory file, containing all the necessary information for the deployment of the database nodes and the benchmark Clients, as well as for the execution of each Client’s workload. All the available options for this configuration file are listed in A.1. Since the execution Orchestrator and the Clients can also be run manually, as well as for the understanding of each component’s specific options, the available options for the execution of the 34
execution Orchestrator and the Clients are listed in A.2 and A.3, respectively. Usage To specify a test configuration, the user needs only to configure the aforementioned Ansible inventory file, specifying hosts and their specific variables, both for the database nodes and the benchmark Clients. Several examples of pre-defined tests have already been specified, and these can also be used as a guide for the user to more easily understand how this file should be organized. To perform a test run, the user should run the Ansible playbook, specifying the inventory file with the configured hosts and variables. After a run has been performed, its logs are stored on the /tmp/ directory in the machine where the Orchestrator was run. These logs contain the benchmark execution logs, as well as the monitoring logs for each Client and database node. These can be used by the included plot generation script to generate graphs of the results. There is also a provided script for the automatic execution of multiple runs and generation of result graphs, with optional arguments which can be used to specify the number of runs, number of the starting run, whether to skip cleaning or whether to only perform cleaning. 3.4 Summary MulletBench is a Time Series Database Management System (TSDBMS) benchmark specifically designed to evaluate how these systems perform when deployed in hybrid Edge-Cloud deployments and how their replication mechanisms perform. By featuring user-configurable workloads, MulletBench can be used to generate a wide range of realistic scenarios. Its deployment features allow for not only in ease-of-use but also consistent results. In addition, the orchestration of several Client instances, allows for loads to be generated in a distributed way, more accurately representing real scenarios while simultaneously enabling the generation of more intensive loads by leveraging multiple machines. While MulletBench’s per-node, per-layer and per-stage metrics give relevant values for comparison and assessment of a system’s performance, its continuous logging of performance and resource monitoring metrics enable a detailed analysis of a system’s continuous performance throughout the whole test, more easily identifying problems or bottlenecks that may occur in any node present in the system under test (SUT). 35
Chapter 4 Evaluation and Results This chapter presents results from experiments performed using the specified and developed benchmark, including its 2 main components: the Orchestrator, and Client. Testing scripts were used to help automate the testing as well as generating all figures presented in this chapter, except for bar charts. These tests allow for the validation of some of the benchmark’s mechanisms, such as its limitation features and observing the benchmarks general functioning. Then, the results from the tests are used to not only compare the two DBMSs, InfluxDB and Apache IoTDB, but also to compare different deployments and data retention strategies using downsampling. Using the benchmark’s capabilities allows for the detailed analysis of performance metrics and resources during the test’s execution. This chapter aims to answer the following questions: • Can MulletBench limit the rate at which its Clients submit a workload? • Can MulletBench limit the resources of the edge database nodes? • Can MulletBench limit the network conditions of its testing environment? • How does the performance of InfluxDB and Apache IoTDB compare in the same deployment scenario for each workload type? • How does the performance of each database vary in different deployment scenarios for each workload type? Section 4.2 presents the results from the tests performed to validate the benchmark’s limiting features, including the Client’s rate limiting, the limiting of edge node resources and the limiting of network conditions. Section 4.3 presents the results from the tests performed in a cloud scenario, Section 4.4 presents the results from the tests performed in an edge-cloud hybric scenario with 2 edge nodes, and Section 4.5 presents the results from the tests performed in an edge-cloud hybric scenario with 4 edge 37
nodes. Throughout these sections, comparisons are made between both databases and between the different deployment scenarios. Finally, Section 4.6 presents the main points that can be drawn from the presented results. 4.1 Methodology Each test was run 3 times as a way to ensure both the consistency of the benchmark and the conclusions taken from analysing the results. Between each run, all database nodes were reset, removing all directories generated by the database’s execution as well as removing the used Docker container instances. For all tests, the Clients’ resources were also monitored to ensure that they were never the limiting factor. 4.1.1 Experimental Setup A total of 8 machines were used for testing, 5 of which were used for all Clients and the Orchestrator, sharing the same specifications. The remaining 3 machines were used for the database nodes, also sharing the same specifications between them. The specifications and software versions for the Client and Orchestrator machines are as listed in Table 7, while those of the database machines are listed in Table 8. All Client and database machines are connected through a switched Gigabit network at a max distance of 1 hop. 4.1.2 Workloads The testing considers 3 main workload categories: insertion or ingestion, query and mixed workloads. Insertion The insertion workload aims to evaluate a database’s pure insertion performance using batch insertions, enabling the analysis of a database’s behaviour when submitted to a controlled insertion load. All insertion tests were configured with multiple Client instances and each with multiple workers. In each test, each worker inserts a batch of 100 lines every second. The scaling of this test is achieved by increasing the number of workers and Client instances. The target runtime length for each test was chosen to be 16 minutes, which means each worker inserts 96,000 lines, at the aforementioned rate. For this workload, the main metric is the number of lines per second successfully inserted by all Client instances, corresponding 38
Table 7: Specifications of database machines Item Description CPU Intel Core i5-9500 RAM 16 GB DDR3 2666MT/s Storage 240 GB NVME SSD Operating System Ubuntu 20.04 Docker version 24.0.4 Docker Compose version v2.11.1 IoTDB version 1.1.1 InfluxDB version 2.7 Table 8: Specifications of Client and Orchestrator machines Item Description CPU Intel Core i3-4170 RAM 16 GB DDR3 1600MT/s Storage 120 GB SATA SSD Operating System Ubuntu 20.04 Docker version 24.0.4 Docker Compose version v2.11.1 Ansible version 2.15.1 Python version 3.10.12 to how many lines from the selected dataset are being inserted, regardless of how many columns each line contains. Query The configured query workloads aim to evaluate a database’s query performance by submitting common IoT queries such as outlier value filters, aggregations and downsampling queries, which are generated based on an analysis of the used dataset. These enable the analysis of a database’s ability to handle each type of query individually, without having the resulting metrics affected by its ingestion mechanism. For this workload, the main metrics are the latencies for each query type. In this workload, the dataset is loaded with its original timestamps due to the benchmark’s mechanism for knowing the time range for queries. The used query distribution is 40% for outlier filters, 40% for aggregations and 20% for downsampling queries, representing a plausible distribution, in which downsampling queries are less common since they correspond to continuous aggregations. Mixed Finally, mixed workloads combine the previous types, performing the insertion of a dataset while submitting relevant queries for that dataset, to simulate a workload which the database might be subjected to when deployed in the real-world. For this workload, the main metrics are both the number of lines per second successfully inserted by all Client instances and the latencies for each query type. For tests using the replication of downsampled data to a cloud node, using the outlier filter query for that node no longer makes sense, since extreme values have been smoothed out. Therefore, in tests including downsampling, the cloud layer are subjected only to downsampling and aggregation queries. The total amount of queries is the same while keeping the proportion between downsampling and aggregation queries. The distribution 39
of query types is the same as in the query workload, except for the cloud layer when downsampling is used, where the distribution is 66.7% aggregations and 33.3% downsampling queries. 4.1.3 Data All tests used the Passive Vehicular Sensors gps_mpu dataset Menegazzo and von Wangenheim [2020]. This dataset is approximately 81 MB, containing 144,036 lines with 30 values each. As mentioned before in Section 3.2.2, the size of the used dataset is not a limitation since the Clients automatically loop over when the end of the dataset is reached, injecting the correct timestamps. This way, the used dataset can essentially be inserted as many times as necessary, each time with updated timestamps which follow the same pattern as the original dataset. 4.2 Feature Validation This section presents a series of tests performed to demonstrate MulletBench’s limiting capabilities. These tests were performed using a single database node without any resource limiting (except when stated). By validating each limiting feature individually, it is possible to ensure that the benchmark is capable of limiting the resources of a database node, or network connection between node, as well as the rate at which data is submitted to it, something which is necessary for the validity of the tests that follow. 4.2.1 Control Test First, as control, a standard insertion test was conducted, showing what throughput each database can achieve without any limiting. For this test, the target throughput was controlled by varying the number of Client instances as well as the number of workers submitting data in each Client instance. For InfluxDB, this test used 8 Client instances configured as listed in Table 26, resulting in a target throughput of 64,000 lines per second. The average throughput is 41,094.186 lines per second, showing the database was overloaded, and not able to handle the requests at the rate they were submitted. This resulted in the tests taking on average roughly 24.9 minutes. For IoTDB, as initial runs show it could handle higher throughputs, the number of clients was increased, with 8 Client instances configured as listed in Table 27, totalling a target throughput of 360,000 lines per second. The average throughput for this test is 320,640.8 lines per second, a significantly higher value than what InfluxDB was able to achieve, which is also clearly visible by looking at the comparison of both database’s throughput in Figure 6. 40
0 200 400 600 800 1000 1200 1400 1600 Elapsed time (s) 0 2 4 6 8 GB Test finish RAM Figure 13: Memory usage for limited resources test targeting InfluxDB Apache IoTDB 0 250 500 750 1000 1250 1500 1750 Elapsed time (s) 0 20000 40000 60000 80000 100000 Throughput (lines/s) Figure 14: Throughput for limited resources test targeting IoTDB 47
Figure 15: Insertion latency for limited resources test targeting IoTDB For IoTDB, the observed average throughput is 49,677.15 lines per second. Looking solely at this value, one might assume that the database’s performance is normal, with an expected reduction in performance. A closer look reveals all runs have over 100 failed insertion operations and the throughput plot in Figure 14 and latency scatter plot in Figure 15 show rather unusual behaviours. This is most likely due to a severe CPU bottleneck, evidenced by the almost constant 100% CPU usage as seen in Figure 16. Figure 17 shows that memory usage is also a possible limiting factor, staying at its maximum of 8 GB throughout the test. 0 250 500 750 1000 1250 1500 1750 Elapsed time (s) 0 50 100 150 200 250 300 % up to 100 * number of cores Test finish cpu-total Figure 16: CPU usage for limited resources test targeting IoTDB 48
0 250 500 750 1000 1250 1500 1750 Elapsed time (s) 0 2 4 6 8 GB Test finish RAM Figure 17: Memory usage for limited resources test targeting IoTDB These results show that MulletBench is capable of limiting the resources available to a deployed database node, as well as submitting an expressive workload, allowing for a fine-grained analysis of the SUT’s behaviour. 4.3 Cloud Scenario All the tests performed in this section use a cloud scenario configuration, meaning a single database node is used and all generated load is submitted to that node. This allows for the evaluation of each system in a scenario where data is offloaded directly to the cloud, without the intervention of any edge nodes. Tests were run for each of the workloads mentioned in Section 4.1.2. 4.3.1 Insertion Workload This workload’s goal is to evaluate a database’s ability to handle an insertion-only scenario in which data is concurrently inserted from several sources and no queries are performed. This test corresponds to the control test previously presented. For both databases, all workers insert the same amount of data at the same rates. To vary the generated load, the number of workers is varied, simulating different amounts of devices accessing and inserting data. 49
300000 350000 400000 InfluxDB IoTDB 0 200 400 600 800 1000 1200 1400 Elapsed time (s) 0 50000 100000 Throughput (lines/s) Figure 18: Throughput comparison for insertion workload in Cloud-only scenario The throughput comparison in Figure 18 shows that IoTDB achieves a significantly higher throughput than InfluxDB, with an average of 320,640.8 lines per second compared to 41,094.19 lines per second, corresponding 7.8x the throughput. InfluxDB’s test has a longer duration since its average throughput is further from its target throughput. InfluxDB Influx’s throughput plot in Figure 18 shows that the target throughput is matched for a few seconds before dropping to around 40,000 lines per second, remaining there with some variation until the end of the test. Although the presented figure pertains to a single run, the same behaviour was present in all runs. Looking at the resource monitoring graphs, it is not easy to tell which is the limiting factor for the database’s performance. Since the IO usage monitoring in Figure 19 shows high IO throughput, the limit may lie in the number of IOps. 50
0 200 400 600 800 1000 1200 1400 1600 Elapsed time (s) 0 25 50 75 100 125 150 175 Throughput (MB/s) Test finish io-read io-write Figure 19: Disk IO usage for insertion workload targeting InfluxDB in Cloud-only scenario IoTDB IoTDB’s disk IO usage in Figure 20 and memory usage in Figure 21 reveal that the memory limit is quickly reached and that there are constant disk reading operations, even though this is an insertion workload. This indicates the performance limitation may be coming from the limited memory, likely resulting in the use of swap memory. 0 200 400 600 800 1000 1200 Elapsed time (s) 0 50 100 150 200 Throughput (MB/s) Test finish io-read io-write Figure 20: Disk IO usage for insertion workload targeting IoTDB in Cloud-only scenario 51
0 200 400 600 800 1000 1200 Elapsed time (s) 0 2000 4000 6000 8000 10000 12000 14000 MB Test finish RAM Figure 21: Memory usage for insertion workload targeting IoTDB in Cloud-only scenario This workload shows that IoTDB is capable of handling a significantly higher insertion throughput than InfluxDB, in a single node deployment scenario. 4.3.2 Query Workload The query workload aims to evaluate a database’s ability to handle an analytical scenario in which previously inserted data needs to be analysed using common time series queries such as aggregation, downsampling and filtering queries. For both databases, the same amount of data is initially ingested, ensuring the observed performance is dependant solely on the database’s querying performance. This initial ingestion has a volume close to 120 GB assuming the original dataset’s text format. In total, 207.4 million lines are inserted, corresponding to 6,222 million values. For this workload, the database is loaded once and the 3 tests were performed in sequence without reloading the database. During test execution, queries were submitted at the same rate and with the same distribution for both databases. The used test configuration consisted of a single Client with a configuration as listed in Table 28. This configuration results in a query being submitted by the Client every 2.5 seconds and for an estimated length of 16 minutes. With this workload, the analytical performance of a database can be measured by the latency with which it responds to each type of query, meaning a database’s performance for a particular kind of query type can be measured and compared. 52
Filter Aggregation Downsampling 0 2,000 4,000 6,000 8,000 Latency (ms) IoTDB InfluxDB Figure 22: Query latency bar chart for query workload in Cloud-only scenario Figure 22 shows that InfluxDB has higher average latency for filter queries, while being only slightly worse for aggregations and slightly better for downsampling, when compared to IoTDB. Although InfluxDB’s outlier filtering queries take on average more than 10x longer than IoTDB’s (8038 ms compared to 778 ms), looking at the median and 10th percentile values reveals something different, with these being, respectively, only 80% and 28% higher for InfluxDB than IoTDB (1433.99 ms compared to 798.50 ms and 818.33 ms compared to 638.67 ms). This indicates that some queries with very high latency are skewing the average values. Another interesting observation from the results is that, although InfluxDB performed worse in outlier filter and aggregations, results show on average 32% lower latency on downsampling queries. 53
InfluxDB 15000 16000 17000 18000 Query Type AGGREGATION DOWNSAMPLING OUTLIER_FILTER 0 200 400 600 800 1000 Elapsed time (s) 0 1000 2000 3000 Latency (ms) Figure 23: Query scatter for query workload targeting InfluxDB in Cloud-only scenario The query scatter plot in Figure 23 shows how InfluxDB’s outlier filtering queries are clearly divided into two layers by their latency, one of them having significantly higher latencies than any other queries. 54
IoTDB 0 200 400 600 800 1000 Elapsed time (s) 0 500 1000 1500 2000 2500 Query Type AGGREGATION DOWNSAMPLING OUTLIER_FILTER Latency (ms) Figure 24: Query scatter for query workload targeting IoTDB in Cloud-only scenario IoTDB’s query scatter plot is also presented in Figure 24. It shows that all query types consistently have similar latencies, with no clear outliers. Unlike InfluxDB, there is no clear split in outlier filtering query latencies. This workload shows that IoTDB has a significant lower average latency for filter queries and a slight advantage for aggregations, while InfluxDB has slightly lower downsampling query average latency. Results also show that both databases should be able to deal with a higher rate of queries at the cost of higher latencies as they both were capable of responding to the submitted query rate. 4.3.3 Mixed Workload This workload corresponds to a realistic scenario in which queries are being performed in parallel with new data being inserted. The workload also includes an initial population stage, consisting of insertion only, and allowing for the database to already populated when the mixed stage starts and queries start being performed. The same exact configuration was used for both databases, both in volume of data and in rate of insertion and querying. Since the benchmark submits a constant rate of insertions, the database’s performance can be measured by whether it can successfully handle that insertion throughput. At the same 55
time, its performance can also be evaluated by the latency with which it responds to querying operations. Since the benchmark allows the analysis of querying latency per query type, a database may be better in certain types of queries while being worse in others. The preload stage consists of 4 Client instances with a configuration as listed in Table 29. The mixed stage consists of 4 insert Clients and 1 query Clients. The insert Clients in the mixed stage use a configuration as listed in Table 30. The query Client uses a configuration as listed in Table 31. The throughput comparison in Figure 25 shows that both databases are able to handle the constant data throughput during the initial population stage, but InfluxDB is unable to do so during the mixed stage. 0 500 1000 1500 2000 Elapsed time (s) 0 5000 10000 15000 20000 25000 30000 35000 Pre-population Mixed stage InfluxDB IoTDB Throughput (lines/s) Figure 25: Throughput comparison for mixed workload in Cloud-only scenario 56
0 200 400 600 800 1000 1200 1400 Elapsed time (s) 0 5000 10000 15000 20000 25000 30000 Throughput (lines/s) Figure 33: Throughput plot for insertion workload targeting InfluxDB with no downsampling in 2 Edge + Cloud scenario No downsampling Replicating raw data instead of downsampled data, InfluxDB manages to achieve a throughput of 20,129.40 lines per second. Looking at the disk IO usage in Figure 34 and memory usage in Figure 35 allows us to clearly identify the TSM compactions InfluxDB performs during execution, as they coincide with the spikes in latency in Figure 36. These compactions are a mechanism used by InfluxDB (and other databases using Log-Structured Merge-Tree (LSM) tree-based storage engines) to free memory or disk by converting write-optimized files to compressed and read-optimized files in disk. This can be seen through the regular interval spikes in IO writes. These spikes can have different lengths depending on the compaction levels, which are increasingly more demanding, since the higher the compaction level, the more data they need to compact. The relation between compactions in LSM-based stores and throughput latency spikes has already been widely studied (Balmau et al. [2019]). 63
0 200 400 600 800 1000 1200 1400 1600 Elapsed time (s) 0 20 40 60 80 Throughput (MB/s) Test finish io-read io-write Figure 34: Edge node disk IO usage for insertion workload targeting InfluxDB in 2 Edge + Cloud scenario 0 200 400 600 800 1000 1200 1400 1600 Elapsed time (s) 0 1000 2000 3000 4000 MB Test finish RAM Figure 35: Edge node memory usage for insertion workload targeting InfluxDB in 2 Edge + Cloud scenario 64
Figure 36: Latency scatter plot for insertion workload targeting InfluxDB in 2 Edge + Cloud scenario Looking at the network usage for an edge node in Figure 37 and for the cloud node in Figure 38 on the section after the test execution has finished shows there is still activity between the edge and cloud nodes. 0 200 400 600 800 1000 1200 1400 1600 Elapsed time (s) 0 5 10 15 20 25 30 Throughput (MB/s) Test finish eth0-in eth0-out Figure 37: Edge node network usage for insertion workload targeting InfluxDB in 2 Edge + Cloud scenario 65
0 200 400 600 800 1000 1200 1400 1600 Elapsed time (s) 0 1 2 3 4 Throughput (MB/s) Test finish eth0-in eth0-out Figure 38: Cloud node memory usage for insertion workload targeting InfluxDB in 2 Edge + Cloud scenario This is due to InfluxDB using a replication queue, where ingested data is placed to later be replicated. A further analysis of InfluxDB’s logs revealed this queue is filling up during the test and not all data inserted at the edge layer is being replicated to the cloud node. This proves to be an issue as from the client’s side all insertions are confirmed by the database, even though the data has not been replicated as expected, leading to data loss if the edge nodes’ bucket is configured to only retain data for a certain period of time through the setting bucket-retention-period . This behaviour has been reported to the developers through a GitHub issue1. Using downsampling may mitigate the replication queue issue as the volume of replicated data decreases significantly. Even so, there is no guarantee of replication if the volume of downsampled data still exceeds what the replication queue can handle. 2 second downsampling Using downsampling, the same patterns in disk IO usage in Figure 39 and memory usage in Figure 40 can be observed on the edge nodes, indicating the same compaction limitations. However, the observed throughput is significantly higher, reaching 33,921.17 lines per second, a 68.6% increase over the no downsampling scenario, indicating that the replication of high volumes of data has a significant performance impact. 1InfluxDB replication issue: https://github.com/influxdata/influxdb/issues/24434 66
0 200 400 600 800 1000 1200 Elapsed time (s) 0 1000 2000 3000 4000 MB Test finish RAM Figure 39: Edge node memory usage for insertion workload targeting InfluxDB with 2 second downsampling in 2 Edge + Cloud scenario 0 200 400 600 800 1000 1200 Elapsed time (s) 0 20 40 60 80 Throughput (MB/s) Test finish io-read io-write Figure 40: Edge node disk IO usage for insertion workload targeting InfluxDB with 2 second downsampling in 2 Edge + Cloud scenario IoTDB For IoTDB the test configuration contained 4 Clients with configurations as listed in Table 33 and totalling a target throughput of 280,000 lines per second. 67
0 200 400 600 800 1000 Elapsed time (s) 0 20 40 60 80 Throughput (MB/s) Test finish io-read io-write Figure 41: Edge node disk IO usage for insertion workload targeting IoTDB in 2 Edge + Cloud scenario 0 200 400 600 800 1000 Elapsed time (s) 0 1000 2000 3000 4000 MB Test finish RAM Figure 42: Edge node memory usage for insertion workload targeting IoTDB in 2 Edge + Cloud scenario This configuration seems to be limited by the edge nodes’ disk IO and memory limits, as shown in Figures 41 and 42, respectively, with both resources being fully utilized for the duration of the test. As such, a deployment with more edge nodes should be able to yield improved results. 68
0 200 400 600 800 1000 Elapsed time (s) 0 10 20 30 40 50 Throughput (MB/s) Test finish eth0-in eth0-out Figure 43: Edge node network usage for insertion workload targeting IoTDB in 2 Edge + Cloud scenario 0 200 400 600 800 1000 Elapsed time (s) 0 20 40 60 80 Throughput (MB/s) Test finish eth0-in eth0-out Figure 44: Cloud node network usage for insertion workload targeting IoTDB in 2 Edge + Cloud scenario Network usage in an edge node and the cloud node, shown in Figure 43 and 44 respectively, shows that after the test execution is finished there is still activity between the edge and cloud nodes. This indicates that replication is still ongoing. Looking at the data present in the cloud node reveals that not all data is successfully replicated, even though all insertion operations are confirmed by the database. This results in some data loss at the cloud node. When compared to the cloud scenario, this deployment does not prove to be beneficial for insertion 69
performance, resulting in a 17.5% and a 12.0% drop in average throughput for InfluxDB and IoTDB, respectively. However, this workload shows how downsampling can be used to increase throughput by reducing the amount of replicated data, with InfluxDB showing a 53% increase in throughput when using 2 second downsampling. 4.4.2 Query Workload For the query workload, the database nodes are initially loaded with data, over which queries performed. The amount of data used for pre-population is the same between the tested databases, as to allow comparisons between them, but also the same as in the previously conducted Cloud scenario tests. This achieves a scenario with the same amount of sensors and queries but a different deployment strategy, enabling their comparison. To enable both of these comparison vectors, the queries performed and their rate of submission needs to be the same, so that the load induced is the same and the resulting query latencies can be compared. However, since this scenario contains heterogeneous nodes, this does not translate directly to the same amount of queries. As each edge node contains half the ingested data, each query performed in an edge node only comprises half as much data as if it is performed in the cloud node. To generate an equivalent query workload, the query load is split equally between the edge and cloud layers, with all nodes on both layers being submitted the same amount of queries. The configuration used by each Client is listed in Table 34. The end result is that 50% more queries are performed when compared to the Cloud scenario, but of which 66% only comprise half the data inserted, essentially distributing the same load over more nodes. 70
Filter Aggregation Downsampling 0 1,000 2,000 3,000 4,000 Latency (m s) Cloud layer Filter Aggregation Downsampling 0 1,000 2,000 3,000 4,000 5,000 Latency (ms) Edge layer IoTDB InfluxDB Figure 45: Query latency for query workload in 2 Edge + Cloud scenario Looking at the edge layer average latencies in Figure 45, InfluxDB’s average filter latency is 115.5% higher than IoTDB’s, while its aggregation and downsampling latencies are 42.4% and 30.7% lower, respectively. InfluxDB Since in this workload the dataset is initially loaded with its original timestamps, testing with downsampling is not possible as the downsampling tasks are applied to data which have timestamps in the past x units of time, and not data which were inserted in the past x units of time. As such, the only possible test was performed without downsampling. Due to the issue with the replication queue previously identified in Section 4.4.1, the initial data population does not result in all the data reaching the cloud node. This can also be confirmed by looking at the average volume returned by each filter query at the cloud layer, which dropped 48.8% from 18,944.4 entries per query to 9,689.9 entries per query, indicating this database node is not completely populated. The cloud layer’s query latency performance should then not be considered valid for comparison. At the edge layer, the average latency for filter queries decreased 47.9% when compared to the cloudonly deployment, while for aggregation and downsampling queries average latency increased 160.78% and 141.37%, respectively. 71
IoTDB When analysing the metrics corresponding to the cloud layer, even though the average aggregation latency is lower when compared to the cloud scenario, at 165.24 ms versus 234.28 ms, the same cannot be said for the filter and downsampling queries. This is unexpected since in this scenario the cloud layer is subjected to half the load when compared to the cloud scenario. Further investigation on how the population data is stored on all nodes, shows that on the cloud node, all data received from the edge nodes is stored on delayed or unsequenced TsFiles. Delayed TsFiles are used when IoTDB identifies received data as being out of order compared to what is stored on the normal TsFiles, using these secondary files to avoid a more expensive write operation at the cost of more expensive query operations and the need for compactions. This mechanism can be beneficial when a small portion of data is received out-of-order. It is currently unclear as to why data received from the edge nodes is identified by the cloud node as being delayed. This behaviour has been reported to the developers in a GitHub issue2. At the edge layer, query performance is worse for all types when compared to the cloud-only deployment, with 150.4%, 626.8%, and 135.8% increase in outlier value filters, aggregations, and downsampling queries, respectively. This workload shows that InfluxDB achieves lower average latencies than IoTDB for aggregations and downsampling queries, while IoTDB has lower latency for filter operations, when subjected to a query load. 4.4.3 Mixed Workload As with the query workload, the mixed workload is configured such that its results could be used for comparison not only between databases but also between deployment configurations, i.e., between cloud and hybrid edge-cloud deployments. To achieve this, the insertion load is simply being split between both edge nodes, while the query load is split between the edge and cloud nodes, in the same way as in the previously presented query workload. Once again, this results in a 50% increase in overall submitted queries, but with two thirds of them comprising half the data achieving essentially the same overall load. The goal is that the same overall load is distributed over a higher number of nodes. The resulting workload is therefore composed of 4 identical insertion Clients, 2 for each edge node, with configurations as listed in Table 36, and 3 identical query Clients, 1 for each node, configured as listed in Table 37. As mentioned in Section 4.1.2, when using downsampling, outlier filter queries are not be submitted 2IoTDB unsequenced files issue: https://github.com/apache/iotdb/issues/11441 72
replicated successfully, with a loss of up to 1.6% of data in the case of some devices. There is no indication of this behaviour returned to the client, and it can only be observed either by looking at the data present in the cloud node or by looking at the status of the replication pipe on the edge node, which was WARN , indicating ‘data loss or other error has occurred‘ (iot [2023]). Looking at the node’s logs also shows multiple errors. This behaviour has been documented and reported to the developers in a GitHub issue 3. This workload shows significant throughput increases when compared to the Cloud-only scenario, showing how the edge layer can be scaled independently from the cloud layer in these deployments. 4.5.2 Query Workload As with the previous query workloads, data is pre-populated through each of the edge nodes, which replicate it to the cloud node. The volume of data loaded is the same as in the previous tests, as to allow comparisons between them. When compared to the 2 edge deployment, the cloud load is kept the same, while the edge load is now distributed over the 4 edge nodes. All nodes are subjected to the same amount and rate of queries, with Clients configured as listed in Table 41, although for the cloud node this translates to a more demanding load as it contains more data. This test is composed of five Clients, one for each edge node and one for the cloud node, with the following configuration. Filter Aggregation Downsampling 0 2,000 4,000 6,000 Latency (ms) IoTDB cloud layer IoTDB edge layer InfluxDB cloud layer InfluxDB edge layer Figure 53: Query latency bar chart for query workload in 4 Edge + Cloud scenario Figure 53 shows the average latency for each query type at the cloud and edge layers for both 3IoTDB replication issue: https://github.com/apache/iotdb/issues/11428 79
databases. In it can be seen that while IoTDB presents significantly lower average latencies for filter queries on both layers, InfluxDB achieves lower latencies for both aggregation and downsampling queries. InfluxDB Distributing the load during the population phase over a higher number of edge nodes mitigates the identified replication queue issue, resulting in more data reaching the cloud layer. This, in addition to performing the loading phase at a slower rate allowed for all records to be replicated correctly. With this deployment we can observe query latency improvements in both layers when compared to the original cloud-only deployment. At the cloud layer there is an 18.32%, 23.07% and 26.02% decrease in average latency, for outlier filters, aggregations and downsampling respectively, when compared to the cloud scenario, while the edge layer sees a 77.17%, 11.48% and an 18.83% decrease, respectively. IoTDB For IoTDB, only downsampling operations at the edge layer show lower latency than the original Cloud-only scenario, with all other query types for both the edge and cloud layer having lower latencies in the original deployment. When compared to the 2 edge deployment, there are considerable improvements in query latency at the edge layer, with 57.2%, 54.0%, 64.4% reductions for filters, aggregation and downsampling queries, respectively. Just like for the 2 Edge + Cloud scenario in Section 4.4.2, looking at the results for the cloud layer reveals some unexpected values. Once again, although the cloud layer is subjected to the same load, the average latency observed has increased significantly for all query types when compared to the 2 edge deployment. This workload shows how InfluxDB can benefit from additional edge nodes when subjected to a query load, showing lower latencies at both layers than in the original Cloud-only scenario. IoTDB, however, failed to show any significant improvements when compared to the Cloud-only scenario. 4.5.3 Mixed Workload The mixed workload is configured to perform the same throughput and queries as in the mixed workload tests, but now distributing the load over 4 edge nodes. The resulting workload is composed of 4 identical insertion Clients, 1 for each edge node, with configurations as listed in Table 43, and 5 identical query Clients, configured in the same way as in Table 44, 1 for each edge node. 80
Filter Aggregation Downsampling 102 103 Latency (ms) IoTDB Cloud Layer IoTDB Edge Layer InfluxDB Cloud Layer InfluxDB Edge Layer Figure 54: Query latency bar chart for mixed workload in 4 Edge + Cloud scenario 0 500 1000 1500 2000 Elapsed time (s) 0 5000 10000 15000 20000 25000 30000 35000 Pre-population Mixed stage InfluxDB IoTDB Throughput (lines/s) Figure 55: Throughput comparison for mixed workload in 4 Edge + Cloud scenario The throughput comparison in Figure 55 shows that IoTDB is able to match the targeted throughput, achieving an average of 27,902.4 lines per second, while InfluxDB achieves 26,792.61467 lines per second, 4.3% below the target of 28,000. By looking at the query latencies for each database’s optimal configuration, InfluxDB using 2 second downsampling, and IoTDB using no downsampling, in Figure 54, we can see how beneficial having the edge layer downsampling data is, as query latency on the cloud layer is significantly lower for InfluxDB, having 86.5% and 88.6% lower average query latencies for aggregations and downsampling queries, respectively. 81
On the other hand, edge layer query performance favours IoTDB, with 83.7% lower average outlier filter latency, 48.4% lower aggregation latency, but 12.6% higher downsampling latency. InfluxDB This scenario results in significantly improved edge layer query performance when compared to the 2 edge deployment, with 76.8%, 86.1%, and 86.3% lower average latencies for filters, aggregations, and downsampling queries, respectively. When comparing to the original cloud deployment, the edge layer also achieves 86.5%, 67.2%, and 67.9% lower average latencies for filters, aggregations, and downsampling queries, respectively, while the cloud layer lower achieves 96.3%, and 92.1% lower average latencies for aggregations, and downsampling queries, respectively. This deployment therefore shows the benefits not only of performing downsampling at the edge layer, resulting in lower latency at the cloud layer, but also how distributing the load over a higher number of nodes can improve performance at the edge layer. IoTDB At the cloud layer, this deployment yields very similar results to the 2 Edge + Cloud scenario, with 4.1%, 3.8%, and 3.4% higher average latencies for filter queries, aggregations and downsampling queries, respectively. This is expected since this node is subjected to the same load and contains the same data, simply being replicated from more edge nodes. At the edge layer, there is a significant improvement in average query latencies when comparing to the deployment with 2 edge nodes, dropping by 78.9%, 83.4%, and 79.7% for filters, aggregations, and downsampling queries, respectively. This is also expected behaviour, since each edge node now receives and queries over a smaller volume of data. Query latencies for both layers represent significant improvements over the cloud-only deployment. Overall, this workload shows how this scenario can enable greatly reduced query latencies for mixed workloads when compared to both the Cloud-only and the 2 Edge + Cloud scenarios, while maintaining a similar throughput. 4.6 Summary The testing of MulletBench’s rate limiting shows that it is able to successfully limit its Clients’ rate of operations, with the achieved throughput being very close to the targeted throughput. Its resource limiting 82
has also been validated, showing clear bottlenecks in resource usage and subsequent performance degradation on the limited nodes. Limiting network conditions by introducing latency between the Clients and database nodes shows an observable increase in latency of operations for both databases. When deployed in a Cloud-only configuration, with a single database node, IoTDB shows significantly better performance than InfluxDB with an insertion workload, achieving a 6.8x throughput increase. The query workload tests show IoTDB with a significant lower average latency for filter queries and a slight advantage for aggregations, while InfluxDB has a slight advantage in downsampling queries. When subjected to a mixed workload, InfluxDB struggles to handle the targeted throughput and performs worse than IoTDB for all query types, including 1413.6% higher average filter latency, while IoTDB is able to keep a constant throughput. Tests in a 2 Edge + Cloud scenario were composed of 2 resource limited nodes simulating the edge layer, with replication to a single cloud node. In this scenario, IoTDB achieves 8.0x the throughput of InfluxDB’s with the best tested downsampling configuration, when subjected to an insertion workload. When subjected to a query workload, at the edge nodes, InfluxDB achieves lower average latencies than IoTDB for aggregations and downsampling queries, while IoTDB has lower latency for filter operations. An issue with InfluxDB’s replication system has been identified, which prevented a correct comparison of the cloud layer query performance. This workload also allowed for the identification of a potential issue in IoTDB, in which replicated data is stored in unsequenced TsFiles on the Cloud node, resulting in degraded querying performance. Under a mixed workload, InfluxDB still struggles to handle the targeted throughput, even when using downsampling, while IoTDB matches the targeted throughput while also having lower latency for all query types at the edge layer. At the cloud layer, InfluxDB can achieve lower average latencies for aggregations and downsampling queries when using downsampling. Tests using downsampling at the edge layer with InfluxDB show it can achieve a 53% increase in throughput under an insertion workload, and a 40% increase in throughput, while increasing query latency at the edge up to 11% and decreasing latencies at the cloud, when under a mixed workload. Increasing the number of edge nodes to 4 results in significant performance improvements for both databases for all workloads. When using 2 second downsampling at the edge layer, InfluxDB achieves a 32.7% increase in throughput, while IoTDB achieves a 67.9% increase in throughput, when compared to the Cloud-only deployment, and even greater increases when compared to the 2 Edge + Cloud deployment. When subjected to a query workload, InfluxDB achieves lower latencies for all queries and at both layers, when compared to the Cloud-only deployment. Overall, IoTDB is unable to show any performance improvements for query workloads when compared to the Cloud-only deployment, apart from slightly lower 83
latencies for downsampling queries at the edge layer. Under a mixed workload, InfluxDB achieves similar throughput to the Cloud-only deployment, when using 2 second downsampling, and shows significant latency improvements for all query types on both layers. IoTDB achieves similar throughput to the Cloud-only deployment, while showing an up to 66.9% reduction in average latency at the edge layer, and up to 30.7% lower average latency at the cloud layer. Comparing both databases using this scenario does not yield significant differences, when compared to the 2 Edge + Cloud scenario These results provide insight into the performance of both databases under different workloads, deployment configurations, and downsampling options, through metrics which were not previously obtainable. In sum, they show how the edge layer can be used to improve performance, both by distributing the load over a higher number of nodes, and by performing downsampling, resulting in lower latencies at the cloud layer and higher throughput. 84
Chapter 5 Conclusions and future work In this thesis, a benchmarking system capable of providing metrics that accurately measure the performance of a TSDBMS deployed in an Edge-Cloud distributed architecture was developed. An in-depth study and survey of currently existing TSDBMS benchmarks builds a new taxonomy of the existing benchmarks, analysing their workloads, data, metrics and orchestration features, showing that existing TSDBMS benchmarking solutions lack metrics and orchestration features suitable for the evaluation of Edge-Cloud hybrid TSDBMSs. MulletBench addresses this problem by enabling the evaluation of TSDBMSs in Edge-Cloud hybrid deployments through the introduction of novel features such as per-node and per-layer metrics, automatic multi-layer database deployment, and orchestration of distributed clients. These features enable a finegrained analysis of the database’s performance in a multi-layer configuration, improve the reproducibility of the benchmark, ease the evaluation of different database configurations and allow for scaling workloads. The evaluation confirms MulletBench’s ability to limit its clients’ operation submission rate, the database nodes’ resources and the tests’ network conditions. An in-depth evaluation of InfluxDB and Apache IoTDB is presented, where MulletBench allowed for a thorough comparison of the performance of the two, as well as for great insight into the individual contributions of each edge node. This evaluation includes a comparison of the targeted databases in multiple workloads, downsampling configurations and deployment scenarios. Results show that the edge layer can be used to improve performance by distributing the workload over multiple layers and performing downsampling at the edge layer, increasing overall throughput and reducing cloud latency. Results also show that IoTDB generally presents superior performance to InfluxDB for all workloads. This evaluation also allowed for the identification of flaws in both systems’ replication mechanisms when these are subjected to high throughput workloads. The conclusions taken from the results are enabled by MulletBench’s novel features, and it would not have been possible to identify them with previously existing solutions, therefore demonstrating its 85
usefulness. 5.1 Future work MulletBench can currently be used to evaluate the performance of TSDBMSs in multi-layer deployments, but there are still some features that could be added to improve its functionality and usability. As highlighted by the testing campaign, databases sometimes fail to successfully replicate data, without any indication of this failure. This results in the need to manually check the validity of the results, by querying how much data the nodes contain. A validation phase could be implemented, in which, after the execution phase is completed, MulletBench would check if the data was successfully inserted into all nodes, allowing for the detection of replication failures, and eliminating the need for manual verification. Additional resource profiles for edge nodes could be added, allowing for the evaluation of databases in different hardware configurations. Support for additional TSDBMSs, including those which do not include explicit Edge-Cloud replication features, could also be added, allowing for the evaluation and comparison of more systems in multi-layer deployments. Finally, tests at a greater scale could also be conducted, including a greater number of edge nodes as well as multiple cloud nodes, allowing for the evaluation of more complex scenarios and approximating production environments. 86
Bibliography Ranking of time series dbms, 2022a. URL https://db-engines.com/en/ranking/time+ series+dbms. Accessed on December 2022. Dbms popularity broken down by database model, 2022b. URL https://db-engines.com/en/ ranking_categories. Accessed on December 2022. Edge-cloud collaboration, Oct 2023. URL https://iotdb.apache.org/UserGuide/V1.1.x/ Edge-Cloud-Collaboration/Sync-Tool.html#show-pipe. Accessed on October 2023. Michael P Andersen and David E Culler. Btrdb: Optimizing storage system design for timeseries processing. In 14th USENIX Conference on File and Storage Technologies (FAST 16) , pages 39–52, 2016. Andreas Bader. Comparison of time series databases. Master’s thesis, University of Stuttgart, 2016. Oana Balmau, Florin Dinu, Willy Zwaenepoel, Karan Gupta, Ravishankar Chandhiramoorthi, and Diego Didona. Ilk: Preventing latency spikes in log-structured merge key-value stores. In 2019 USENIX Annual Technical Conference (USENIX ATC 19) , pages 753–766, 2019. Maxim Buevich, Anne Wright, Randy Sargent, and Anthony Rowe. Respawn: A distributed multi-resolution time-series datastore. In 2013 IEEE 34th Real-Time Systems Symposium , pages 288–297. IEEE, 2013. Hongming Cai, Boyi Xu, Lihong Jiang, and Athanasios V Vasilakos. Iot-based big data storage systems in cloud computing: perspectives and challenges. IEEE Internet of Things Journal , 4(1):75–87, 2016. Stephan Cejka, Ralf Mosshammer, and Alfred Einfalt. Java embedded storage for time series and meta data in smart grids. In 2015 IEEE International Conference on Smart Grid Communications (SmartGridComm) , pages 434–439. IEEE, 2015. Yanpei Chen, Francois Raab, and Randy Katz. From tpc-c to big data benchmarks: A functional workload model. In Specifying Big Data Benchmarks , pages 28–43. Springer, 2012. 87
Brian F Cooper, Adam Silberstein, Erwin Tam, Raghu Ramakrishnan, and Russell Sears. Benchmarking cloud serving systems with ycsb. In Proceedings of the 1st ACM symposium on Cloud computing , pages 143–154, 2010. Ousmane Diallo, Joel JPC Rodrigues, Mbaye Sene, and Jaime Lloret. Distributed database management techniques for wireless sensor networks. IEEE Transactions on Parallel and Distributed Systems , 26 (2):604–620, 2013. Bernhard Firner, Robert S Moore, Richard Howard, Richard P Martin, and Yanyong Zhang. Poster: Smart buildings, sensor networks, and the internet of things. In Proceedings of the 9th ACM Conference on Embedded Networked Sensor Systems , pages 337–338, 2011. Peeyush Gupta, Michael J Carey, Sharad Mehrotra, and oberto Yus. Smartbench: a benchmark for data management in smart spaces. Proceedings of the VLDB Endowment , 13(12):1807–1820, 2020. Susan Hall. Iotdb provides data management for industrial edge it, Oct 2020. URL https:// thenewstack.io/iotdb-provides-data-management-for-industrial-edge-it/. Rui Han, Lizy Kurian John, and Jianfeng Zhan. Benchmarking big data systems: A review. IEEE Transactions on Services Computing , 11(3):580–597, 2017. Yuanzhe Hao, Xiongpai Qin, Yueguo Chen, Yaru Li, Xiaoguang Sun, Yu Tao, Xiao Zhang, and Xiaoyong Du. Ts-benchmark: A benchmark for time series databases. In 2021 IEEE 37th International Conference on Data Engineering (ICDE) , pages 588–599. IEEE, 2021. HashiCorp. Vagrant by hashicorp, 2010. URL https://www.vagrantup.com/. Accessed on January 2023. InfluxData. Influxdb time series data platform, 2013. URL https://www.influxdata.com/ products/influxdb/. Accessed on January 2023. Carlos Kamienski, Juha-Pekka Soininen, Markus Taumberger, Ramide Dantas, Attilio Toscano, Tullio Salmon Cinotti, Rodrigo Filev Maia, and André Torre Neto. Smart water management platform: Iotbased precision irrigation for agriculture. Sensors , 19(2):276, 2019. Roman Karlstetter, Robert Widhopf-Fenk, and Martin Schulz. Querying distributed sensor streams in the edge-to-cloud continuum. In 2022 IEEE International Conference on Edge Computing and Communications (EDGE) , pages 192–197. IEEE, 2022. 88
Table 19: Execution Orchestrator stage options Option Description stage integer to identify stage in order clients list of Client instance identifiers Table 20: Execution Orchestrator node options Option Description name identifier for database node layer which layer node belongs to, EDGE or CLOUD address IP address where database node is running monitor whether database node should be monitored with PCP, boolean A.3 Client Options Table 21: Client options Option Description clientId identifier for Client instance dataFile name of data file to be used in test dataset name of dataset to be used (GPSMPU) target database targeted by test numWorkers number of workers to use type type of load to generate (QUERY or INSERT) currentTime whether operations should use current time or original dataset time, boolean sharedDataset whether the loaded dataset should be shared by all workers, boolean sharedConnection whether database connection should be shared by all workers, boolean query query options insertion insertion options iotdb IoTDB options influx InfluxDB options orchestratorAddress IP address where orchestrator node is running 95
Table 22: Client query options Option Description count total number of queries to be submitted (per worker) rate number of queries to be submitted per second (per worker) aggChance weight for probability of aggregation filterChance weight for probability of filter outlierChance weight for probability of outlierFilter downsampleChance weight for probability of downsample numberOfLoops number of loops of data previously inserted (only for currentTime at false) countOutlierFilter whether outlier filter queries should use count aggregation filterZScore target z-score used for outlier filters Table 23: Client insertion options Option Description workerSectionSize total number of entries to be submitted batchSize number of batched entries per insertion rate number of insertions to be submitted per second (per worker) Table 24: Client InfluxDB options Option Description serverURL URL of server in format http://SERVER_IP:SERVER_PORT token security token used in instance orgID organization id used in instance bucket data bucket writeTimeout number of seconds before insertions should timeout readTimeout number of seconds before queries should timeout Table 25: Client IoTDB options Option Description host IP address of IoTDB instance port port of IoTDB instance username database username password database password devicePath path for devices in database queryAlignByDevice whether to queries should use align by device instead of time, boolean 96
Appendix B Test Configurations B.1 Cloud Scenario B.1.1 Insertion Workload InfluxDB Table 26: Client configuration for insertion workload targeting InfluxDB in 2 Edge + Cloud scenario Variable Value num_workers 80 insert_volume 96000 batch_size 100 insert_rate 1 IoTDB Table 27: Client configuration for insertion workload targeting IoTDB in 2 Edge + Cloud scenario Variable Value num_workers 300 insert_volume 96000 batch_size 100 insert_rate 1 97
B.1.2 Query Workload Table 28: Client configuration for query workload Cloud-only scenario Variable Value num_workers 2 query_count 192 query_rate 0.2 agg_weight 40 filter_weight 0 downsample_weight 20 outlier_filter_weight 40 count_outlier_filter false filter_z_score 4.0 B.1.3 Mixed Workload Table 29: Pre-population Client configuration for mixed workload Cloud-only scenario Variable Value num_workers 70 insert_volume 128000 batch_size 100 insert_rate 1 Table 30: Insert Client configuration for mixed workload Cloud-only scenario Variable Value num_workers 70 insert_volume 96000 batch_size 100 insert_rate 1 98
Table 31: Query Client configuration for mixed workload Cloud-only scenario Variable Value num_workers 2 query_count 240 query_rate 0.25 agg_weight 40 filter_weight 0 downsample_weight 20 outlier_filter_weight 40 count_outlier_filter false filter_z_score 4.0 B.2 2 Edge + Cloud Scenario B.2.1 Insertion Workload InfluxDB Table 32: Client configuration for insertion workload targeting InfluxDB in 2 Edge + Cloud scenario Variable Value num_workers 80 insert_volume 96000 batch_size 100 insert_rate 1 IoTDB Table 33: Client configuration for insertion workload targeting IoTDB in 2 Edge + Cloud scenario Variable Value num_workers 700 insert_volume 96000 batch_size 100 insert_rate 1 99
B.2.2 Query Workload Table 34: Query Client configuration for query workload in 2 Edge + Cloud scenario Variable Value num_workers 1 query_count 192 query_rate 0.2 agg_weight 40 filter_weight 0 downsample_weight 20 outlier_filter_weight 40 count_outlier_filter false filter_z_score 4.0 B.2.3 Mixed Workload Table 35: Pre-population Client configuration for mixed workload 2 Edge + Cloud scenario Variable Value num_workers 70 insert_volume 128000 batch_size 100 insert_rate 1 Table 36: Insert Client configuration for mixed workload in 2 Edge + Cloud scenario Variable Value num_workers 70 insert_volume 96000 batch_size 100 insert_rate 1 100
Table 37: Query Client configuration for mixed workload in 2 Edge + Cloud scenario Variable Value num_workers 1 query_count 240 query_rate 0.25 agg_weight 40 filter_weight 0 downsample_weight 20 outlier_filter_weight 40 count_outlier_filter false filter_z_score 4.0 Table 38: Cloud query Client configuration for mixed workload with downsampling in 2 Edge + Cloud scenario Variable Value num_workers 1 query_count 240 query_rate 0.25 agg_weight 40 filter_weight 0 downsample_weight 20 outlier_filter_weight 0 count_outlier_filter false filter_z_score 4.0 B.3 4 Edge + Cloud Scenario B.3.1 Insertion Workload InfluxDB Table 39: Client configuration for insertion workload targeting InfluxDB in 4 Edge + Cloud scenario Variable Value num_workers 80 insert_volume 96000 batch_size 100 insert_rate 1 101
IoTDB Table 40: Client configuration for insertion workload targeting IoTDB in 4 Edge + Cloud scenario Variable Value num_workers 700 insert_volume 96000 batch_size 100 insert_rate 1 B.3.2 Query Workload Table 41: Query Client configuration for query workload in 4 Edge + Cloud scenario Variable Value num_workers 1 query_count 192 query_rate 0.2 agg_weight 40 filter_weight 0 downsample_weight 20 outlier_filter_weight 40 count_outlier_filter false filter_z_score 4.0 B.3.3 Mixed Workload Table 42: Pre-population Client configuration for mixed workload 4 Edge + Cloud scenario Variable Value num_workers 70 insert_volume 128000 batch_size 100 insert_rate 1 102
Table 43: Insert Client configuration for mixed workload in 4 Edge + Cloud scenario Variable Value num_workers 70 insert_volume 96000 batch_size 100 insert_rate 1 Table 44: Query Client configuration for mixed workload in 4 Edge + Cloud scenario Variable Value num_workers 1 query_count 240 query_rate 0.25 agg_weight 40 filter_weight 0 downsample_weight 20 outlier_filter_weight 40 count_outlier_filter false filter_z_score 4.0 103