1 of 29

NO COMPROMISES: DISTRIBUTED TRANSACTIONS WITH CONSISTENCY, AVAILABILITY, AND PERFORMANCE

Members:

This presentation uses a free template provided by FPPT.com

www.free-power-point-templates.com

2 of 29

Introduction

  • Transactions with strong consistency and high availability enhance distributed systems.
  • They provide simple and powerful abstraction that simplifies distributed systems.
  • Earlier efforts to implement the abstraction led to poor performance and compromise.
  • However, the new FaRM software in modern data centers eliminates the need to compromise.
  • This presentation reviews an article on FaRM software to describe how it eliminates the need to compromise.

This presentation uses a free template provided by FPPT.com

www.free-power-point-templates.com

3 of 29

Background

  • FaRM is a main memory distributed computing platform.
  • It provides transaction, replication and recovery protocols that eliminates the need to compromise.
  • FaRM offers distributed ACID transactions with high availability, strong consistency, high throughput and low latency.
  • FaRM protocols were developed to:
    • Facilitate fast commodity networks with remote direct memory access (RDMA).
    • Offer cheaper approach for non-volatile Dynamic-RAM (DRAM).

This presentation uses a free template provided by FPPT.com

www.free-power-point-templates.com

4 of 29

Background Contd..

  • FaRM’s protocols reduce message counts, use one-sided RDMA reads and writes and exploit parallelism maximumly.
  • They eliminate and expose storage, network & CPU bottlenecks.
  • FaRM scales out distributed systems in data center by allowing transactions to span across numerous machines.
  • It uses precise membership and reservations to ensure high availability and strong consistency.
  • The fast failure recovery protocol help FaRM recover within 50 ms.
  • FaRM thus provides consistency, high availability, and performance.

This presentation uses a free template provided by FPPT.com

www.free-power-point-templates.com

5 of 29

Hardware Trends of FaRM

  • The plentiful availability of cheap DRAM in data center machines encourages FaRM’s design.
  • Data center’s DRAM configuration is 128-512 GB/ 2-socket machines.
  • The DRAM costs below $12/GB, making FaRM design affordable.
  • 2000 machines are enough for a petabyte of DRAM & applications.
  • To eliminate storage and network bottlenecks, FaRM exploits:
    • Non-volatile DRAM
    • Fast commodity networks with RDMA

This presentation uses a free template provided by FPPT.com

www.free-power-point-templates.com

6 of 29

Non-Volatile DRAM and RDMA Networking

  • A distributed uninterruptible power supply (D-UPS) lowers the energy cost in data center with its Lithium-ion batteries.
  • Also makes DRAM durable by saving data to SSD in case of power failure.
  • D-UPS extends the SSD lifetime by writing on it during failure.
  • Less energy is needed to copy data from DRAM to SSD.
  • Treating all machine memory as NVRAM is feasible & cheaper.
  • FaRM uses one-sided RDMA since they do not use remote CPU.

This presentation uses a free template provided by FPPT.com

www.free-power-point-templates.com

7 of 29

Programming Model and Architecture

  • FaRM offers abstraction of global address space to applications.
  • This spans machines in a cluster and improve address storage.
  • FaRM API offers access to remote and local transaction objects.
  • At the end, FaRM commits the transaction.
  • It also offers strict serializability of all successful transactions.
  • FaRM API offers lock-free reads that allow programmers to locate related objects on common machines.
  • Zookeeper coordination service of FaRM streamlines the machines’ current configuration and storage.

This presentation uses a free template provided by FPPT.com

www.free-power-point-templates.com

8 of 29

FaRM Architecture

This presentation uses a free template provided by FPPT.com

www.free-power-point-templates.com

9 of 29

Programming Model and Architecture

  • Replication of data occurs on back-ups to attain consistency.
  • The available 2 GB regions provide up to 250 regions.
  • This allows a single CM handles region allocation for thousands of machines.
  • The centralized model enhances flexibility to satisfy failure independence and locality constraints in machines.
  • The ring buffers in machines serve as transaction logs or message queues.

