Software-Defined Caching in Data Centers
Software-Defined Caching in Data Centers
Ioan Stefanovici⋆ , Eno Thereska, Greg O’Shea, Bianca Schroeder⋆ , Hitesh Ballani, Thomas
Karagiannis, Antony Rowstron, Tom Talpey†
University of Toronto⋆, Microsoft Research, Microsoft†
Abstract VM VM VM VM VM VM
In data centers, caches work both to provide low IO laten-
cies and to reduce the load on the back-end network and Hypervisor Hypervisor Hypervisor
storage. But they are not designed for multi-tenancy; system
level caches today cannot be configured to match tenant or
provider objectives. Exacerbating the problem is the increas-
ing number of un-coordinated caches on the IO data plane. Storage Storage
The lack of global visibility on the control plane to coor-
dinate this distributed set of caches leads to inefficiencies,
increasing cloud provider cost. Figure 1: Simplified IO stack in a multi-tenant data cen-
We present Moirai, a tenant and workload aware system ter. Two tenants, a green and red one are shown, with 3
that allows data center providers to control their distributed VMs each spread over 3 hypervisors. The circles repre-
caching infrastructure. Moirai can help ease the management sent typical caches on the IO stack.
of the cache infrastructure and achieve various objectives,
such as improving overall resource utilization or providing
tenant isolation and QoS guarantees, as we show through • Lack of performance isolation Since caches are not
several use cases. A key benefit of Moirai is that it is trans- tenant- or workload-aware, applications with different IO
parent to applications or VMs deployed in data centers. Our patterns and request rates sharing the same cache will impact
prototype runs unmodified OSes and databases providing each other’s cache performance. For example, depending on
immediate benefit to existing applications. the cache eviction policy, one application’s large sequen-
tial reads can blast away another workload’s working set.
Even with scan-resistant cache management policies, such
1. Introduction as ARC [27], aggressive clients with higher request rates will
An increasing number of enterprise applications have mi- still be allocated larger portions of the cache.
grated to hosted platforms in private enterprise and public • Lack of customization Since caches are not tenant
cloud data centers. Such platforms are typically virtualized, aware, the entire cache is treated as a single pool with one
i.e., tenants deploy applications in virtual machines (VMs) cache write policy (write-through, write-back, etc), despite
whose access to the underlying resources (memory, storage, different durability requirements of different applications,
network) is shared with other tenants, and mediated by hy- and one eviction policy, despite the fact that different work-
pervisors such as Hyper-V, VMware ESX, or Xen. Uninhib- loads benefit from different cache eviction policies. For ex-
ited sharing of such resources in a multi-tenant environment ample, Figure 2(a) shows two IOMeter workloads under two
leads to poor and variable application performance. While different eviction policies, LRU and MRU [9] respectively.
recent efforts give providers control over how resources like The workload on the left performs at its peak with an MRU
network [1, 17, 22, 30, 33] and storage [2, 15, 16, 34, 38] policy, while the one on the right performs best with LRU.
are shared, there is no coordinated end-to-end control of the Today, if both workloads were running atop the same hyper-
distributed caching infrastructure, made up of storage caches visor, they would have to follow the same eviction policy,
at multiple places along the IO stack (inside VMs, hypervi- leading to performance penalties on the order of 4-5x.
sors, storage servers; see Figure 1). Today, storage caches • Lack of coordination Each cache in the IO stack makes
along the IO stack are transparent to both applications and its decisions locally, agnostic to the state of other caches in
cloud providers, lack workload-aware mechanisms, and are the stack, leading to inefficiencies, such as double caching,
each managed in isolation, leading to multiple problems: as was also noted by Wong and Wilkes [46].
1
20000 T1 T2 T3 T4
Throughput (MB/s)
18000
16000 ŽŶƚƌŽůWůĂŶĞ sD sD
14000
12000
Naive Paroning ĂƚĂWůĂŶĞ
10000
,LJƉĞƌǀŝƐŽƌ
8000 DĞƚƌŝĐƐ
6000
4000
ŶŐŝŶĞ ŽŵƉƵƚĞ
2000
Op mal Par oning ƐĞƌǀĞƌ
Workload 2
0
ĐĂĐŚĞ
Workload 1
130
173
216
259
302
345
388
431
474
517
560
603
646
689
732
775
818
861
904
947
44
87
Time (s) ŽŶƚƌŽůůĞƌ
ĂĐŚĞĂůůŽĐĂƚŝŽŶ͕
(a) The effect of eviction policy (b) The effect of cache size ĐŽŶĨŝŐƵƌĂƚŝŽŶ͕
ƉůĂĐĞŵĞŶƚ
Figure 2: Performance depends on the cache policy (a) ĐĂĐŚĞ ^ƚŽƌĂŐĞ
and allocation (b). ƐĞƌǀĞƌ
• Lack of adaptability Currently, the organization and Figure 3: The Moirai architecture.
configuration of caches is fixed. Caches cannot be added,
removed, or resized on the fly to adapt to changes in the
workload or in provider objectives. requests serviced from cache as a function of the cache size.
• Waste of system resources Simple solutions for par- We use phantom caches, which inspect IO headers (with
titioning caches along the IO stack are not sufficient. For fields such as accessed file name, offset, length, etc.) and ex-
example, Figure 2(b) shows that the observed performance ploit techniques from recent work [32, 44, 45] to generate hit
triples when cache space is optimally allocated according ratio curves efficiently at runtime. The Metrics Engine peri-
to workload characteristics (the workload consists of 4 ten- odically sends these performance metrics to the centralized
ants using 120 VMs in total), compared to the case when controller.
caches are naively allocated across tenants. We will describe
2.2 Programmable Caches
the details of this experiment in Section 5, but note that all
workloads’ throughputs benefit when the right cache size is Caches along the IO stack are programmable through a sim-
chosen. This is true even for tenants that receive less total ple API shown in Table 1. Caches are created at the desired
cache, as the contention at the storage device is reduced. position in the IO stack by sending a createCache call to
While some of these problems have been tackled in iso- the appropriate level in the stack (more details in Section 4).
lation, there is no comprehensive framework for the end-to- A cache c is made workload-aware using the createRule
end management of caches that allows providers to address call, which installs a rule to specify the IOs that should
the major issues they are facing today. We present Moirai1 , be cached in c. If the header of an incoming I/O matches
a tenant- and workload-aware system that allows data center one of c’s rules, the IO (header+data) is sent through the
providers to control their distributed caching infrastructure cache. The controller can also configure cache properties
to achieve provider objectives, such as improving resource (configureCache) to set the size, eviction, and write poli-
utilization and request latency, achieving tenant isolation and cies. Similarly, cache performance metrics are obtained via
QoS guarantees. Moirai does not require changes to the IO the getCacheStats call.
stack architecture, is transparent to applications and VMs, Care must be taken to maintain consistency semantics
and does not change cache consistency semantics. when the location of a cache changes. For example, the con-
troller could decide to cache at the storage server rather than
2. Design at the hypervisor. In order to maintain consistency, Moirai
first removes the caches on the old path, which automatically
Figure 3 shows the architecture of Moirai, which comprises
triggers the eviction of all cached state, including writing any
three key components. At the core is a logically-centralized
dirty blocks to the back-end storage, and then installs caches
controller that uses information on workload characteris-
on the new path. We considered other options, such as keep-
tics maintained by the metrics engine to configure the pro-
ing the old caches until all accesses eventually move to the
grammable caches to achieve provider objectives. Details on
new caches, but they add complexity and require maintain-
each of the three components are provided next.
ing extra metadata.
2.1 The Metrics Engine
2.3 Controller
The Metrics Engine is a hypervisor-based module that main-
tains key characteristics for each workload running on the The centralized controller uses the API described in Sec-
system, such as throughput, number of reads vs. writes, etc., tion 2.2 and information provided by the Metrics Engine to
but also hit ratio curves, which describe the percentage of create and configure caches in order to implement a set of
objectives specified by the provider, as illustrated in the next
1 Moirai (Ancient Greek for “apportioner”) in Greek mythology are the section.
three personifications of fate, who control the thread of life of every mortal
from birth to death, analogously to the end-to-end control of caches by the
three components that comprise Moirai.
2
createCache (<size,eviction pol,write pol>) tϭ ϭ tϭ͘ĨŝůĞ
returns a reference to the newly created cache c ͘͘͘ ĞĨĂƵůƚ/K
removeCache (Cache c) Ŷ tŶ͘ĨŝůĞ
createRule (IO Header h, Cache c)
tŶ
ĞĨĂƵůƚ/K
creates cache rule <src,op,file,range>→ c sD ,LJƉĞƌǀŝƐŽƌ ^ƚŽƌĂŐĞ
removeRule (IO Header h, Cache c)
configureCache (<size,eviction pol,write pol>, Cache c) To answer this question, the controller uses information
getCacheStats (Cache c) from the Metrics Engine to first determine the hit ratio
returns cache statistics Hiticache required for workload Wi to meet a certain band-
width guarantee, and then allocates the workload Wi cache
Table 1: Moirai’s API for a configurable cache.
space ai , such that U(ai ) = Hiticache , where U is the work-
load’s hit ratio function (provided by the Metrics Engine).
3. Data Plane Transformations More precisely, note that if the total bandwidth achievable
storage
from the storage back-end 2 is BWi and main memory
In this section, we explore Moirai’s ability to program bandwidth is BW memory
, a workload’s bandwidth depends on
and transform the data plane to implement various cloud its hit ratio Hiticache as follows:
provider objectives and improve workload performance. For
each goal, we illustrate how the controller effects the neces- storage
SLABW
i ≤ Hiticache × BW memory + (1 − Hiticache) × BWi
sary changes on the data plane. (1)
That means the cache hit ratio in order to achieve a band-
3.1 Prioritizing a Workload width SLABW
i needs to be at least:
It’s often desirable to be able to isolate the performance of a storage
particular (high-priority) application A from that of another SLABW
i − BWi
Hiticache ≥ storage (2)
application B sharing a cache in the same VM. The con- BW memory − BWi
troller can achieve this by configuring a dedicated cache C (a
50GB LRU write-through cache in this particular example) After the min. data bandwidth guarantees SLABW BW
1 , . . . , SLAn
inside the hypervisor, which is exclusive to workload A: are met for all n workloads, the leftover cache space can be
allocated based on priorities or using approaches highlighted
1: C = createCache (< 50GB, LRU, write-through>) in Section 3.3 to optimize for global utility.
2: createRule (< V M, *, [Link], *>, C)
The createRule call configures the cache to accept all 3.3 Maximizing Global Workload Utility
R/W IOs originating from the VM, that access any part Rather than per-workload guarantees, a provider might strive
of [Link]. The figure below shows the resulting data plane. to maximize the global workload utility, i.e., the sum of
Workload A flows through its own cache C in the hyper- the utilities across all workloads in the system. Utility of
visor, while workload B continues along its previous path, a workload could be measured by hit ratio, or be defined
bypassing that cache, effectively isolating A’s traffic from it. more generally in terms of bytes per second (Bps) satisfied
by the cache, or by extending the notion of hit ratio by
͘ĨŝůĞ
introducing weights to account for the type of IO (i.e. reads
tŽƌŬůŽĂĚ
tŽƌŬůŽĂĚ ĞĨĂƵůƚ/K
vs. writes), or even to account for the impact of a workload
sD ,LJƉĞƌǀŝƐŽƌ ^ƚŽƌĂŐĞ on the storage device (e.g. sequential vs. random access).
The choice of definition for utility will be dictated by the
optimization goals of the cloud provider.
3.2 Providing Per-Workload Bandwidth Guarantees Using the example of hit ratios as the utility function,
Next we extend the objectives beyond simple priorities, and the controller can create a separate cache for each work-
examine how Moirai allocates cache space to several arbi- load (similar to Section 3.2) and then use a classic result [37]
trary workloads W1 ,W2 , . . . ,Wn , all running on the the same to determine the cache allocations a1 , ..., an . The algorithm,
system, in order to guarantee each workload Wi a partic- shown in Algorithm 3.1, uses a water-filling approach, i.e,
ular bandwidth Bi . Similar to Section 3.1, the controller it allocates the cache to workloads in small increments. The
passes each workload’s traffic through its own dedicated basic idea at each step is to allocate the next increment of
cache Ci at the hypervisor (see figure on the following col- cache to the workload that will achieve the highest hit rate
umn), but the question now becomes what the size each of out of the allocation. When the hit rate curves of workloads
the caches needs to be. In this section, we focus on hyper- are concave functions, this algorithm will achieve an alloca-
visor level caches only, but the techniques can be expanded tion that maximizes the total hit rate, i.e., total hit rate at the
to include simultaneous allocation of hypervisor and storage 2 If
the storage back-end is remote, BWi
storage
is the minimum of the net-
level cache space, as we explain in Section 3.5. work, and the back-end storage array’s bandwidth.
3
cache. We are currently investigating meta-heuristics to deal while appearing to the VM and applications as one single
with non-concave hit ratio curves. E.g. Soundararajan [36] aggregate cache. Note that today, workloads do flow through
proposed hill-climbing search, although we find that their both caches (at the hypervisor, and at the storage server),
particular algorithm and implementation is too slow for our but this occurs in an uncontrolled fashion, leading to wasted
system to react dynamically. memory capacity by double-caching of data in both places.
In situations where the hypervisor is hosting several ap-
Algorithm 3.1 Utility-maximizing cache allocation plications and memory is limited, the controller has several
Require: n workloads sharing a cache of capacity C; U1 ,...,Un : hit choices for how to split the cache for a workload A and con-
rate curves for workloads. figure it at the hypervisor(1) and storage server(2). If the
Ensure: Assign cache allocation ai to workload i s.t. ∑ ai = C, and workload access is uniform across the file, one choice is to
max ∑ U(ai ) cache half the file in each of the respective caches:
1: ∀i, ai = 0 //Initialize allocations
2: le f tC =C //Cache left to distribute 1: createRule (< V M, *, [Link], 0, size/2>, C1)
3: ε = 0.001×C //Water-filling constant (as fraction of C) 2: createRule (< V M, *, [Link], size/2+1, size>, C2)
4: while (le f tC > 0) do The resulting data plane is shown in the figure below:
5: cacheAlloc = min(ε, le f tC)
6: j = arg max(Ui (ai + cacheAlloc) −Ui (ai )) //workload with tŽƌŬůŽĂĚ ϭ ͘ĨŝůĞ
i Ϯ
the most utility gained from extra cache
7: a j + = cacheAlloc sD ,LJƉĞƌǀŝƐŽƌ ^ƚŽƌĂŐĞ
8: le f tC− = cacheAlloc The controller can also match workload access patterns
to the way the cache is split based on hot or cold blocks or
One might argue that a standard, workload-agnostic sys- files.
tem that manages the entire cache as a single pool and ap- Another option is to treat both caches as a global LRU
plies its favourite replacement policy to it is also designed to cache. To do that, the controller programs C1 to cache
achieve the same goal of maximizing overall hit ratio. How- the IOs from [Link], and C2 to only cache IOs that were
ever, Moirai can provide generalizations of this goal (e.g. a evicted (or “demoted”) from C1. To provide per-workload
weighted sum of the hit ratios across workloads) and simul- bandwidth guarantees, Moirai extends the cache allocation
taneously provide other goals, such as isolation (e.g. pro- method presented in Section 3.2. The controller now needs
tecting one workload from the effects of workload spikes in to determine two things:
another workload), which a standard system cannot.
1. How much cache space ai to allocate for the global LRU
3.4 Consolidating Memory Over Fast Networks cache made up of both C1 and C2, such that:
As systems are increasingly making use of fast networks U(ai ) = Hiticache (3)
with speeds in excess of 40-100Gbps , and RDMA capabil-
ities [10], use of remote resources is becoming increasingly 2. The individual cache space allocations a1i and a2i , for C1
feasible and can improve overall utilization of resources. and C2 respectively. Thus, for some α:
Consider as an example a read-only dataset [Link] ac-
cessed by N VMs across N hypervisors. Placing one consol- a1i = α × ai
(4)
idated cache at the storage server can result in an Nx reduc- a2i = (1 − α) × ai
tion in total cache space used, with potentially only small in-
creases in latency. The controller can accomplish this as fol- The relationship between these variables is illustrated
lows (using the example of a 100GB MRU write-back cache using a simple, example hit rate utility function in Figure 4.
C as the consolidated cache): A workload’s bandwidth SLABW i now depends on the hit
1: C = createCache (<100GB , MRU, write-back>) ratio Hiticache of the global LRU cache as follows:
2: createRule (< V M1 − N, *, [Link], *>, C)
The resulting data plane is shown in the figure below: storage
SLABW
i ≤ Hiticache ×BWiGlobalLRU +(1−Hiticache )×BWi
sDϭ (5)
͘͘͘ d͘ĨŝůĞ Similar to Equation (2), the cache hit ratio Hiticache needs
sDŶ to be at least:
ĞĨĂƵůƚ/K
sD ,LJƉĞƌǀŝƐŽƌ ^ƚŽƌĂŐĞ storage
SLABW
i − BWi
Hiticache ≥ GlobalLRU storage (6)
3.5 Scaling Out Caches BWi − BWi
In addition to fully-remote caching, caching capacity per Here, BWiGlobalLRU refers to the total bandwidth achiev-
workload can be split across the compute and storage server, able from the global LRU cache.
4
[Link] [Link] [Link] [Link] Exchange
Read % 75% 61% 56% 1% 40%
IO Sizes 64 KB 8 KB 64 KB 64 KB 8 KB
Seq/rand Mixed Rand Rand Seq Rand
# IOs 32M 158M 36M 54M 60M
form IOFlow does keep track of each IO’s tenant class, it was
designed to provide IO queueing and rate limiting based on
IO request headers. By contrast, caching involves inspection
and manipulation of the data associated with an IO request.
Figure 4: Example of a cache hit rate function U, and We implemented an extension of the IOFlow architecture
associated parameters ai and α, used to compute cache to add support for data transformations using a version of
allocations for scaled-out LRU caches with bandwidth the Windows messaging API for filter drivers, in around 500
guarantees. new LOC. IOs are passed to a user-level cache through an
upcall, while a kernel-mode thread handling the I/O request
blocks pending a return code from the cache. The latter de-
Since C1 and C2 form the global LRU cache, the fraction cides whether the request is terminated at the filter driver
α of cache space allocated to C1 will result in U(αai )% (hit), or is sent further down the IO stack.
of the hits, while the rest of the hits, [U(ai ) − U(αai )]%,
will be served from C2. Since C1 is a hypervisor cache, its
achievable bandwidth is BW memory , while C2’s achievable 5. Experimental Evaluation
bandwidth is constrained by the bandwidth of the network This section provides an experimental evaluation of some of
BWinetwork . Thus: Moirai’s use cases presented in Section 3. Our experimen-
tal testbed has 12 servers, each with 16 Intel Xeon 2.4 GHz,
BWiGlobalLRU = U(αai ) × BW memory 384 GB of RAM and three Seagate Constellation 2 disks
(7) or four Intel 520 SSDs in RAID-0. The servers run Win-
+ [U(ai) − U(αai )] × BWinetwork dows Server 2012 R2 operating system and can act as either
Hyper-V hypervisors or as storage servers. Each server has
Simultaneously using Equations (3), (4), (6), and (7), the
a 40 Gbps Mellanox ConnectX-3 NIC supporting RDMA
controller solves for ai , and α, effectively determining the
and connected to a Mellanox MSX1036B-1SFR switch. We
cache allocations a1i and a2i , for C1 and C2 respectively. Fur-
use a combination of real enterprise application traces and
ther constraints can also be added to the problem statement
benchmarks, as specified in more detail below. As Moirai is
(e.g., imposing a maximum size on either C1, or C2) to limit
transparent to applications and VM’s, they can run on our
the solution space.
testbed without modifications
We use a mixture of real enterprise application traces and
4. Implementation benchmarks in the evaluation. For the former, we use pub-
We have implemented and deployed a Moirai prototype, lic traces from an enterprise Exchange email server [35] and
comprising all components described in Section 2, on a Hotmail [39]. Key characteristics of these traces are shown
Windows-based system and made the code publicly avail- in Table 2. The traces are diverse across a number of met-
able [28]. The controller is implemented in around 6000 rics such as the Read-to-Write ratio, IO sizes, sequentiality
LOC of C# and communicates with the caches through RPCs of access and number of IOs which allows for a comprehen-
over TCP. The Metrics Engine is implemented as a user-level sive evaluation across realistic workload mixes. However, a
stage in the hypervisor in around 500 LOC and uses a variant limitation of these workloads is that they were originally col-
of SHARDS [44] to determine hit ratio curves. Cache mod- lected underneath file caches. As such, they under-represent
ules implement the APIs in Table 1 at user-level in around the amount of application reads.
2000 LOC in C#. To account for this limitation, we also use TPC-E [41]
One implementation challenge is how to classify and di- and TPC-H [42] to cover a broad class of workloads, rang-
rect a tenant’s traffic to the configurable caches. We decided ing from transaction processing OLTP operations with small
to build an extension of the IOFLow framework [38] to im- IO sizes (TPC-E) to large streaming IO from data mining
plement this functionality. Note that while in its original queries (TPC-H). They run over unmodified SQL Server
5
1200 Hypervisor Storage 40Gbps Storage 1Gbps
Default caching Moirai
Transactions/minute 1000 100
Latency (seconds)
800
600
10
400
200
0 1
TPC-E alone TPC-E with TPC-E alone TPC-E with Q2 Q6 Q7 Q15 Q19
TPC-H TPC-H
Figure 5: Prioritizing one workload (TPC-E) vs another Figure 6: Latency for 5 TPC-H queries. The controller
(TPC-H). With Moirai, the performance of TPC-E is not can decide to use file caches in the storage server for fast
impacted by TPC-H. In contrast, today, running both RDMA-based networks. Y-axis is log scale.
workloads together would result in a 5x performance hit
for TPC-E. We compare two approaches of dividing up the cache
space. In the first we divide space equally among the four
2012 R2 databases. When error bars are shown they repre- tenants. In the second we use the method described in Sec-
sent the average, minimum and maximum from 5 runs. tion 3.3 to partition the cache and reconfigure the data plane.
The results are shown in Figure 2(b).
5.1 Enforcing Priorities Interestingly, we observe not only that overall throughput
We examine Moirai’s ability to prioritize a workload using increases by more than 2.5x, but also that this improvement
the example of a VM with one SQL Server instance running comes at no cost to any of the individual four tenants. The
both TPC-E and TPC-H. The corresponding database files, reason is that all tenants benefit from the decreased load at
“[Link]” and “[Link]” each have a footprint of the storage back-end.
50GB and are stored on Virtual Hard Drives (VHDs) on two We have experimented with other workload combinations
separate disk-based storage servers. as well. In the worst case across all experiments the overall
We run two experiments, one with default caching and throughput still increased by 35%, but this came at the cost
one where we use Moirai to prioritize the TPC-E workload, of a small penalty to one tenant, whose throughput dropped
as explained in Section 3.1 and measure the throughput by 10%. A cloud provider could feed into the controller a
(transactions/min) for the TPC-E workload. The results are minimum amount of cache space or minimum hit rate it
shown in Figure 5. wants to guarantee each workload, and then ask it to divide
We observe that in the system without Moirai, TPC-E’s the remaining cache space to maximize global utility.
performance drops by more than 5X when TPC-H runs. On
the other hand, we find that with Moirai, TPC-E’s throughput 5.3 Consolidating Memory Over Fast Networks
running alongside TPC-H is within 2.3 % of its throughput In this section we use Moirai on a TPC-H workload running
running by itself. on ten different hypervisors to illustrate the trade-offs for
Note that our current implementation of Moirai results in memory consolidation over fast networks. We compare the
a data plane overhead of 20% (this difference is due to using case where Moirai is used to insert a 50GB cache inside
our user-level cache vs. SQL Server’s native cache, which is each of the 10 hypervisors, to the case where Moirai inserts
heavily optimized). The overhead stems in part from extra one shared 50GB cache at the storage server, which is either
memory copies between the kernel and the user-level cache. accessed at 1Gbps over TCP or at 40Gbps over RDMA. In
However, we believe that this overhead is acceptable com- all cases, all the data resides in memory (100% hit rate). The
pared to the 5x drop in performance with today’s caching in- results are shown in Figure 6.
frastructure. Further, note that the controller can detect when We observe that the average latency overhead when using
no other workloads are running and remove the user-level a consolidated cache over the fast network is around 26%,
cache and thus avoid the extra overhead. compared to using local hypervisor caches. For the slow
network the overheads are 153%. Note that in exchange for
5.2 Maximizing Global Hit Rate paying these overheads one gains a 10X reduction in the
We consider the example of maximizing global hit rate using total amount of cache space allocated for this workload. Also
four tenants with 30 VMs each, spread over 10 hypervisors note that with Moirai a provider has the option to seamlessly
accessing VHDs on an SSD-based storage server. Each ten- switch from one cache configuration to another, depending
ant’s VM uses IOMeter, parameterized with the key charac- on the state of the system. For example, a provider might
teristics of the Hotmail workloads (Tenants 1-4 are running switch to a consolidated cache at the cost of some latency
the Index, Data, Msg and Log workloads respectively). penalties when cache space is scarce.
6
2000 90 8
Cache size Throughput
80 7
Transactions/minute
10 no cache 120 flows allocated 60 flows go idle. 60 active 120 flows allocated 1
0 72GB cache flows allocated 72GB cache 72GB cache
0
24
48
72
96
120
144
168
192
216
240
264
288
312
336
360
384
408
432
456
480
504
528
552
576
600
(1x5GB) caching (2x5GB)
(2x5GB) Time (s)
Figure 7: Splitting IOs for TPC-E across two different Figure 8: Moirai adapting to workloads dynamically
caches. Today, “double caching” occurs since all IOs flow over time. Note there are two y-axis.
through all caches. Moirai can prevent this, and match
the performance of an aggregate cache. 5.6 Control Plane Overheads
We now consider Moirais overheads on the control plane.
5.4 Scaling Out Caches Moirai implements a version of SHARDS [44] in the Met-
In this section, we evaluate Moirai’s ability to scale out rics Engine to construct hit ratio curves at runtime. While our
the storage cache as described in Section 3.5. We con- current implementation is not as highly-optimized, the orig-
sider a TPC-E workload on a machine low on memory. The inal SHARDS paper showed that hit ratio curves with very
provider wishes to scale out TPC-E’s 5GB hypervisor cache high fidelity can be constructed online using less than 10MB
to include another 5GB at the storage server. of memory per workload [44], and marginal (less than 5
Figure 7 shows the results when allocating the split cache We also evaluated the time it takes to compute the opti-
with and without Moirai. Without Moirai there is little ben- mal cache size allocation as a function of the number of VMs
efit to adding a 5GB cache at the storage server (second bar in the system (Algorithm 3.1, described in Section 3.3). We
from the left), compared to having only the 5GB cache at varied the number of VMs from 100 to 5000, and measured
the hypervisor (left-most bar), due to double caching. On the time it took to compute the allocation. For 5000 VMs, it
the other hand when setting up the two caches with Moirai took less than 15s to make that decision, with a water-filling
(third bar from the left), performance is similar to that of an constant ε of 1MB. For ε of 2MB, the time is less than 5s.
aggregate 10GB cache at the hypervisor (right-most bar). This highlights the tradeoff between how fine-grained the
cache allocation is, and the completion time for the alloca-
5.5 Dynamic Workloads tion algorithm. However, since cache allocations can feasi-
bly be done at granularities (ε) coarser than 2MB, and they
The controller in Moirai continuously monitors the metrics
do not need to be re-computed at very short time intervals,
and dynamically reacts to changes after some reaction time
we believe this method of cache allocation is very reason-
s, a configurable parameter. For example, Moirai will detect
able. We are currently exploring optimizations to reduce the
when a cache goes unutilized and reuse the space accord-
algorithms runtime further.
ingly. We have worked with values for s on the order of
15-30 seconds - we believe that this range presents a good
trade-off between responsiveness and unwarranted reconfig- 6. Related work
urations due to momentary changes in workload demand. Application caches. There has been much work recently
To illustrate Moirai’s dynamic capabilities, we evaluate a on caches in data centers. Much of it focused on spe-
setup consisting of 10 hypervisors each with 12 VMs, where cialized application caches, such as Facebook’s photo-
each VM has a 2GB file stored on an SSD back-end that serving stack [20], Facebook’s social graph store [3], mem-
it accesses through IOMeter. The provider uses Moirai to cached [12], or explicit cloud caching services [6, 7]. In
allocate a total of 72GB cache at the hypervisor level, which contrast, our work is on system storage caches for hosting
is evenly split between VMs (i.e., each receiving 72/120 cloud providers that run arbitrary workloads.
GB). System caches. Work on system caches has focused on
Figure 8 shows what happens if half of the VMs on each efficient use of memory for virtual machines through bal-
hypervisor go idle for some time, before they become active looning and sharing techniques [18, 29, 43], which are im-
again at a later point. Moirai detects when the VMs go plemented in state-of-the-art hypervisors like VMware’s and
idle and re-computes the cache allocation for each VM to Hyper-V. Our work focuses on other caches in the system,
distribute spare capacity, hence improving the performance beneath the VM abstraction.
of the active VMs. Once all VMs become active again, Cache replacement policies Some prior work has fo-
Moirai re-computes the initial allocations, and performance cused on isolating the cache effects of streams with differ-
goes back to previous levels. ent access patterns (sequential versus looping) within the
7
same workload from each other [8, 13, 23, 27]. However, Aug. 2011.
these policies are not workload or tenant aware and cannot [2] J.-P. Billaud and A. Gulati. hClock: Hierarchical QoS for
prevent a more aggressive workload from occupying more packet scheduling in a hypervisor. In EuroSys, Apr. 2013.
than its fair share of cache. Moreover, each of these policies [3] N. Bronson, Z. Amsden, G. Cabrera, P. Chakka, P. Di-
might actually work better when applied in the context of mov, H. Ding, J. Ferris, A. Giardullo, S. Kulkarni, H. Li,
Moirai, where a cache policy works on per workload seg- M. Marchukov, D. Petrov, L. Puzar, Y. J. Song, and
regated cache, as patterns of different workloads don’t get V. Venkataramani. Tao: Facebook’s distributed data store for
interspersed and hence might be easier to detect. Others pro- the social graph. In Presented as part of the 2013 USENIX An-
pose methods to detect changes in workload patterns and dy- nual Technical Conference (USENIX ATC 13), pages 49–60,
namically adjust the caching policy used by the system [14]. San Jose, CA, 2013. USENIX.
Moirai provides a perfect vehicle for implementing such an [4] M. Casado, M. J. Freedman, J. Pettit, J. Luo, N. McKeown,
approach and it would be interesting to extend it to support and S. Shenker. Ethane: taking control of the enterprise. In
such functionality. Yet another line of work [19, 25] pro- Proceedings of ACM SIGCOMM, Kyoto, Japan, 2007.
poses that applications explicitly manage their cache space [5] Z. Chen, Y. Zhang, Y. Zhou, H. Scott, and B. Schiefer. Em-
and its contents, while our goal was to provide a solution that pirical evaluation of multi-level buffer cache collaboration for
is transparent to the application. storage systems. In Proceedings of the 2005 ACM SIGMET-
Inefficiencies in cache hierarchies Several other papers RICS International Conference on Measurement and Model-
have addressed the problem of inefficiencies in cache hier- ing of Computer Systems, SIGMETRICS ’05, pages 145–156,
New York, NY, USA, 2005. ACM.
archies, e.g., some [5, 26] pass hints from the client to bet-
ter inform caching decisions at the storage server and oth- [6] G. Chockler, G. Laden, and Y. Vigfusson. Data caching as a
ers [46] extend the SCSI command set by a demote com- cloud service. In Proceedings of the 4th International Work-
shop on Large Scale Distributed Systems and Middleware,
mand to avoid double caching. Our goal was a solution that
LADIS ’10, pages 18–21, New York, NY, USA, 2010. ACM.
does not require application or VM support, or changes to
existing protocols. [7] G. Chockler, G. Laden, and Y. Vigfusson. Design and imple-
mentation of caching services in the cloud. IBM Journal of
Software defined storage. Similar to recent work on
Research and Development, 55(6):9:1–9:11, Nov 2011.
software-defined networking (SDNs) [4, 11, 21, 24, 31, 40,
47] and storage (SDS) [38], our architecture is controller- [8] J. Choi, S. H. Noh, S. L. Min, and Y. Cho. An imple-
mentation study of a detection-based adaptive block replace-
based with a separation between the data and the control
ment scheme. In Proceedings of the Annual Conference on
plane. Moirai’s implementation uses IOFlow’s [38] mecha-
USENIX Annual Technical Conference, ATEC ’99, pages 18–
nisms for traffic classification, however Moirai’s implemen- 18, Berkeley, CA, USA, 1999. USENIX Association.
tation required extensions to IOFlow, e.g., to support arbi-
[9] H.-T. Chou and D. J. DeWitt. An evaluation of buffer man-
trary inspection and manipulation of IO request data, as well
agement strategies for relational database systems. In Pro-
as the implementation of the three core components Moirai ceedings of the 11th International Conference on Very Large
comprises (as described in Section 2 and Section 4). Data Bases - Volume 11, VLDB ’85, pages 127–141, Stock-
holm, Sweden, 1985. VLDB Endowment.
7. Summary [10] A. Dragojevic, D. Narayanan, O. Hodson, and M. Castro.
Caches are a critical resource in data centers. They im- Farm: Fast remote memory. In Proceedings of the 11th
prove latency, throughput and reduce the load on networks USENIX Conference on Networked Systems Design and Im-
plementation, NSDI’14, Seattle, WA, 2014. USENIX Associ-
and storage. But today, caches are implicit, not designed
ation.
for controlled sharing, leading to severe inefficiencies un-
der multi-tenancy. This paper presents Moirai, a software- [11] A. D. Ferguson, A. Guha, C. Liang, R. Fonseca, and S. Krish-
defined caching architecture that enables control of caches namurthi. Participatory networking: An API for application
control of SDNs. In Proceedings of ACM SIGCOMM, Hong
in a multi-tenant data center. Moirai is transparent to hosted
Kong, 2013.
tenants. Their throughput and latency benefit without requir-
ing any tenant input or hints. We show using several differ- [12] B. Fitzpatrick. Distributed caching with memcached. Linux
J., 2004(124):5–, Aug. 2004.
ent use cases how Moirai can help ease the management of
the distributed caching infrastructure and enable the provider [13] C. Gniady, A. R. Butt, and Y. C. Hu. Program-counter-
to achieve a series of different objectives. We hope that our based pattern classification in buffer caching. In Proceedings
of the 6th Conference on Symposium on Opearting Systems
public release of the code [28] implementing Moirai will
Design & Implementation - Volume 6, OSDI’04, pages 27–27,
help foster future work in this area. Berkeley, CA, USA, 2004. USENIX Association.
[14] R. B. Gramacy, M. K. Warmuth, S. A. Brandt, I. Ari, and I. A.
References . Adaptive caching by refetching. In In Advances in Neural
[1] H. Ballani, P. Costa, T. Karagiannis, and A. Rowstron. To- Information Processing Systems 15, pages 1465–1472. MIT
wards predictable datacenter networks. In ACM SIGCOMM, Press, 2002.
8
[15] A. Gulati, I. Ahmad, and C. A. Waldspurger. PARDA: propor- [28] Microsoft. Microsoft research storage toolkit. http://
tional allocation of resources for distributed storage access. In [Link]/en-us/downloads/
FAST, Feb. 2009. 230b8ad4-3340-4a87-8ef0-cf92b376db86/.
[16] A. Gulati, A. Merchant, and P. J. Varman. mClock: Handling [29] G. Miłós, D. G. Murray, S. Hand, and M. A. Fetterman. Satori:
throughput variability for hypervisor IO scheduling. In OSDI, Enlightened page sharing. In Proceedings of the 2009 Confer-
Oct. 2010. ence on USENIX Annual Technical Conference, USENIX’09,
pages 1–1, Berkeley, CA, USA, 2009. USENIX Association.
[17] C. Guo, G. Lu, H. J. Wang, S. Yang, C. Kong, P. Sun, W. Wu,
and Y. Zhang. SecondNet: A data center network virtu- [30] L. Popa, G. Kumar, M. Chowdhury, A. Krishnamurthy, S. Rat-
alization architecture with bandwidth guarantees. In ACM nasamy, and I. Stoica. Faircloud: Sharing the network in cloud
CoNEXT, Nov. 2010. computing. In ACM SIGCOMM, Aug. 2012.
[18] D. Gupta, S. Lee, M. Vrable, S. Savage, A. C. Snoeren, [31] Z. A. Qazi, C.-C. Tu, L. Chiang, R. Miao, S. Vyas, and M. Yu.
G. Varghese, G. M. Voelker, and A. Vahdat. Difference en- SIMPLE-fying middlebox policy enforcement using SDN. In
gine: Harnessing memory redundancy in virtual machines. In Proceedings of the ACM SIGCOMM, Hong Kong, 2013.
Proceedings of the 8th USENIX Conference on Operating Sys- [32] T. Saemundsson, H. Bjornsson, G. Chockler, and Y. Vig-
tems Design and Implementation, OSDI’08, pages 309–322, fusson. Dynamic performance profiling of cloud caches. In
Berkeley, CA, USA, 2008. USENIX Association. Proceedings of the ACM Symposium on Cloud Computing,
[19] K. Harty and D. R. Cheriton. Application-controlled physical SOCC ’14, pages 28:1–28:14, New York, NY, USA, 2014.
memory using external page-cache management. In Proceed- ACM.
ings of ACM ASPLOS, Boston, Massachusetts, USA, 1992. [33] A. Shieh, S. Kandula, A. Greenberg, and C. Kim. Sharing the
[20] Q. Huang, K. Birman, R. van Renesse, W. Lloyd, S. Kumar, datacenter network. In NSDI, Mar. 2011.
and H. C. Li. An analysis of Facebook photo caching. In [34] D. Shue, M. J. Freedman, and A. Shaikh. Performance isola-
Proceedings of ACM SOSP, Farminton, Pennsylvania, 2013. tion and fairness for multi-tenant cloud storage. In OSDI, Oct.
[21] S. Jain, A. Kumar, S. Mandal, J. Ong, L. Poutievski, A. Singh, 2012.
S. Venkata, J. Wanderer, J. Zhou, M. Zhu, J. Zolla, U. Hölzle, [35] SNIA. Exchange server traces. [Link]
S. Stuart, and A. Vahdat. B4: Experience with a globally- traces/130.
deployed software defined wan. In Proceedings of ACM [36] G. Soundararajan, J. Chen, M. A. Sharaf, and C. Amza. Dy-
SIGCOMM, Hong Kong, China, 2013. namic partitioning of the cache hierarchy in shared data cen-
[22] V. Jeyakumar, M. Alizadeh, D. Mazires, B. Prabhakar, and ters. Proc. VLDB Endow., 1(1):635–646, Aug. 2008.
C. Kim. EyeQ: Practical network performance isolation at the [37] H. S. Stone, J. Turek, and J. L. Wolf. Optimal partitioning
edge. In NSDI, Apr. 2013. of cache memory. IEEE Trans. Comput., 41(9):1054–1068,
[23] J. M. Kim, J. Choi, J. Kim, S. H. Noh, S. L. Min, Y. Cho, and Sept. 1992.
C. S. Kim. A low-overhead high-performance unified buffer [38] E. Thereska, H. Ballani, G. O’Shea, T. Karagiannis, A. Row-
management scheme that exploits sequential and looping ref- strow, T. Talpey, R. Black, and T. Zhu. IOFlow: A software-
erences. In Proceedings of the 4th Conference on Symposium defined storage architecture. In Proceedings of ACM SOSP,
on Operating System Design & Implementation - Volume 4, 2013.
OSDI’00, pages 9–9, Berkeley, CA, USA, 2000. USENIX As-
[39] E. Thereska, A. Donnelly, and D. Narayanan. Sierra: practical
sociation.
power-proportionality for data center storage. In Proceedings
[24] T. Koponen, M. Casado, N. Gude, J. Stribling, L. Poutievski, of Eurosys, pages 169–182, Salzburg, Austria, 2011.
M. Zhu, R. Ramanathan, Y. Iwata, H. Inoue, T. Hama, and
[40] N. Tolia, M. Kaminsky, D. G. Andersen, and S. Patil. An
S. Shenker. Onix: a distributed control platform for large-
architecture for internet data transfer. In Proceedings of
scale production networks. In Proceedings of USENIX OSDI,
USENIX NSDI, NSDI’06, San Jose, CA, 2006.
Vancouver, BC, Canada, 2010.
[41] TPC Council. TPC-E. [Link]
[25] C.-H. Lee, M. C. Chen, and R.-C. Chang. Hipec: High per-
formance external virtual memory caching. In Proceedings of [42] TPC Council. TPC-H. [Link]
USENIX OSDI, Monterey, California, 1994. [43] C. A. Waldspurger. Memory resource management in
[26] X. Li, A. Aboulnaga, K. Salem, A. Sachedina, and S. Gao. VMware ESX server. SIGOPS Oper. Syst. Rev., 36(SI):181–
Second-tier cache management using write hints. In Proceed- 194, Dec. 2002.
ings of the 4th Conference on USENIX Conference on File [44] C. A. Waldspurger, N. Park, A. Garthwaite, and I. Ahmad. Ef-
and Storage Technologies - Volume 4, FAST’05, pages 9–9, ficient MRC construction with SHARDS. In 13th USENIX
Berkeley, CA, USA, 2005. USENIX Association. Conference on File and Storage Technologies (FAST 15),
[27] N. Megiddo and D. S. Modha. Arc: A self-tuning, low over- pages 95–110, Santa Clara, CA, Feb. 2015. USENIX Asso-
head replacement cache. In Proceedings of the 2Nd USENIX ciation.
Conference on File and Storage Technologies, FAST ’03, [45] J. Wires, S. Ingram, Z. Drudi, N. J. A. Harvey, and
pages 115–130, Berkeley, CA, USA, 2003. USENIX Asso- A. Warfield. Characterizing storage workloads with counter
ciation. stacks. In 11th USENIX Symposium on Operating Systems
9
Design and Implementation (OSDI 14), Broomfield, CO, Oct.
2014. USENIX Association.
[46] T. M. Wong and J. Wilkes. My cache or yours? making storage
more exclusive. In Proceedings of USENIX ATC, Monterey,
California, 2002.
[47] H. Yan, D. A. Maltz, T. S. E. Ng, H. Gogineni, H. Zhang, and
Z. Cai. Tesseract: a 4D network control plane. In Proceedings
of USENIX NSDI, Cambridge, MA, 2007.
10