K. Psarakis
info
Please Note
<p>This page displays the records of the person named above and is not linked to a unique person identifier. This record may need to be merged to a profile.</p>
7 records found
1
Uncoordinated Checkpointing in Stateful Transactional Systems
Decoupling Fault Tolerance from Coordination in Styx
Styx is a distributed runtime for stateful transactional functions. It executes transactions in deterministic, lockstep epochs and periodically writes state to
stable storage for recovery. In the original design, each checkpoint depends on both the workers’ state files and a coordinator-written global sequencer file.
Consequently, a delayed or failed coordinator write can make otherwise complete worker checkpoints unusable. This thesis investigates whether checkpoint
persistence can be moved from the coordinator to the workers without harming system performance or recovery.
The proposed design uses Styx’s shared epoch boundaries as consistent recovery points. Workers persist their partition checkpoints independently, while
recovery selects the newest complete epoch across all partitions. Transaction execution remains coordinated; only checkpoint persistence is decentralized. The
evaluation considers steady-state performance, checkpoint-file compaction, and partition rebalancing during degraded operation after a worker failure.
Across approximately 1,670 single-node runs and 334 configurations, coordinated and uncoordinated checkpoint persistence show no systematic performance
difference. Their saturation points differ by -5.3% to +6.3%, with no consistent winner. Recovery performance depends more strongly on how stored checkpoints
are managed. Compacting every 100 checkpoints reduces state restoration time from 4,971 ms to 474 ms and provides the best tested balance between recovery
speed and foreground latency. More aggressive compaction restores state faster but becomes disruptive as load increases.
Rebalancing is similarly workload-dependent. After a worker failure, the surviving workers operate with an uneven partition distribution. Redistributing those
partitions reduces P95 latency from 288 ms to 146 ms at 40% utilization without affecting throughput. At loads of 50% and above, however, the newly restored
worker limits epoch progress and lowers throughput during the 90-second observation window.
These results show that checkpoint persistence can be decentralized without measurable steady-state cost under the tested conditions. Compaction should balance
recovery objectives against foreground load, while post-failure rebalancing should be applied selectively rather than as a universal recovery rule. The
findings are limited to one workload and a single-node deployment.
Related dataset 4TU.ResearchData: https://doi.org/10.4121/a18e4390-8d18-4a2a-883c-0fe6a103bda1 ...
stable storage for recovery. In the original design, each checkpoint depends on both the workers’ state files and a coordinator-written global sequencer file.
Consequently, a delayed or failed coordinator write can make otherwise complete worker checkpoints unusable. This thesis investigates whether checkpoint
persistence can be moved from the coordinator to the workers without harming system performance or recovery.
The proposed design uses Styx’s shared epoch boundaries as consistent recovery points. Workers persist their partition checkpoints independently, while
recovery selects the newest complete epoch across all partitions. Transaction execution remains coordinated; only checkpoint persistence is decentralized. The
evaluation considers steady-state performance, checkpoint-file compaction, and partition rebalancing during degraded operation after a worker failure.
Across approximately 1,670 single-node runs and 334 configurations, coordinated and uncoordinated checkpoint persistence show no systematic performance
difference. Their saturation points differ by -5.3% to +6.3%, with no consistent winner. Recovery performance depends more strongly on how stored checkpoints
are managed. Compacting every 100 checkpoints reduces state restoration time from 4,971 ms to 474 ms and provides the best tested balance between recovery
speed and foreground latency. More aggressive compaction restores state faster but becomes disruptive as load increases.
Rebalancing is similarly workload-dependent. After a worker failure, the surviving workers operate with an uneven partition distribution. Redistributing those
partitions reduces P95 latency from 288 ms to 146 ms at 40% utilization without affecting throughput. At loads of 50% and above, however, the newly restored
worker limits epoch progress and lowers throughput during the 90-second observation window.
These results show that checkpoint persistence can be decentralized without measurable steady-state cost under the tested conditions. Compaction should balance
recovery objectives against foreground load, while post-failure rebalancing should be applied selectively rather than as a universal recovery rule. The
findings are limited to one workload and a single-node deployment.
Related dataset 4TU.ResearchData: https://doi.org/10.4121/a18e4390-8d18-4a2a-883c-0fe6a103bda1 ...
Styx is a distributed runtime for stateful transactional functions. It executes transactions in deterministic, lockstep epochs and periodically writes state to
stable storage for recovery. In the original design, each checkpoint depends on both the workers’ state files and a coordinator-written global sequencer file.
Consequently, a delayed or failed coordinator write can make otherwise complete worker checkpoints unusable. This thesis investigates whether checkpoint
persistence can be moved from the coordinator to the workers without harming system performance or recovery.
The proposed design uses Styx’s shared epoch boundaries as consistent recovery points. Workers persist their partition checkpoints independently, while
recovery selects the newest complete epoch across all partitions. Transaction execution remains coordinated; only checkpoint persistence is decentralized. The
evaluation considers steady-state performance, checkpoint-file compaction, and partition rebalancing during degraded operation after a worker failure.
Across approximately 1,670 single-node runs and 334 configurations, coordinated and uncoordinated checkpoint persistence show no systematic performance
difference. Their saturation points differ by -5.3% to +6.3%, with no consistent winner. Recovery performance depends more strongly on how stored checkpoints
are managed. Compacting every 100 checkpoints reduces state restoration time from 4,971 ms to 474 ms and provides the best tested balance between recovery
speed and foreground latency. More aggressive compaction restores state faster but becomes disruptive as load increases.
Rebalancing is similarly workload-dependent. After a worker failure, the surviving workers operate with an uneven partition distribution. Redistributing those
partitions reduces P95 latency from 288 ms to 146 ms at 40% utilization without affecting throughput. At loads of 50% and above, however, the newly restored
worker limits epoch progress and lowers throughput during the 90-second observation window.
These results show that checkpoint persistence can be decentralized without measurable steady-state cost under the tested conditions. Compaction should balance
recovery objectives against foreground load, while post-failure rebalancing should be applied selectively rather than as a universal recovery rule. The
findings are limited to one workload and a single-node deployment.
Related dataset 4TU.ResearchData: https://doi.org/10.4121/a18e4390-8d18-4a2a-883c-0fe6a103bda1
stable storage for recovery. In the original design, each checkpoint depends on both the workers’ state files and a coordinator-written global sequencer file.
Consequently, a delayed or failed coordinator write can make otherwise complete worker checkpoints unusable. This thesis investigates whether checkpoint
persistence can be moved from the coordinator to the workers without harming system performance or recovery.
The proposed design uses Styx’s shared epoch boundaries as consistent recovery points. Workers persist their partition checkpoints independently, while
recovery selects the newest complete epoch across all partitions. Transaction execution remains coordinated; only checkpoint persistence is decentralized. The
evaluation considers steady-state performance, checkpoint-file compaction, and partition rebalancing during degraded operation after a worker failure.
Across approximately 1,670 single-node runs and 334 configurations, coordinated and uncoordinated checkpoint persistence show no systematic performance
difference. Their saturation points differ by -5.3% to +6.3%, with no consistent winner. Recovery performance depends more strongly on how stored checkpoints
are managed. Compacting every 100 checkpoints reduces state restoration time from 4,971 ms to 474 ms and provides the best tested balance between recovery
speed and foreground latency. More aggressive compaction restores state faster but becomes disruptive as load increases.
Rebalancing is similarly workload-dependent. After a worker failure, the surviving workers operate with an uneven partition distribution. Redistributing those
partitions reduces P95 latency from 288 ms to 146 ms at 40% utilization without affecting throughput. At loads of 50% and above, however, the newly restored
worker limits epoch progress and lowers throughput during the 90-second observation window.
These results show that checkpoint persistence can be decentralized without measurable steady-state cost under the tested conditions. Compaction should balance
recovery objectives against foreground load, while post-failure rebalancing should be applied selectively rather than as a universal recovery rule. The
findings are limited to one workload and a single-node deployment.
Related dataset 4TU.ResearchData: https://doi.org/10.4121/a18e4390-8d18-4a2a-883c-0fe6a103bda1
Master thesis
(2025)
-
Smruti Kshirsagar, A. Katsifodimos, K. Psarakis, G.C. Christodoulou, G. Iosifidis, B. Özkan
Stateful Functions-as-a-Service (SFaaS) platforms, such as Styx, are emerging as powerful abstractions for building distributed, serverless cloud applications. By combining the abilities of FaaS with strong transactional guarantees, they enable complex, stateful workflows without requiring developers to manage infrastructure. However, they lack built-in support for analytical queries across distributed function state. This thesis addresses that gap by proposing H-Styx, whose hybrid architecture extends Styx with a snapshot-based Query Engine, enabling near-real-time OLAP queries over global state while maintaining performance isolation for transactions. The Query Engine integrates seamlessly into the Styx architecture, leveraging periodic snapshots transmitted via a loosely-coupled, asynchronous interface. It ingests partitioned state from object store MinIO into columnar database DuckDB, supports incremental delta loads, and delivers results over a Kafka-based interface to achieve scalable, low-latency analytical querying while employing robust fault tolerance.
Empirical evaluation demonstrates that H-Styx preserves transactional throughput and latency under hybrid workloads, while significantly outperforming a baseline HTAP architecture (Postgres with Streaming Replication) on analytical throughput and providing superior workload isolation. These results validate the feasibility of supporting hybrid transactional and analytical processing in SFaaS environments. Overall, H-Styx bridges a crucial capability gap in SFaaS, enabling more powerful data-driven applications in distributed, event-driven architectures. ...
Empirical evaluation demonstrates that H-Styx preserves transactional throughput and latency under hybrid workloads, while significantly outperforming a baseline HTAP architecture (Postgres with Streaming Replication) on analytical throughput and providing superior workload isolation. These results validate the feasibility of supporting hybrid transactional and analytical processing in SFaaS environments. Overall, H-Styx bridges a crucial capability gap in SFaaS, enabling more powerful data-driven applications in distributed, event-driven architectures. ...
Stateful Functions-as-a-Service (SFaaS) platforms, such as Styx, are emerging as powerful abstractions for building distributed, serverless cloud applications. By combining the abilities of FaaS with strong transactional guarantees, they enable complex, stateful workflows without requiring developers to manage infrastructure. However, they lack built-in support for analytical queries across distributed function state. This thesis addresses that gap by proposing H-Styx, whose hybrid architecture extends Styx with a snapshot-based Query Engine, enabling near-real-time OLAP queries over global state while maintaining performance isolation for transactions. The Query Engine integrates seamlessly into the Styx architecture, leveraging periodic snapshots transmitted via a loosely-coupled, asynchronous interface. It ingests partitioned state from object store MinIO into columnar database DuckDB, supports incremental delta loads, and delivers results over a Kafka-based interface to achieve scalable, low-latency analytical querying while employing robust fault tolerance.
Empirical evaluation demonstrates that H-Styx preserves transactional throughput and latency under hybrid workloads, while significantly outperforming a baseline HTAP architecture (Postgres with Streaming Replication) on analytical throughput and providing superior workload isolation. These results validate the feasibility of supporting hybrid transactional and analytical processing in SFaaS environments. Overall, H-Styx bridges a crucial capability gap in SFaaS, enabling more powerful data-driven applications in distributed, event-driven architectures.
Empirical evaluation demonstrates that H-Styx preserves transactional throughput and latency under hybrid workloads, while significantly outperforming a baseline HTAP architecture (Postgres with Streaming Replication) on analytical throughput and providing superior workload isolation. These results validate the feasibility of supporting hybrid transactional and analytical processing in SFaaS environments. Overall, H-Styx bridges a crucial capability gap in SFaaS, enabling more powerful data-driven applications in distributed, event-driven architectures.
MovR as a Benchmark for Geo-Distributed Databases
Performance Evaluation and Insights
Bachelor thesis
(2025)
-
W.P.A. Marcu, A. Katsifodimos, O. Mráz, G.C. Christodoulou, K. Psarakis, K.G. Langendoen
Distributed systems are vital for handling large-scale data and rely on geo-distributed databases to ensure low latency and high availability. Traditional benchmarks, such as TPC-C and YCSB-T, are not designed to handle the complexities of geo-distributed environments and do not allow for configuration of multi-home transaction ratios or dynamic data access patterns. To fill this gap, we implement a benchmark based on the MovR workload and assess its performance on the Detock, Janus, SLOG, and Calvin geo-distributed database systems. Key insights revealed through experiments are that network conditions act as a major bottleneck and high concurrency leads to unsustainable latency spikes which severely limits scalability.
...
Distributed systems are vital for handling large-scale data and rely on geo-distributed databases to ensure low latency and high availability. Traditional benchmarks, such as TPC-C and YCSB-T, are not designed to handle the complexities of geo-distributed environments and do not allow for configuration of multi-home transaction ratios or dynamic data access patterns. To fill this gap, we implement a benchmark based on the MovR workload and assess its performance on the Detock, Janus, SLOG, and Calvin geo-distributed database systems. Key insights revealed through experiments are that network conditions act as a major bottleneck and high concurrency leads to unsustainable latency spikes which severely limits scalability.
Serverless computing has allowed developers to write pieces of code comprising solely of the necessary functionality whilst not having to think about the underlying infrastructure. One prominent model is Function-as-a-Service (FaaS), where the code is structured into functions that run based on incoming events. This model was initially stateless, as new function calls can be instantiated at any location and there is no clear consensus on state whenever multiple instances are running simultaneously. Access to external persistent state is slow, making FaaS not suitable for low latency applications. Recent works have found different ways of incorporating state, resulting in Stateful FaaS (SFaaS). With the addition of state, these components are perfectly suited for distributed transactions. SFaaS frameworks try to outperform one another on metrics such as throughput and latency, but less work is performed on consistency.
In this thesis we look at the work on consistency of SFaaS transactions that has been done. We take Jepsen, a framework for testing transactions in distributed systems, and show that it can also be applied to SFaaS transactions. We then proceed to apply it to three SFaaS frameworks. We use Elle as a consistency checker to verify that the three frameworks comply with the consistency level they are advertised as. We have found that two of the tested frameworks do not have the consistency level promised. Facets of the SFaaS frameworks seem to have been overlooked, and diverging from the laid out benchmarking path quickly results in unintended behaviour. ...
In this thesis we look at the work on consistency of SFaaS transactions that has been done. We take Jepsen, a framework for testing transactions in distributed systems, and show that it can also be applied to SFaaS transactions. We then proceed to apply it to three SFaaS frameworks. We use Elle as a consistency checker to verify that the three frameworks comply with the consistency level they are advertised as. We have found that two of the tested frameworks do not have the consistency level promised. Facets of the SFaaS frameworks seem to have been overlooked, and diverging from the laid out benchmarking path quickly results in unintended behaviour. ...
Serverless computing has allowed developers to write pieces of code comprising solely of the necessary functionality whilst not having to think about the underlying infrastructure. One prominent model is Function-as-a-Service (FaaS), where the code is structured into functions that run based on incoming events. This model was initially stateless, as new function calls can be instantiated at any location and there is no clear consensus on state whenever multiple instances are running simultaneously. Access to external persistent state is slow, making FaaS not suitable for low latency applications. Recent works have found different ways of incorporating state, resulting in Stateful FaaS (SFaaS). With the addition of state, these components are perfectly suited for distributed transactions. SFaaS frameworks try to outperform one another on metrics such as throughput and latency, but less work is performed on consistency.
In this thesis we look at the work on consistency of SFaaS transactions that has been done. We take Jepsen, a framework for testing transactions in distributed systems, and show that it can also be applied to SFaaS transactions. We then proceed to apply it to three SFaaS frameworks. We use Elle as a consistency checker to verify that the three frameworks comply with the consistency level they are advertised as. We have found that two of the tested frameworks do not have the consistency level promised. Facets of the SFaaS frameworks seem to have been overlooked, and diverging from the laid out benchmarking path quickly results in unintended behaviour.
In this thesis we look at the work on consistency of SFaaS transactions that has been done. We take Jepsen, a framework for testing transactions in distributed systems, and show that it can also be applied to SFaaS transactions. We then proceed to apply it to three SFaaS frameworks. We use Elle as a consistency checker to verify that the three frameworks comply with the consistency level they are advertised as. We have found that two of the tested frameworks do not have the consistency level promised. Facets of the SFaaS frameworks seem to have been overlooked, and diverging from the laid out benchmarking path quickly results in unintended behaviour.
Today's need for highly available systems leads to data partitioning and replication across multiple nodes. Providing strong transactional consistency in a distributed database requires extensive communication. For this, algorithms such as two phase commit are used. These communication algorithms add extra network latency's. For application developers and database systems, this is the reason for lowering the isolation level of a database. Deterministic databases run transactions effectively without communication between replicas. Most deterministic databases need the read write sets of a transaction prior to execution to calculate a deterministic execution schedule. Aria does not need the read write sets a priory but uses an epoch based commit protocol. The commit protocol is an optimistic concurrency control algorithm that executes all transactions against a snapshot in the execution phase and determines which transactions can commit in the commit phase. For most workloads Aria outperforms state of the art deterministic databases. However, for high contention workloads Aria suffers performance because of high abort rates. To overcome this problem this thesis proposes two solutions: 1) Lowering the isolation level to snapshot isolation. 2) Reordering the input sequence of transactions on transaction degree. We have found that lowering the isolation level to snapshot isolation allows for 3\% less aborts per epoch and reduces latency from 210 ms to 170ms. Reordering the transaction sequence allows for 5 percent less aborts per epoch for snapshot isolation and serializable isolation level. Reordering transactions on degree for serializable isolation level reduces the average latency from 210 ms to 150 ms. Snapshot isolation reordering transactions on degree reduces the average latency from 170 ms to 120 ms.
...
...
Today's need for highly available systems leads to data partitioning and replication across multiple nodes. Providing strong transactional consistency in a distributed database requires extensive communication. For this, algorithms such as two phase commit are used. These communication algorithms add extra network latency's. For application developers and database systems, this is the reason for lowering the isolation level of a database. Deterministic databases run transactions effectively without communication between replicas. Most deterministic databases need the read write sets of a transaction prior to execution to calculate a deterministic execution schedule. Aria does not need the read write sets a priory but uses an epoch based commit protocol. The commit protocol is an optimistic concurrency control algorithm that executes all transactions against a snapshot in the execution phase and determines which transactions can commit in the commit phase. For most workloads Aria outperforms state of the art deterministic databases. However, for high contention workloads Aria suffers performance because of high abort rates. To overcome this problem this thesis proposes two solutions: 1) Lowering the isolation level to snapshot isolation. 2) Reordering the input sequence of transactions on transaction degree. We have found that lowering the isolation level to snapshot isolation allows for 3\% less aborts per epoch and reduces latency from 210 ms to 170ms. Reordering the transaction sequence allows for 5 percent less aborts per epoch for snapshot isolation and serializable isolation level. Reordering transactions on degree for serializable isolation level reduces the average latency from 210 ms to 150 ms. Snapshot isolation reordering transactions on degree reduces the average latency from 170 ms to 120 ms.
Serverless computing is an increasingly popular paradigm in cloud computing where many of the operational challenges of running cloud applications, like server provi- sioning and management, are left to the cloud provider. A popular form of server- less computing is Functions-as-a-Service (FaaS), where the user submits functions for which the resources are automatically provisioned and scaled. However, FaaS functions are traditionally stateless, and thus rely on external services to handle state. Stateful FaaS systems are an extension to traditional FaaS offerings which have built-in function state management and function addressability, which allows for functions to communicate and to form complex operations. This work introduces a performance benchmark for stateful functions systems, based on an e-commerce application. The benchmark includes two workflows based on two complex oper- ations that involve multiple stateful functions. The workflows can be dynamically altered to form different function calling structures. We provide reference imple- mentations of the benchmark application for four current stateful FaaS systems, and a benchmark client which runs the two benchmark workloads on these implementa- tions. We show that our benchmark can be used to test and compare stateful systems on performance, networking, cost and scalability, by running various experiments on our reference implementations.
...
Serverless computing is an increasingly popular paradigm in cloud computing where many of the operational challenges of running cloud applications, like server provi- sioning and management, are left to the cloud provider. A popular form of server- less computing is Functions-as-a-Service (FaaS), where the user submits functions for which the resources are automatically provisioned and scaled. However, FaaS functions are traditionally stateless, and thus rely on external services to handle state. Stateful FaaS systems are an extension to traditional FaaS offerings which have built-in function state management and function addressability, which allows for functions to communicate and to form complex operations. This work introduces a performance benchmark for stateful functions systems, based on an e-commerce application. The benchmark includes two workflows based on two complex oper- ations that involve multiple stateful functions. The workflows can be dynamically altered to form different function calling structures. We provide reference imple- mentations of the benchmark application for four current stateful FaaS systems, and a benchmark client which runs the two benchmark workloads on these implementa- tions. We show that our benchmark can be used to test and compare stateful systems on performance, networking, cost and scalability, by running various experiments on our reference implementations.
As serverless computing grows in popularity, developers are demanding more from existing serverless models. One example is the emergence of Stateful Function as a Service (SFaaS), in which state is added to operators in existing Function as a Service (FaaS) models, to support microservice-type applications while utilising the benefits of a serverless architecture. However, current SFaaS systems cannot provide performant transactions across operators with strong semantics. In this thesis, we present three conceptual transaction protocols for SFaaS dataflow systems based on Two-phase Commit (2PC), Deterministic Databases and Conflict-free Replicated Datatype (CRDT)s. Based on these insights, we implement Rhea, a deterministic transaction protocol in the prototype SFaaS execution engine, Universalis. Two optimizations are implemented for Rhea: deterministic reordering and a fallback mechanism, to reduce aborts caused by Read-after-write (RAW) and Write-after-write (WAW) dependencies, respectively. We present a transaction benchmarking client for Universalis that supports workloads from the Transaction Processing Performance Council Benchmark C (TPC-C) and the Yahoo! Cloud Serving Benchmark (YCSB). Using the client, we compare Rhea to a 2PC baseline microservice application. We found that Rhea can provide more than twice the throughput and half the latency of 2PC in the baseline application across a variety of workloads.
...
As serverless computing grows in popularity, developers are demanding more from existing serverless models. One example is the emergence of Stateful Function as a Service (SFaaS), in which state is added to operators in existing Function as a Service (FaaS) models, to support microservice-type applications while utilising the benefits of a serverless architecture. However, current SFaaS systems cannot provide performant transactions across operators with strong semantics. In this thesis, we present three conceptual transaction protocols for SFaaS dataflow systems based on Two-phase Commit (2PC), Deterministic Databases and Conflict-free Replicated Datatype (CRDT)s. Based on these insights, we implement Rhea, a deterministic transaction protocol in the prototype SFaaS execution engine, Universalis. Two optimizations are implemented for Rhea: deterministic reordering and a fallback mechanism, to reduce aborts caused by Read-after-write (RAW) and Write-after-write (WAW) dependencies, respectively. We present a transaction benchmarking client for Universalis that supports workloads from the Transaction Processing Performance Council Benchmark C (TPC-C) and the Yahoo! Cloud Serving Benchmark (YCSB). Using the client, we compare Rhea to a 2PC baseline microservice application. We found that Rhea can provide more than twice the throughput and half the latency of 2PC in the baseline application across a variety of workloads.