This presentation uses a free template provided by FPPT.com

www.free-power-point-templates.com

10 of 29

Distributed Transactions and Replication

  • FaRM integrates transaction and replication protocols to enhance performance.
  • Uses fewer messages, promotes CPU efficiency and lo latency.
  • Primary-backup in NV-DRAM help FaRM achieve replication of data and transaction logs.
  • FaRM also uses unreplicated transaction coordinators to improve the consistency.

This presentation uses a free template provided by FPPT.com

www.free-power-point-templates.com

11 of 29

Distributed Transactions and Replication

FaRM commit protocol with a coordinator C, primaries on P1, P2, P3, and backups on B1, B2, B3. P1 and P2 are read and written. P3 is only read. We use dashed lines for RDMA reads, solid ones for RDMA writes, dotted ones for hardware acks, and rectangles for object data.

This presentation uses a free template provided by FPPT.com

www.free-power-point-templates.com

12 of 29

Failure Recovery

  • FaRM’s replication offers durability and high availability.
  • Assumption: machine can fail by crashing but can recover without losing NV-DRAM contents.
  • Bounded clock drift and message delay provides safety and liveness.
  • Durability for all committed transactions is guaranteed even if the entire cluster fails or loses power.
  • Recovery phases: failure detection, reconfiguration, transaction state recovery, bulk data recovery, and allocator state recovery.

This presentation uses a free template provided by FPPT.com

www.free-power-point-templates.com

13 of 29

a) Failure Detection

  • FaRM relies on leases to detect failures.
  • All machines hold a lease at the CM and the CM holds a lease at every machine.
  • Expiry of any lease triggers the failure recovery.
  • In every 1/5 of lease expiry period, lease renewal is attempted.
  • FaRM uses dedicated lease manager thread to improve the recovery speed and message latency.
  • All memory used by lease manager is preallocated during initialization to avoid delay.

This presentation uses a free template provided by FPPT.com

www.free-power-point-templates.com

14 of 29

b) Reconfiguration

Reconfiguration protocol move the FaRM instance from one configuration to another as shown.

This presentation uses a free template provided by FPPT.com

www.free-power-point-templates.com

15 of 29

c) Transaction State Recovery

Transaction state recovery showing a coordinator C, primary P, and two backups B1 and B2.

This presentation uses a free template provided by FPPT.com

www.free-power-point-templates.com

16 of 29

d) Recovering Data

  • FaRM recovers data at new backups for a region.
  • This promotes toleration of f replica failures in future.
  • Data recovery is usually delayed to until all regions are activated.
  • This minimizes impact on latency-critical lock recovery.
  • Active machines send “REGIONS-ACTIVE” message to the CM.
  • CM sends “ALL-REGIONS-ACTIVE” to all machines in configuration.
  • FaRM then begins data recovery for new backups in parallelism.
  • Recovered data are evaluated for normalcy before copying.

This presentation uses a free template provided by FPPT.com

www.free-power-point-templates.com

17 of 29

e) Recovering Allocator State

  • FaRM allocator splits regions into blocks used as slabs for allocating small objects in a transaction.
  • It stores block headers (object size) and slab free lists.
  • Block headers are replicated on allocated new block to ensure availability on new primary after the failure.
  • The new primary sends block headers to all backups to ensure consistency if failure occurs during replication.
  • Slab free lists are stored in the primary only to reduce overheads of object allocation.
  • Free lists are scanned and recovered on new primary after failure.

This presentation uses a free template provided by FPPT.com

www.free-power-point-templates.com

18 of 29

Evaluation

Setup

  • Testbed consisted of 90 machines used for FaRM cluster and 5 for a replicated Zookeeper instance.
  • Specifications: 256 GB DRAM and two 8-core Intel E5-2650 CPUs running Windows Server 2012 R2.
  • Hyper-threading enabled; 30 - foreground work & 2 -lease manager.
  • Machines have 2 Mellanox ConnectX-3 56 Gbps Infiniband NICs used in different sockets; connected by a single Mellanox SX6512 switch.
  • FaRM configuration: 3-way replication and 10 ms lease time.

