Bettina Kemme

dblp:k/BettinaKemme · DBLP profile ↗
← Back
74ranked-venue papers
7as first author
8since 2021 · last 2025
0000-0003-4694-4923ORCID · verified

Domains — the database's venue-derived domains; a paper can count in several

Databases, data management, data science and information retrieval · 27 · 4 first-author · 2 since 2021Systems, architecture and hardware · 17 · 3 first-author · 1 since 2021Software engineering, systems software and programming languages · 14 · 4 since 2021Security and privacy · 10 · 1 first-authorComputer networks · 5 · 1 since 2021Graphics, computer vision, multimedia, augmented reality and games · 2Human-computer interaction and ubiquitous computing · 1
YearPublicationVenuePosition
2025 G-View: View Management for Graph Databases
abstract
Graph database systems (GDBS) have become popular for representing real-world entities and their relationships, and offering convenient query languages based on graph pattern matching. As graphs increase in size and complexity, GDBS need to provide the appropriate support for abstraction for which views have demonstrated to be an effective tool, facilitating query writing and improving query execution time via materialization techniques. This paper explores how views can be defined and used in GDBS. We propose view-based extensions to the widely used graph query language Cypher, explore a wide range of possible view types, and outline several implementation strategies for view materialization. Using a set of micro- and macro-benchmarks, we provide insight into how expressive different view types are and how effective the proposed implementation strategies are for different GDBS. Our results show that views can be a powerful tool for GDBS, offering great flexibility in query expression and providing performance improvements if materialized.
Yunjia Zheng, Charlotte Sacré, Mohanna Shahrad, Owen Lipchitz, Yuting Gu, Bettina Kemme
Proc. VLDB Endow.6
2024 Give me some REST: A Controlled Experiment to Study Effects and Perception of Model-Driven Engineering with a Domain-Specific Language
abstract
Domain-Specific Languages (DSLs) are an efficient means to counter accidental complexity and are therefore a key technology for Model-Driven Engineering (MDE). Despite DSLs' potential, there is a lack of empirical research regarding the practical effects and developer perception of DSL-driven tools. In this paper, we present a controlled experiment with 28 participants around a previously developed DSL-based toolchain, which assists the migration of legacy software to REST. A direct comparison of developer performance for a) "DSL toolchain" and b) "classic manual software migration" allowed for analysis, quantification of effects and developer perception, as well as reasoning on general advantages, and DSL-related challenges. In certain cases, we measured a significant correlation between toolchain use and performance gains for developers. Detailed analysis of developer activities suggests the DSL toolchain alleviates tasks which show error-prone or time-consuming in the manual alternative. We then extracted acceptance-hindering factors from participant feedback and derived a series of recommendations for MDE practitioners who seek to develop DSL-based tools.
Maximilian Schiedermeier, Jörg Kienzle, Bettina Kemme
MODELS3
2023 Dynamic Application Call Graph Formation and Service Identification in Cloud Data Centers
abstract
Monitoring distributed service-based cloud applications and understanding the interactions among the different components are crucial to diagnose and resolve performance issues. However, many existing cloud monitoring systems require sophisticated application and/or platform instrumentation and cannot be deployed on-demand, or they provide only partial functionality. In most cases, monitoring comes with a significant overhead. To overcome these shortcomings, this paper presents DyMonD, a holistic framework that Dynamically Monitors an application, Discovers the service components, and visualizes them together with some performance metrics such as throughput in the form of a call graph. DyMonD is completely decoupled from the internals of the applications and the services themselves, as it deploys monitoring agents transparently at the software switches within the network, extracts all necessary information from the messages exchanged by the components, and performs service identification by using deep learning on the network flows. Our evaluation shows that DyMonD has significantly less overhead than existing tools, reducing the monitoring-induced impact on response time by up to 89%, and reducing resource consumption such as CPU and memory usage by up to 75%. Furthermore, DyMond’s deep learning module identifies services up to 11% more accurately than competing models.
Mona Elsaadawy, Mohamed F. Younis, Xinchen Hou, Bettina Kemme
IEEE Trans. Netw. Serv. Manag.5
2022 Global Decision Making Over Deep Variability in Feedback-Driven Software Development
abstract
To succeed with the development of modern software, organizations must have the agility to adapt faster to constantly evolving environments to deliver more reliable and optimized solutions that can be adapted to the needs and environments of their stakeholders including users, customers, business, development, and IT. However, stakeholders do not have sufficient automated support for global decision making, considering the increasing variability of the solution space, the frequent lack of explicit representation of its associated variability and decision points, and the uncertainty of the impact of decisions on stakeholders and the solution space. This leads to an ad-hoc decision making process that is slow, error-prone, and often favors local knowledge over global, organization-wide objectives. The Multi-Plane Models and Data (MP-MODA) framework explicitly represents and manages variability, impacts, and decision points. It enables automation and tool support in aid of a multi-criteria decision making process involving different stakeholders within a feedback-driven software development process where feedback cycles aim to reduce uncertainty. We present the conceptual structure of the framework, discuss its potential benefits, and enumerate key challenges related to tool supported automation and analysis within MP-MODA.
Jörg Kienzle, Benoît Combemale, Gunter Mussbacher, Omar Alam, Francis Bordeleau, Loli Burgueño, Gregor Engels, Jessie Galasso, Jean-Marc Jézéquel, Bettina Kemme, Sébastien Mosser 0001, Houari Sahraoui, Maximilian Schiedermeier, Eugene Syriani
ASE10
2021 HorsePower: Accelerating Database Queries for Advanced Data Analytics
Hanfeng Chen, Joseph Vinish D'silva, Laurie J. Hendren, Bettina Kemme
EDBT4
2021 Demo: Application Monitoring as a Network Service
abstract
The recent rise of cloud applications, representing large complex modern distributed services, has made performance monitoring a major issue and a critical process for both cloud providers and cloud customers. Many different monitoring techniques are used such as tracking resource consumption, performing application-specific measures or analyzing message exchanges. Typically the collected data is logged at the host on which the application is deployed, then either analyzed locally or forwarded to a remote analysis host. In contrast, this demonstration paper presents a Monitoring as a Service prototype that uses the advances in Software Defined Networking (SDN) to move some of the logging functionality into the network. The core of our MaaS is implemented as a virtual network function where agents are co-located with software switches in order to extract performance metrics from the message flows between components in a non-intrusive manner and send the calculated measures to the clients for visualization in near real-time. The MaaS has a lot of flexibility in how it is deployed and does not require to instrument software or platforms. In our demo we show the tool in action demonstrating how users can choose to monitor different service types and performance metrics in a user-friendly manner.
Mona Elsaadawy, Laetitia Fesselier, Bettina Kemme
ICDCS3
2021 Flow-based Service Type Identification using Deep Learning
abstract
Automatic identification of the service type used by network flows (e.g., HTTP and MySQL) is an essential part of many cloud management and monitoring tasks for quality of service, security monitoring, resource allocation, etc. Several studies have adapted deep learning models for accurate service type identification of network traffic. These models vary in how the message flow data is used and what datasets are considered. There are no published guidelines on selecting the best approach for automating the service identification process. In this paper, we opt to fill such a technical gap and provide a detailed study of the trade-offs of different deep-learning based approaches for service type identification of network traffic. Towards this end, we generate flow-based datasets for a wide range of service types that are commonly deployed in the cloud. We consider two different deep learning models that have shown promising results in this context, and show their performance for both payload- and header-based datasets, considering fundamental parameters such as dynamic service port configuration, flow direction and the packet order in the flow stream.
Mona Elsaadawy, Petar Basta, Yunjia Zheng, Bettina Kemme, Mohamed F. Younis
NetSoft4
2021 FIDDLR: streamlining reuse with concern-specific modelling languages
abstract
Model-Driven Engineering (MDE) reduces complexity, improves Separation of Concerns and promotes reuse by structuring software development as a process of model production and refinement. Domain-Specific Modelling Languages and Aspect-Oriented Modelling techniques can reduce complexity and improve modularization of crosscutting concerns in situations where the features of general purpose modelling languages are not well aligned with the subject of study. In this article we present FIDDLR, a novel framework that integrates the ideas of Domain-Specific Modelling Languages, Concern-Oriented Reuse and MDE to modularize concerns that cross-cut multiple levels of abstraction of the software development process and streamline the reuse process. It also prescribes the integration of the different tooling along this process. We demonstrate the effectiveness of our framework and the potential for reduced complexity and leveraged reuse by building a reusable concern that exposes the services a system offers through a REST interface.
Maximilian Schiedermeier, Jörg Kienzle, Bettina Kemme
SLE3
2020 Improving database query performance with automatic fusion
abstract
Array-based programming languages have shown significant promise for improving performance of column-based in-memory database systems, allowing elegant representation of query execution plans that are also amenable to standard compiler optimization techniques. Use of loop fusion, however, is not straightforward, due to the complexity of built-in functions for implementing complex database operators. In this work, we apply a compiler approach to optimize SQL query execution plans that are expressed in an array-based intermediate representation. We analyze this code to determine shape properties of the data being processed, and use a subsequent optimization phase to fuse multiple database operators into single, compound operations, reducing the need for separate computation and storage of intermediate values. Experimental results on a range of TPC-H queries show that our fusion technique is effective in generating efficient code, improving query time over a baseline system.
Hanfeng Chen, Alexander Krolik, Bettina Kemme, Clark Verbrugge, Laurie J. Hendren
CC3
2020 Enabling Efficient Application Monitoring in Cloud Data Centers using SDN
abstract
Nowadays, many cloud applications can be considered large complex distributed services. The increasing sophistication and complexity have made performance monitoring a major issue and a critical process for both cloud providers and cloud customers. Existing monitoring techniques instrument the applications to collect measurements and log them at the host nodes on which the application is deployed. Such an approach introduces overhead and can slow down the application. New paradigms such as Software Defined Networking (SDN) and Network Function Virtualization (NFV) show promise for moving some of the measurement collection and logging functionality into the network as a lot of information can be extracted from messages exchanged between application components. Such a methodology could enable the cloud infrastructure to provide a Monitoring as a Service to applications in a transparent manner without software instrumentation and allowing for a more flexible placement of logging functionality. In this paper, we explore mechanisms to integrate application monitoring into SDN. In particular, we analyze whether switch based message filtering is feasible and we propose a customized port sniffing approach. We discuss the implementation aspects using OVS. The results confirm that moving application monitoring to the network is indeed an attractive option.
Mona El Saadawy, Bettina Kemme, Mohamed F. Younis
ICC2
2019 Keep Your Host Language Object and Also Query it: A Case for SQL Query Support in RDBMS for Host Language Objects
abstract
As a result of prolific growth in data science and machine learning applications, modern relational database management systems (RDBMS) are experimenting with various approaches to facilitate advanced analytical computations, in addition to the relational operations that they traditionally support. The most common approach has been to integrate an embedded high level language (HLL) interpreter into the RDBMS along with any additional libraries that specialize in numerical computations. Such implementations, e.g., user defined functions (UDFs), follow generally a black-box setup, and for many complex workflows that require datasets to be passed and processed back-and-forth between the query execution engine and the embedded HLL interpreter, optimization opportunities are not fully explored yet. In this paper, we propose and implement the concept of virtual tables that can be used to expose data set objects maintained by the embedded HLL interpreter to the query engine for executing relational operations. Unlike prevalent solutions, our approach minimizes the need for performing data copies and conversions, performing them lazily when required. It also facilitates better optimization opportunities for the execution of SQL queries as the RDBMS is able to analyze the data characteristics of the HLL objects before producing an execution plan. The approach is also programmer friendly, allowing for a more intuitive implementation of computational workflows. We perform evaluations over a variety of workloads which demonstrate the performance and programming benefits of virtual tables.
Joseph Vinish D'silva, Florestan De Moor, Bettina Kemme
SSDBM3
2019 A scalable network-aware framework for cloud monitoring orchestration
Masoume Jabbarifar, Alireza Shameli-Sendi, Bettina Kemme
J. Netw. Comput. Appl.3
2019 Making an RDBMS Data Scientist Friendly: Advanced In-database Interactive Analytics with Visualization Support
abstract
We are currently witnessing the rapid evolution and adoption of various data science frameworks that function external to the database. Any support from conventional RDBMS implementations for data science applications has been limited to procedural paradigms such as user-defined functions (UDFs) that lack exploratory programming support. Therefore, the current status quo is that during the exploratory phase, data scientists usually use the database system as the "data storage" layer of the data science framework, whereby the majority of computation and analysis is performed outside the database, e.g., at the client node. We demonstrate AIDA, an in-database framework for data scientists. AIDA allows users to write interactive Python code using a development environment such as a Jupyter notebook. The actual execution itself takes place inside the database (near-data), where a server component of AIDA, that resides inside the embedded Python interpreter of the RDBMS, manages the data sets and computations. The demonstration will also show the visualization capabilities of AIDA where the progress of computation can be observed through live updates. Our evaluations show that AIDA performs several times faster compared to contemporary external data science frameworks, but is much easier to use for exploratory development compared to database UDFs.
Joseph Vinish D'silva, Florestan De Moor, Bettina Kemme
Proc. VLDB Endow.3
2018 HorseIR: bringing array programming languages together with database query processing
abstract
Relational database management systems (RDBMS) are operationally similar to a dynamic language processor. They take SQL queries as input, dynamically generate an optimized execution plan, and then execute it. In recent decades, the emergence of in-memory databases with columnar storage, which use array-like storage structures, has shifted the focus on optimizations from the traditional I/O bottleneck to CPU and memory. However, database research so far has primarily focused on CPU cache optimizations. The similarity in the computational characteristics of such database workloads and array programming language optimizations are largely unexplored. We believe that these database implementations can benefit from merging database optimizations with dynamic array-based programming language approaches. Therefore, in this paper, we propose a novel approach to optimize database query execution using a new array-based intermediate representation, HorseIR, that resides between database queries and compiled code. Furthermore, we provide a translator to generate HorseIR from database execution plans and a compiler that optimizes HorseIR and generates efficient code. We compare HorseIR with the MonetDB RDBMS, by testing standard SQL queries, and show how our approach and compiler optimizations improve the runtime of complex queries.
Hanfeng Chen, Joseph Vinish D'silva, Hongji Chen 0001, Bettina Kemme, Laurie J. Hendren
DLS4
2018 AIDA - Abstraction for Advanced In-Database Analytics
abstract
With the tremendous growth in data science and machine learning, it has become increasingly clear that traditional relational database management systems (RDBMS) are lacking appropriate support for the programming paradigms required by such applications, whose developers prefer tools that perform the computation outside the database system. While the database community has attempted to integrate some of these tools in the RDBMS, this has not swayed the trend as existing solutions are often not convenient for the incremental, iterative development approach used in these fields. In this paper, we propose AIDA - an abstraction for advanced in-database analytics. AIDA emulates the syntax and semantics of popular data science packages but transparently executes the required transformations and computations inside the RDBMS. In particular, AIDA works with a regular Python interpreter as a client to connect to the database. Furthermore, it supports the seamless use of both relational and linear algebra operations using a unified abstraction. AIDA relies on the RDBMS engine to efficiently execute relational operations and on an embedded Python interpreter and NumPy to perform linear algebra operations. Data reformatting is done transparently and avoids data copy whenever possible. AIDA does not require changes to statistical packages or the RDBMS facilitating portability.
Joseph Vinish D'silva, Florestan De Moor, Bettina Kemme
Proc. VLDB Endow.3
2017 Powering Archive Store Query Processing via Join Indices
Joseph Vinish D'silva, Bettina Kemme, Richard Grondin, Evgueni Fadeitchev
EDBT2
2017 Self-Evolving Subscriptions for Content-Based Publish/Subscribe Systems
abstract
Traditional pub/sub systems cannot adequately handle workloads of applications with dynamic, short-lived subscriptions such as location-based social networks, predictive stock trading, and online games. Subscribers must continuously interact with the pub/sub system to remove and insert subscriptions, thereby inefficiently consuming network and computing resources, and sacrificing consistency. In the aforementioned applications, we recognize that the changes in the subscriptions can follow a predictable pattern over some variable (e.g., time). In this paper, we present a new type of subscription, called evolving subscription, which encapsulates these patterns and allow the pub/sub system to autonomously adapt to the dynamic interests of the subscribers without incurring an expensive re-subscription overhead. We propose a general model for expressing evolving subscriptions and a framework for supporting them in a pub/sub system. To this end, we propose three different designs to support evolving subscriptions, which are evaluated and compared to the traditional resubscription approach in the context of two use cases: online games and high-frequency trading. Our evaluation shows that our solutions can reduce subscription traffic by 96.8% and improve delivery accuracy when compared to the baseline resubscription mechanism.
César Cañas, Kaiwen Zhang 0001, Bettina Kemme, Jörg Kienzle, Hans-Arno Jacobsen
ICDCS3
2017 MultiPub: Latency and Cost-Aware Global-Scale Cloud Publish/Subscribe
abstract
Topic-based pub/sub is a widely used communication mechanism in distributed systems for targeted information dissemination between loosely coupled entities. To scale dynamically depending on the current communication demands, pub/services can be conveniently deployed in the cloud. To provide fast dissemination, the service can be distributed across multiple cloud regions. The architectural design and run-time deployment of such a middleware is tricky, though, as it can have a significant effect on communication latency and cloud-based cost. In this paper, we propose MultiPub, a flexible pub/sub middleware for latency-constrained, world-wide distributed applications that dynamically reconfigures the communication layer to ensure a predefined maximum latency for publication dissemination while minimizing cloud-based costs. This is achieved by routing publications either through a single or across multiple cloud regions. We demonstrate the effectiveness of MultiPub by presenting a set of experiments that report on the achieved communication latency and cost savings compared to traditional approaches, as well as a performance evaluation.
Julien Gascon-Samson, Jörg Kienzle, Bettina Kemme
ICDCS3
2016 AdaptCache: Adaptive Data Partitioning and Migration for Distributed Object Caches
Omar Asad, Bettina Kemme
Middleware2
2015 Dynamoth: A Scalable Pub/Sub Middleware for Latency-Constrained Applications in the Cloud
abstract
This paper presents Dynamoth, a dynamic, scalable, channel-based pub/sub middleware targeted at large scale, distributed and latency constrained systems. Our approach provides a software layer that balances the load generated by a high number of publishers, subscribers and messages across multiple, standard pub/sub servers that can be deployed in the Cloud. In order to optimize Cloud infrastructure usage, pub/sub servers can be added or removed as needed. Balancing takes into account the live characteristics of each channel and is done in an hierarchical manner across channels (macro) as well as within individual channels (micro) to maintain acceptable performance and low latencies despite highly varying conditions. Load monitoring is performed in an unintrusive way, and rebalancing employs a lazy approach in order to minimize its temporal impact on performance while ensuring successful and timely delivery of all messages. Extensive real-world experiments that illustrate the practicality of the approach within a massively multiplayer game setting are presented. Results indicate that with a given number of servers, Dynamoth was able to handle 60% more simultaneous clients than the consistent hashing approach, and that it was properly able to deal with highly varying conditions in the context of large workloads.
Julien Gascon-Samson, Franz-Philippe Garcia, Bettina Kemme, Jörg Kienzle
ICDCS3
2015 Monitoring Large-Scale Location-Based Information Systems
abstract
Monitoring the state of a distributed virtual world is challenging for several reasons: 1) the distributed information must be gathered in real-time without affecting the performance of the information system, 2) in large-scale systems it is impossible for a single node to collect and process all the data, 3) the vast information must be filtered and aggregated according to what the human observer wants to focus on, and 4) the point of interest of the observer can change frequently. In this paper we present and evaluate a non-intrusive monitoring middleware that addresses these challenges by dynamically partitioning the geographic map (e.g., of the virtual world or the game) in terms of map objects and (expected) state changes. We assign a different collector node to each of these partitions to collect and pre-process the data, and forward it to a central monitoring node. Furthermore, we provide mechanisms to efficiently filter and aggregate location changes, the pre-dominant changes in location-based information systems. We describe a specific monitoring setup that takes advantage of the replication model that is common in many virtual worlds and multiplayer games to collect the data. Finally, we present extensive performance results that show the trade-offs between scalability, precision, and real-time performance.
Hammad Khan, Julien Gascon-Samson, Jörg Kienzle, Bettina Kemme
IPDPS4
2015 GraPS: A Graph Publish/Subscribe Middleware
abstract
Pub/sub is an elegant paradigm for disseminating information efficiently and anonymously among producers (publishers) and consumers (subscribers). However, with current topic and content-based pub/sub approaches it is difficult to formulate subscriptions that adequately and accurately express the interest of the consumers in a semantic information domain. In this paper we introduce GraPS, a pub/sub middleware that provides a publication model based on graphs. Points of interest in the information domain are mapped to nodes, and relationships between points of interest are mapped to edges. Consumers can effectively express their interest in publications by means of graph subscriptions that exploit the properties of nodes and the semantics of the edge relationships. Graph subscriptions do not require complete knowledge of the graph and can be updated whenever the consumer's interest changes. Furthermore, graph subscriptions are automatically updated whenever the information domain changes. We illustrate GraPS by means of three application scenarios and present a set of experiments with an implementation of GraPS based on standard pub/sub middleware.
César Cañas, Eduardo Pacheco, Bettina Kemme, Jörg Kienzle, Hans-Arno Jacobsen
Middleware3
2015 Compaction Management in Distributed Key-Value Datastores
abstract
Compactions are a vital maintenance mechanism used by datastores based on the log-structured merge-tree to counter the continuous buildup of data files under update-intensive workloads. While compactions help keep read latencies in check over the long run, this comes at the cost of significantly degraded read performance over the course of the compaction itself. In this paper, we offer an in-depth analysis of compaction-related performance overheads and propose techniques for their mitigation. We offload large, expensive compactions to a dedicated compaction server to allow the datastore server to better utilize its resources towards serving the actual workload. Moreover, since the newly compacted data is already cached in the compaction server's main memory, we fetch this data over the network directly into the datastore server's local cache, thereby avoiding the performance penalty of reading it back from the filesystem. In fact, pre-fetching the compacted data from the remote cache prior to switching the workload over to it can eliminate local cache misses altogether. Therefore, we implement a smarter warmup algorithm that ensures that all incoming read requests are served from the datastore server's local cache even as it is warming up. We have integrated our solution into HBase, and using the YCSB and TPC-C benchmarks, we show that our approach significantly mitigates compaction-related performance problems. We also demonstrate the scalability of our solution by distributing compactions across multiple compaction servers.
Muhammad Yousuf Ahmad, Bettina Kemme
Proc. VLDB Endow.2
2014 Publish/subscribe network designs for multiplayer games
abstract
Massively multiplayer online games (MMOGs), which are typically supported by large distributed systems, require a scalable, low latency messaging middleware that supports the location-based semantics and the loosely coupled interaction of multiplayer games components. In this paper, we present three different pub/sub-driven designs for a MMOG networking engine that account for the highly interactive and massive nature of these games. Each design uses not only different pub/sub approaches (from topic-based to content-based) but also serves varying degrees of responsibilities. In particular, some of them integrate game functionality, such as interest management, into the network engine. We implement, evaluate, and compare our proposed designs in the MMOG prototype Mammoth. Our real-world results show the viability of pub/sub while at the same time highlighting clear trade-offs between the different designs used, especially in the number and frequency of the various message types, such as subscriptions.
César Cañas, Kaiwen Zhang 0001, Bettina Kemme, Jörg Kienzle, Hans-Arno Jacobsen
Middleware3
2014 Consistency anomalies in multi-tier architectures: automatic detection and prevention
Kamal Zellag, Bettina Kemme
VLDB J.2
2013 Watchmen: Scalable Cheat-Resistant Support for Distributed Multi-player Online Games
abstract
Multi-player online games are inherently distributed applications, and a wide range of distributed architectures have been proposed. However, only few successful commercial systems follow such approaches, even given their benefits, due to one main hurdle: the easiness with which cheaters can disrupt the game state computation and dissemination, perform illegal actions, or unduly gain access to sensitive information. The challenge is that any measures used to address cheating must meet the heavy scalability and tight latency requirements of fast paced games. We propose Watchmen, the first distributed scalable protocol designed with cheat detection and prevention in mind that supports fast paced games. It is based on a randomized dynamic proxy scheme for both the dissemination and verification of actions. Furthermore, Watchmen reduces the information exposed to players close to the minimum required to render the game. We build our proof-of-concept prototype on top of Quake III. We show that Watchmen, while scaling to hundreds of players and meeting the tight latency requirements of first person shooter games, is able to significantly reduce opportunities to cheat, even in the presence of collusion.
Amir Yahyavi, Kévin Huguenin, Julien Gascon-Samson, Jörg Kienzle, Bettina Kemme
ICDCS5
2013 Transactional Failure Recovery for a Distributed Key-Value Store
Muhammad Yousuf Ahmad, Bettina Kemme, Ivan Brondino, Marta Patiño-Martínez, Ricardo Jiménez-Peris
Middleware2
2013 Interest modeling in games: the case of dead reckoning
Amir Yahyavi, Kévin Huguenin, Bettina Kemme
Multim. Syst.3
2012 How consistent is your cloud application?
abstract
Current cloud datastores usually trade consistency for performance and availability. However, it is often not clear how an application is affected when it runs under a low level of consistency. In fact, current application designers have basically no tools that would help them to get a feeling of which and how many inconsistencies actually occur for their particular application. In this paper, we propose a generalized approach for detecting consistency anomalies for arbitrary cloud applications accessing various types of cloud datastores in transactional or non-transactional contexts. We do not require any knowledge on the business logic of the studied application nor on its selected consistency guarantees. We experimentally verify the effectiveness of our approach by using the Google App Engine and Cassandra datastores.
Kamal Zellag, Bettina Kemme
SoCC2
2012 ConsAD: a real-time consistency anomalies detector
abstract
In this demonstration, we present ConsAD, a tool that detects consistency anomalies for arbitrary multi-tier applications that use lower levels of isolation than serializability. As the application is running, ConsAD detects and quantifies anomalies indicating exactly the transactions and data items involved. Furthermore, it classifies the detected anomalies into patterns showing the business methods involved as well as their occurrence frequency. ConsAD can guide designers to either choose an isolation level for which their application shows few anomalies or change their transaction design to avoid the anomalies. Its graphical interface shows detailed information about detected anomalies as they occur and analyzes their patterns as well as their distribution.
Kamal Zellag, Bettina Kemme
SIGMOD Conference2
2011 Distributed data management in 2020?
abstract
Work on distributed data management commenced shortly after the introduction of the relational model in the mid-1970's. 1970's and 1980's were very active periods for the development of distributed relational database technology, and claims were made that in the following ten years centralized databases will be an “antique curiosity” and most organizations will move toward distributed database managers [1]. That prediction has certainly become true, and all commercial DBMSs today are distributed.
M. Tamer Özsu, Patrick Valduriez, Serge Abiteboul, Bettina Kemme, Ricardo Jiménez-Peris, Beng Chin Ooi
ICDE4
2011 Real-time quantification and classification of consistency anomalies in multi-tier architectures
abstract
While online transaction processing applications heavily rely on the transactional properties provided by the underlying infrastructure, they often choose to not use the highest isolation level, i.e., serializability, because of the potential performance implications of costly strict two-phase locking concurrency control. Instead, modern transaction systems, consisting of an application server tier and a database tier, offer several levels of isolation providing a trade-off between performance and consistency. While it is fairly well known how to identify the anomalies that are possible under a certain level of isolation, it is much more difficult to quantify the amount of anomalies that occur during run-time of a given application. In this paper, we address this issue and present a new approach to detect, in realtime, consistency anomalies for arbitrary multi-tier applications. As the application is running, our tool detect anomalies online indicating exactly the transactions and data items involved. Furthermore, we classify the detected anomalies into patterns showing the business methods involved as well as their occurrence frequency. We use the RUBiS benchmark to show how the introduction of a new transaction type can have a dramatic effect on the number of anomalies for certain isolation levels, and how our tool can quickly detect such problem transactions. Therefore, our system can help designers to either choose an isolation level where the anomalies do not occur or to change the transaction design to avoid the anomalies.
Kamal Zellag, Bettina Kemme
ICDE2
2011 Transaction Models for Massively Multiplayer Online Games
abstract
Massively Multiplayer Online Games are considered large distributed systems where the game state is partially replicated across the server and thousands of clients. Given the scale, game engines typically offer only relaxed consistency without well-defined guarantees. In this paper, we leverage the concept of transactions to define consistency models that are suitable for gaming environments. We define game specific levels of consistency that differ in the degree of isolation and atomicity they provide, and demonstrate the costs associated with their execution. Each action type within a game can then be assigned the appropriate consistency level, choosing the right trade-off between consistency and performance. The issue of durability and fault-tolerance of game actions is also discussed.
Kaiwen Zhang 0001, Bettina Kemme
SRDS2
2011 Building a peer-to-peer content distribution network with high performance, scalability and robustness
Manal El Dick, Esther Pacitti, Reza Akbarinia, Bettina Kemme
Inf. Syst.4
2011 Elastic SI-Cache: consistent and scalable caching in multi-tier architectures
Francisco Perez-Sorrosal, Marta Patiño-Martínez, Ricardo Jiménez-Peris, Bettina Kemme
VLDB J.4
2010 Database Replication: a Tale of Research across Communities
abstract
Replication is a key mechanism to achieve scalability and fault-tolerance in databases. Its importance has recently been further increased because of the role it plays in achieving elasticity at the database layer. In database replication, the biggest challenge lies in the trade-off between performance and consistency. A decade ago, performance could only be achieved through lazy replication at the expense of transactional guarantees. The strong consistency of eager approaches came with a high cost in terms of reduced performance and limited scalability. Postgres-R combined results from distributed systems and databases to develop a replication solution that provided both scalability and strong consistency. The use of group communication primitives with strong ordering and delivery guarantees together with optimized transaction handling (tailored locking, transferring logs instead of re-executing updates, keeping the message overhead per transaction constant) were a drastic departure from the state-of-the-art at the time. Ten years later, these techniques are widely used in a variety of contexts but particularly in cloud computing scenarios. In this paper we review the original motivation for Postgres-R and discuss how the ideas behind the design have evolved over the years.
Bettina Kemme, Gustavo Alonso
Proc. VLDB Endow.1
2009 Flower-CDN: a hybrid P2P overlay for efficient query processing in CDN
abstract
Many websites with a large user base, e.g., websites of nonprofit organizations, do not have the financial means to install large web-servers or use specialized content distribution networks such as Akamai. For those websites, we have developed Flower-CDN, a locality-aware P2P based content-distribution network (CDN) in which the users that are interested in a website support the distribution of its content. The idea is that peers keep the content they retrieve and later serve it to other peers that are close to them in locality. Our architecture is a hybrid between structured and unstructured networks. When a new client requests some content from a website, a locality-aware DHT quickly finds a peer in its neighborhood that has the content available. Additionally, all peers in a given locality that maintain content of a particular website build an unstructured content overlay. Within this overlay, peers gossip information about their content allowing the system to maintain accurate information despite churn. In our performance evaluation, we compare Flower-CDN with an existing P2P-CDN strictly based on DHT and not locality aware. Flower-CDN reduces lookup latency by a factor of 9 and transfer distance by a factor of 2. We also show that Flower-CDN's gossip has low overhead and can be adjusted according to hit ratio requirements and bandwidth availability.
Manal El Dick, Esther Pacitti, Bettina Kemme
EDBT3
2009 A Unified Framework for Load Distribution and Fault-Tolerance of Application Servers
Huaigu Wu, Bettina Kemme
Euro-Par2
2009 Mammoth: a massively multiplayer game research framework
abstract
This paper presents Mammoth, a massively multiplayer game research framework designed for experimentation in an academic setting. Mammoth provides a modular architecture where different components, such as the network engine, the replication engine, or interest management, can easily be replaced. Subgames allow a researcher to define different game goals, for instance, in order to evaluate the effects of different team-play tactics on the game performance. Mammoth also offers a modular and flexible infrastructure for the definition of non-player characters with behavior controlled by complex artificial intelligence algorithms. This paper focuses on the Mammoth architecture, demonstrating how good design practices can be used to create a modular framework where researchers from different research domains can conduct their experiments. The effectiveness of the architecture is demonstrated by several successful research projects accomplished using the Mammoth framework.
Jörg Kienzle, Clark Verbrugge, Bettina Kemme, Alexandre Denault, Michael Hawker
FDG3
2009 Snapshot isolation and integrity constraints in replicated databases
abstract
Database replication is widely used for fault tolerance and performance. However, it requires replica control to keep data copies consistent despite updates. The traditional correctness criterion for the concurrent execution of transactions in a replicated database is 1-copy-serializability. It is based on serializability, the strongest isolation level in a nonreplicated system. In recent years, however, Snapshot Isolation (SI), a slightly weaker isolation level, has become popular in commercial database systems. There exist already several replica control protocols that provide SI in a replicated system. However, most of the correctness reasoning for these protocols has been rather informal. Additionally, most of the work so far ignores the issue of integrity constraints. In this article, we provide a formal definition of 1-copy-SI using and extending a well-established definition of SI in a nonreplicated system. Our definition considers integrity constraints in a way that conforms to the way integrity constraints are handled in commercial systems. We discuss a set of necessary and sufficient conditions for a replicated history to be producible under 1-copy-SI. This makes our formalism a convenient tool to prove the correctness of replica control algorithms.
Yi Lin 0005, Bettina Kemme, Ricardo Jiménez-Peris, Marta Patiño-Martínez, José Enrique Armendáriz-Iñigo
ACM Trans. Database Syst.2
2008 Maintaining replicas in unstructured P2P systems
abstract
Replication is widely used in unstructured peer-to-peer systems to improve search or achieve availability. We identify and solve a subclass of replication problems where each object is associated with a maintainer node, and its replicas should only be available as long as its maintainer is part of the network. Such requirement can be found in various applications, e.g., when objects are directory lists, service lists, or subscriptions of a publish/subscribe system. We provide maintainers with proven guarantees on the number of replicas, in spite of network churn and crash failures. We also tackle the related problems of changing the number of replicas, updating replicas, balancing storage load in a heterogeneous network, and eliminating replicas left by crashing maintainers. Our algorithm is based on probabilistic methods and is simple to implement. We show by simulation and formal proof that our algorithm is correct. 1.
Christof Leng, Wesley W. Terpstra, Bettina Kemme, Wilhelm Stannat, Alejandro P. Buchmann
CoNEXT3
2008 Online recovery in cluster databases
abstract
Cluster based replication solutions are an attractive mechanism to provide both high-availability and scalability for the database backend within the multi-tier information systems of service-oriented businesses. An important issue that has not yet received sufficient attention is how database replicas that have failed can be reintegrated into the system or how completely new replicas can be added in order to increase the capacity of the system. Ideally, recovery takes place online, i.e, while transaction processing continues at the replicas that are already running. In this paper we present a complete online recovery solution for database clusters. One important issue is to find an efficient way to transfer the data the joining replica needs. In this paper, we present two data transfer strategies. The first transfers the latest copy of each data item, the second transfers the updates a rejoining replica has missed during its downtime. A second challenge is to coordinate this transfer with ongoing transaction processing such that the joining node does not miss any updates. We present a coordination protocol that can be used with Postgres-R, a replication tool which uses a group communication system for replica control. We have implemented and compared our transfer solutions against a set of parameters, and present heuristics which allow an automatic selection of the optimal strategy for a given configuration.
WeiBin Liang, Bettina Kemme
EDBT2
2008 Showing correctness of a replication algorithm in a component based system
abstract
Reasoning about the correctness of a replication algorithm is a difficult endeavor. If correctness has to be shown for a component based architecture where a client request can lead to execution across different components or tiers, this is even more difficult. Existing formalisms are either restricted to systems with only one component, or make strong assumptions about the setup of the system. In this paper, we present a flexible framework to reason about exactly-once execution in a failure-prone replicated component based system. Our approach allows us to reason about the execution across the entire system, e.g., application server and database tier. If a given replication algorithm makes assumptions about some of the components, then those can be easily integrated into the reasoning process.
Huaigu Wu, Bettina Kemme
IDEAS2
2008 pSense - Maintaining a Dynamic Localized Peer-to-Peer Structure for Position Based Multicast in Games
abstract
This paper presents an algorithm for creating and maintaining a dynamic localized peer-to-peer overlay network with its main application to massively multiplayer games. In these games, players reside in a large game world with many thousands of players but each player has typically a limited vision range. In our solution, players join the network as peers and mainly connect to neighbor peers that are close to them in the virtual game world. As players move in the game they change their neighbors dynamically with very little overhead. Peers can multicast messages that are received by peers in their locality very fast (often faster than in client-server solutions) while players that are further away receive them later or not at all. Not receiving messages from remote players is important in order to not cause the load on each peer to grow with the number of players in the game. Our performance analysis confirms that our solution allows for dynamic game worlds of practically unlimited size, only limited in scale by the number of players within the vision range.
Arne Schmieg, Michael Stieler, Sebastian Jeckel, Patric Kabus, Bettina Kemme, Alejandro P. Buchmann
Peer-to-Peer Computing5
2008 An Autonomic Approach for Replication of Internet-based Services
abstract
As more and more applications are deployed as Internet-based services, they have to be available anytime anywhere in a seamless manner. This requires the underlying infrastructure to provide scalability, fault tolerance and fast response times. While replicating the services and the data they access across sites that are located in different geographic regions is a promising means to achieve these requirements, data consistency is challenging if data continuously changes and queries are dynamic by nature, as is typical for e-commerce applications.Thus, current WAN replication solutions either trade performance for data consistency or are notable to scale in wide-area settings. In this paper, we present a novel approach to provide performance and consistency for Internet services. One of the main contributions is an autonomic replica placement module that places data copies only on servers close to clients that actually need them. The goal is to find the right trade-off between fast local access and the overhead of keeping data copies consistent. As data access patterns might change over time, reconfiguration is done periodically and online, i.e., allowing sites to receive new data copies or drop data copies while at the same time transaction processing continues in the system.
Damián Serrano, Marta Patiño-Martínez, Ricardo Jiménez-Peris, Bettina Kemme
SRDS4
2007 A Recovery Protocol for Middleware Replicated Databases Providing GSI
abstract
Middleware database replication is a way to increase availability and afford site failures for dynamic content Websites. There are several replication protocols that ensure data consistency for these systems. The most attractive ones are those providing generalized snapshot isolation (GSI), as read operations never block. These replication protocols are based on the certification process, however, up to our knowledge, they do not cope with the recovery of a replica. In this paper we propose a recovery protocol that ensures GSI (we provide an outline of its correctness) that does not interfere with user transactions and permits the execution of transactions in the recovering node, even though the recovery process has not finished
José Enrique Armendáriz-Iñigo, Francesc D. Muñoz-Escoí, José Ramón Juárez-Rodríguez, José Ramón González de Mendívil, Bettina Kemme
ARES5
2007 Topic 5 Parallel and Distributed Databases
Marta Patiño-Martínez, Genoveva Vargas-Solar, Elena Baralis, Bettina Kemme
Euro-Par4
2007 Consistent and Scalable Cache Replication for Multi-tier J2EE Applications
Francisco Perez-Sorrosal, Marta Patiño-Martínez, Ricardo Jiménez-Peris, Bettina Kemme
Middleware4
2007 Boosting Database Replication Scalability through Partial Replication and 1-Copy-Snapshot-Isolation
abstract
Databases have become a crucial component in modern information systems. At the same time, they have become the main bottleneck in most systems. Database replication protocols have been proposed to solve the scalability problem by scaling out in a cluster of sites. Current techniques have attained some degree of scalability, however there are two main limitations to existing approaches. Firstly, most solutions adopt a full replication model where all sites store a full copy of the database. The coordination overhead imposed by keeping all replicas consistent allows such approaches to achieve only medium scalability. Secondly, most replication protocols rely on the traditional consistency criterion, 1-copy-serializability, which limits concurrency, and thus scalability of the system. In this paper, we first analyze analytically the performance gains that can be achieved by various partial replication configurations, i.e., configurations where not all sites store all data. From there, we derive a partial replication protocol that provides 1-copy-snapshot isolation as correctness criterion. We have evaluated the protocol with TPC-W and the results show better scalability than full replication.
Damián Serrano, Marta Patiño-Martínez, Ricardo Jiménez-Peris, Bettina Kemme
PRDC4
2007 Enhancing Edge Computing with Database Replication
abstract
As the use of the Internet continues to grow explosively, edge computing has emerged as an important technique for delivering Web content over the Internet. Edge computing moves data and computation closer to end-users for fast local access and better load distribution. Current approaches use caching, which does not work well with highly dynamic data. In this paper, we propose a different approach to enhance edge computing. Our approach lies in a wide area data replication protocol that enables the delivery of dynamic content with full consistency guarantees and with all the benefits of edge computing, such as low latency and high scalability. What is more, the proposed solution is fully transparent to the applications that are brought to the edge. Our extensive evaluations in a real wide area network using TPC-W show promising results.
Yi Lin 0005, Bettina Kemme, Marta Patiño-Martínez, Ricardo Jiménez-Peris
SRDS2
2007 Enterprise Grids: Challenges Ahead
Ricardo Jiménez-Peris, Marta Patiño-Martínez, Bettina Kemme
J. Grid Comput.3
2006 Don't be a Pessimist: Use Snapshot based Concurrency Control for XML
abstract
As native XML database systems (e.g., [3, 7, 8]) get increasingly popular, fine-granularity concurrency control becomes imperative in order to allow different clients to concurrently access the same documents. Existing concurrency control approaches for XML are mainly based on locking [2, 3, 4, 6, 5]. However, the experiments of [5] have shown that the locking overhead, especially for read operations, can be tremendous. In this paper, we present two snapshot based concurrency control mechanisms that avoid locking. Instead, transactions access a committed snapshot of the data.
Zeeshan Sardar, Bettina Kemme
ICDE2
2006 Data Mining Using Relational Database Management Systems
Beibei Zou, Xuesong Ma, Bettina Kemme, Glen Newton, Doina Precup
PAKDD3
2006 Lightweight Reflection for Middleware-based Database Replication
abstract
Middleware-based database replication approaches have emerged in the last few years as an alternative to traditional database replication implemented within the database kernel. A middleware approach enables third party vendors to provide high availability solutions, a growing practice nowadays in the software industry. However, middleware solutions often lack scalability and exhibit a number of consistency and performance issues. The reason is that in most cases the middleware has to handle the database as a black box, and hence, cannot take advantage of the many optimizations implemented in the database kernel. Thus, middleware solutions often reimplement key functionality but cannot achieve the same efficiency as a kernel implementation. Reflection has been proposed during the last decade as a fruitful paradigm to separate non-functional aspects from functional ones, simplifying software development and maintenance whilst fostering reuse. However, fully reflective databases are not feasible due to the high cost of reflection. Our claim is that by exposing some minimal database functionality through a lightweight reflective interface, efficient and scalable middleware database replication can be attained. In this paper we explore a wide variety of such lightweight reflective interfaces and discuss what kind of replication algorithms they enable. We also discuss implementation alternatives for some of these interfaces and evaluate their performance
Jorge Salas, Ricardo Jiménez-Peris, Marta Patiño-Martínez, Bettina Kemme
SRDS4
2005 Consistent Data Replication: Is It Feasible in WANs?
Yi Lin 0005, Bettina Kemme, Marta Patiño-Martínez, Ricardo Jiménez-Peris
Euro-Par2
2005 Postgres-R(SI): Combining Replica Control with Concurrency Control based on Snapshot Isolation
abstract
Replicating data over a cluster of workstations is a powerful tool to increase performance, and provide fault-tolerance for demanding database applications. The big challenge in such systems is to combine replica control (keeping the copies consistent) with concurrency control. Most of the research so far has focused on providing the traditional correctness criteria serializability. However, more and more database systems, e.g., Oracle and PostgreSQL, use multi-version concurrency control providing the isolation level snapshot isolation. In this paper, we present Postgres-R(SI), an extension of PostgreSQL offering transparent replication. Our replication tool is designed to work smoothly with PostgreSQL's concurrency control providing snapshot isolation for the entire replicated system. We present a detailed description of the replica control algorithm, and how it is combined with PostgreSQL's concurrency control component. Furthermore, we discuss some challenges we encountered when implementing the protocol. Our performance analysis based on the TPC-W benchmark shows that this approach exhibits excellent performance for real-life applications even if they are update intensive.
Shuqing Wu, Bettina Kemme
ICDE2
2005 Fine-Granularity Access Control in 3-Tier Laboratory Information Systems
abstract
Laboratory information systems (LIMS) are used in life science research to manage complex experiments. Since LIMS systems are often shared by different research groups, powerful access control is needed to allow different access rights to different records of the same table. Traditional access control models that define a permission as the right of a user/role to perform a specific operation on a specific object cannot handle the enormous amount of objects and user/roles. In this paper, we propose an enhancement to role-based access control by introducing conditions that can be added to the traditional concept of permissions in order to keep the number of permissions small. Furthermore, we present an implementation of our access control model at the application programming level. Although access control is performed for every single database access, our solution completely separates access control from the application logic by using aspect-oriented programming. With this, access control can be integrated into a legacy 3-tier information system without changing the application programs.
Xueli Li, Nomair A. Naeem, Bettina Kemme
IDEAS3
2005 Middleware based Data Replication providing Snapshot Isolation
abstract
Many cluster based replication solutions have been proposed providing scalability and fault-tolerance. Many of these solutions perform replica control in a middleware on top of the database replicas. In such a setting concurrency control is a challenge and is often performed on a table basis. Additionally, some systems put severe requirements on transaction programs (e.g., to declare all objects to be accessed in advance). This paper addresses these issues and presents a middleware-based replication scheme which provides the popular snapshot isolation level at the same tuple-level granularity as database systems like PostgreSQL and Oracle, without any need to declare transaction properties in advance. Both read-only and update transactions can be executed at any replica while providing data consistency at all times. Our approach provides what we call "1-copy-snapshot-isolation" as long as the underlying database replicas provide snapshot isolation. We have implemented our approach as a replicated middleware on top of PostgreSQL replicas. By providing a standard JDBC interface, the middleware is completely transparent to the client program. Fault-tolerance is provided by automatically reconnecting clients in case of crashes. Our middleware shows good performance in terms of response times and scalability.
Yi Lin 0005, Bettina Kemme, Marta Patiño-Martínez, Ricardo Jiménez-Peris
SIGMOD Conference2
2005 Fault-tolerance for Stateful Application Servers in the Presence of Advanced Transactions Patterns
abstract
Replication is widely used in application server products to tolerate faults. An important challenge is to correctly coordinate replication and transaction execution for stateful application servers. Many current solutions assume that a single client request generates exactly one transaction at the server. However, it is quite common that several client requests are encapsulated within one server transaction or that a single client request can initiate several server transactions. In this paper, we propose a replication tool that is able to handle these variations in request/transaction association. We have integrated our approach into the J2EE application server JBoss. Our evaluation using the ECPerf benchmark shows a low overhead of the approach.
Huaigu Wu, Bettina Kemme
SRDS2
2005 MIDDLE-R: Consistent database replication at the middleware level
abstract
The widespread use of clusters and Web farms has increased the importance of data replication. In this article, we show how to implement consistent and scalable data replication at the middleware level. We do this by combining transactional concurrency control with group communication primitives. The article presents different replication protocols, argues their correctness, describes their implementation as part of a generic middleware, Middle-R, and proves their feasibility with an extensive performance evaluation. The solution proposed is well suited for a variety of applications including Web farms and distributed object platforms.
Marta Patiño-Martínez, Ricardo Jiménez-Peris, Bettina Kemme, Gustavo Alonso
ACM Trans. Comput. Syst.3
2004 Comparison of UDDI Registry Replication Strategies
abstract
UDDI registries are intended to become the world-wide lookup mechanism for Web-services. As such, the registry has to provide high throughput, low response times, high availability, and access to accurate data. Replication is often used to satisfy such requirements. Various replication strategies exist, favoring different subsets of the above performance metrics. In this paper we have a closer look at two very different replication strategies. One strategy follows the UDDI specification, the second uses a middleware based replication tool. In this paper, we provide a comparison of these two approaches focusing on performance and ease of integration with an existing UDDI implementation.
Chenliang Sun, Yi Lin 0005, Bettina Kemme
ICWS3
2004 Adaptive Middleware for Data Replication
Jesús M. Milán-Franco, Ricardo Jiménez-Peris, Marta Patiño-Martínez, Bettina Kemme
Middleware4
2003 Using Optimistic Atomic Broadcast in Transaction Processing Systems
abstract
Atomic broadcast primitives are often proposed as a mechanism to allow fault-tolerant cooperation between sites in a distributed system. Unfortunately, the delay incurred before a message can be delivered makes it difficult to implement high performance, scalable applications on top of atomic broadcast primitives. Recently, a new approach has been proposed for atomic broadcast which, based on optimistic assumptions about the communication system, reduces the average delay for message delivery to the application. We develop this idea further and show how applications can take even more advantage of the optimistic assumption by overlapping the coordination phase of the atomic broadcast algorithm with the processing of delivered messages. In particular, we present a replicated database architecture that employs the new atomic broadcast primitive in such a way that communication and transaction processing are fully overlapped, providing high performance without relaxing transaction correctness.
Bettina Kemme, Fernando Pedone, Gustavo Alonso, André Schiper, Matthias Wiesmann
IEEE Trans. Knowl. Data Eng.1
2003 Are quorums an alternative for data replication?
abstract
Data replication is playing an increasingly important role in the design of parallel information systems. In particular, the widespread use of cluster architectures often requires to replicate data for performance and availability reasons. However, maintaining the consistency of the different replicas is known to cause severe scalability problems. To address this limitation, quorums are often suggested as a way to reduce the overall overhead of replication. In this article, we analyze several quorum types in order to better understand their behavior in practice. The results obtained challenge many of the assumptions behind quorum based replication. Our evaluation indicates that the conventional read-one/write-all-available approach is the best choice for a large range of applications requiring data replication. We believe this is an important result for anybody developing code for computing clusters as the read-one/write-all-available strategy is much simpler to implement and more flexible than quorum-based approaches. In this article, we show that, in addition, it is also the best choice using a number of other selection criteria.
Ricardo Jiménez-Peris, Marta Patiño-Martínez, Gustavo Alonso, Bettina Kemme
ACM Trans. Database Syst.4
2002 Improving the Scalability of Fault-Tolerant Database Clusters
abstract
Replication has become a central element in modem information systems playing a dual role: increase availability and enhance scalability. Unfortunately, most existing protocols increase availability at the cost of scalability; This paper presents architecture, implementation and performance of a middleware based replication tool that provides both availability and better scalability than existing systems. Main characteristics are the usage of specialized broadcast primitives and efficient data propagation.
Ricardo Jiménez-Peris, Marta Patiño-Martínez, Bettina Kemme, Gustavo Alonso
ICDCS3
2001 Online Reconfiguration in Replicated Databases Based on Group Communication
abstract
Over the last years, many replica control protocols have been developed that take advantage of the ordering and reliability semantics of group communication primitives to simplify database system design and to improve performance. Although current solutions are able to mask site failures effectively, many of them are unable to cope with recovery of failed sites, merging of partitions, or joining of new sites. This paper addresses this important issue. It proposes efficient solutions for online system reconfiguration providing new sites with a current state of the database without interrupting transaction processing in the rest of the system. Furthermore, the paper analyzes the impact of cascading reconfigurations, and argues that they call be handled in an elegant way by extended forms of group communication.
Bettina Kemme, Alberto Bartoli, Özalp Babaoglu
DSN1
2001 How to Select a Replication Protocol According to Scalability, Availability, and Communication Overhead
abstract
Data replication is playing an increasingly important role in the design of parallel information systems. In particular, the widespread use of cluster architectures in high-performance computing has created many opportunities for applying data replication techniques in new areas. For instance, as part of work related to cluster computing in bioinformatics, we have been confronted with the problem of having to choose an optimal replication strategy in terms of scalability, availability and communication overhead. Thus, we have evaluated several representative replication protocols in order to better understand their behavior in practice. The results obtained are surprising in that they challenge many of the assumptions behind existing protocols. Our evaluation indicates that the conventional read-one/write-all approach is the best choice for a large range of applications requiring data replication. We believe this is an important result for anybody developing code for computing clusters as the read-one/write-all strategy is much simpler to implement and more flexible than quorum-based approaches. In this paper we show that, in addition, it is also the best choice using a number of other selection criteria.
Ricardo Jiménez-Peris, Marta Patiño-Martínez, Bettina Kemme, Gustavo Alonso
SRDS3
2000 Understanding Replication in Databases and Distributed Systems
abstract
Replication is an area of interest to both distributed systems and databases. The solutions developed from these two perspectives are conceptually similar but differ in many aspects: model, assumptions, mechanisms, guarantees provided, and implementation. In this paper, we provide an abstract and "neutral" framework to compare replication techniques from both communities. The framework has been designed to emphasize the role played by different mechanisms and to facilitate comparisons. The paper describes the replication techniques used in both communities, compares them, and points out ways in which they can be integrated to arrive to better, more robust replication protocols.
Fernando Pedone, Matthias Wiesmann, André Schiper, Bettina Kemme, Gustavo Alonso
ICDCS4
2000 Database Replication Techniques: A Three Parameter Classification
abstract
Data replication is an increasingly important topic as databases are more and more deployed over clusters of workstations. One of the challenges in database replication is to introduce replication without severely affecting performance. Because of this difficulty, current database products use lazy replication, which is very efficient but can compromise consistency. As an alternative, eager replication guarantees consistency but most existing protocols have a prohibitive cost. In order to clarify the current state of the art and open up new avenues for research, this paper analyses existing eager techniques using three key parameters (server architecture, server interaction and transaction termination). In our analysis, we distinguish eight classes of eager replication protocols and, for each category, discuss its requirements, capabilities and cost. The contribution lies in showing when eager replication is feasible and in spelling out the different aspects a database replication protocol must account for.
Matthias Wiesmann, André Schiper, Fernando Pedone, Bettina Kemme, Gustavo Alonso
SRDS4
2000 Don't Be Lazy, Be Consistent: Postgres-R, A New Way to Implement Database Replication
Bettina Kemme, Gustavo Alonso
VLDB1
2000 Scalable Replication in Database Clusters
Marta Patiño-Martínez, Ricardo Jiménez-Peris, Bettina Kemme, Gustavo Alonso
DISC3
2000 A new approach to developing and implementing eager database replication protocols
abstract
Database replication is traditionally seen as a way to increase the availability and performance of distributed databases. Although a large number of protocols providing data consistency and fault-tolerance have been proposed, few of these ideas have ever been used in commercial products due to their complexity and performance implications. Instead, current products allow inconsistencies and often resort to centralized approaches which eliminates some of the advantages of replication. As an alternative, we propose a suite of replication protocols that addresses the main problems related to database replication. On the one hand, our protocols maintain data consistency and the same transactional semantics found in centralized systems. On the other hand, they provide flexibility and reasonable performance. To do so, our protocols take advantage of the rich semantics of group communication primitives and the relaxed isolation guarantees provided by most databases. This allows us to eliminate the possibility of deadlocks, reduce the message overhead and increase performance. A detailed simulation study shows the feasibility of the approach and the flexibility with which different types of bottlenecks can be circumvented.
Bettina Kemme, Gustavo Alonso
ACM Trans. Database Syst.1
1999 Processing Transactions over Optimistic Atomic Broadcast Protocols
abstract
Atomic broadcast primitives allow fault-tolerant cooperation between sites in a distributed system. Unfortunately, the delay incurred before a message can be delivered makes it difficult to implement high performance, scalable applications on top of atomic broadcast primitives. A new approach has been proposed which, based on optimistic assumptions about the communication system, reduces the average delay for message delivery. We develop this idea further and present a replicated database architecture that employs the new atomic broadcast primitive in such a way that the coordination phase of the atomic broadcast is fully overlapped with the execution of transactions, providing high performance without relaxing transaction correctness.
Bettina Kemme, Fernando Pedone, Gustavo Alonso, André Schiper
ICDCS1
1998 A Suite of Database Replication Protocols based on Group Communication Primitives
abstract
This paper proposes a family of replication protocols based on group communication in order to address some of the concerns expressed by database designers regarding existing replication solutions. Due to these concerns, current database systems allow inconsistencies and often resort to centralized approaches, thereby reducing some of the key advantages provided by replication. The protocols presented in this paper take advantage of the semantics of group communication and use related isolation guarantees to eliminate the possibility of deadlocks, reduce the message overhead, and increase performance. A simulation study shows the feasibility of the approach and the flexibility with which different types of bottlenecks can be circumvented.
Bettina Kemme, Gustavo Alonso
ICDCS1