This presentation uses a free template provided by FPPT.com

www.free-power-point-templates.com

19 of 29

Evaluation

Benchmarks

  • TATP and TPC-C benchmarks used to measure FaRM’s performance.
  • A database with 9.2 billion subscribers was used.
  • TATP was not partitioned, and most operations accessed data on remote machines.
  • With TPC-C, a 16-index schema was used (12 needed point queries and updates while 4 needed range queries).
  • Indexes implemented using FaRM tree; 8 GB/machine cache reserve.
  • Database with 21,600 warehouses used.

This presentation uses a free template provided by FPPT.com

www.free-power-point-templates.com

20 of 29

Evaluation: Normal-case performance

a) TATP performance b)TPC-C performance

This presentation uses a free template provided by FPPT.com

www.free-power-point-templates.com

21 of 29

Evaluation: Failures

  • The FaRM process in one machine was killed 35s into experiment.
  • Timeline with throughput of 89 running machines aggregated at 1 ms intervals.
  • RDMA messaging was used to synchronize the experiment at start.

TATP performance timeline with CM failure

This presentation uses a free template provided by FPPT.com

www.free-power-point-templates.com

22 of 29

Evaluation: Failures

TATP performance timeline with failure

This presentation uses a free template provided by FPPT.com

www.free-power-point-templates.com

23 of 29

Evaluation: Failures

Correlated Failures

This presentation uses a free template provided by FPPT.com

www.free-power-point-templates.com

24 of 29

Evaluation: Failures

Data recovery pacing

This presentation uses a free template provided by FPPT.com

www.free-power-point-templates.com

25 of 29

Evaluation: Failures

Data recovery pacing

  • FaRM paces data recovery to minimize impact on the throughput.
  • This optimizes the time for completing re-replication of regions at new backups.
  • TPC-C is less sensitive to interference than TATP.

This presentation uses a free template provided by FPPT.com

www.free-power-point-templates.com

26 of 29

Evaluation: Lease Times

  • To evaluate the lease manger optimizations, a test was ran.
  • Recovery was disabled and the number of lease expiry events counted across the cluster.
  • Aim: lease manager implementations and lease durations.
  • All optimizations are needed to allow lease times of 10 ms or less without false positives.
  • Shared queue pairs allow frequent expiry of up to 100 ms.
  • Unreliable datagrams reduce the no. of false positives but does not eliminate it due to CPU contention.

This presentation uses a free template provided by FPPT.com

www.free-power-point-templates.com

27 of 29

Related Work

  • FaRM is the first system to provide high availability, high throughput, low latency, and strict serializability simultaneously.
  • Optimized protocol sends up to 44% lesser messages than the transaction protocol.
  • FaRM supports transactions and enhances data availability through quick recovery (tens of milliseconds) from a failure.
  • It also has an order of magnitude higher throughput per machine.
  • FaRM offers strict serializability and general distributed transactions optimized to take advantage of RDMA.

This presentation uses a free template provided by FPPT.com

www.free-power-point-templates.com

28 of 29

Conclusion

  • FaRM is a reliable D-RAM that provides high availability, high throughput, low latency, and strict serializability simultaneously.
  • The new transaction, replication and recovery protocols designed from first principles are essential in attaining the above aspects.
  • FaRM is more effective and efficient in providing higher throughput and lower latency than in-memory databases.
  • It also facilitates quick recovery from a machine failure to peak throughput in less than 50 ms.
  • The quick recovery makes failure transparent to applications.

This presentation uses a free template provided by FPPT.com

www.free-power-point-templates.com

29 of 29

Reference

  • Dragojević, A., Narayanan, D., Nightingale, E. B., Renzelmann, M., Shamis, A., Badam, A., & Castro, M. (2015, October). No compromises: distributed transactions with consistency, availability, and performance. In Proceedings of the 25th symposium on operating systems principles (pp. 54-70).

This presentation uses a free template provided by FPPT.com

www.free-power-point-templates.com