<!DOCTYPE article PUBLIC "-//NLM//DTD JATS (Z39.96) Journal Archiving and Interchange DTD v1.0 20120330//EN" "JATS-archivearticle1.dtd">
<article xmlns:xlink="http://www.w3.org/1999/xlink">
  <front>
    <journal-meta />
    <article-meta>
      <title-group>
        <article-title>Communication-efficient Outlier Detection for Scale-out Systems</article-title>
      </title-group>
      <contrib-group>
        <contrib contrib-type="author">
          <string-name>Moshe Gabel Technion Haifa</string-name>
          <email>assaf@cs.technion.ac.il</email>
          <email>dkeren@cs.haifa.ac.il</email>
          <email>mgabel@cs.technion.ac.il</email>
          <xref ref-type="aff" rid="aff0">0</xref>
          <xref ref-type="aff" rid="aff1">1</xref>
        </contrib>
        <contrib contrib-type="author">
          <string-name>Israel</string-name>
          <xref ref-type="aff" rid="aff0">0</xref>
          <xref ref-type="aff" rid="aff1">1</xref>
        </contrib>
        <aff id="aff0">
          <label>0</label>
          <institution>Assaf Schuster Technion Haifa</institution>
          ,
          <country country="IL">Israel</country>
        </aff>
        <aff id="aff1">
          <label>1</label>
          <institution>Daniel Keren Haifa University Haifa</institution>
          ,
          <country country="IL">Israel</country>
        </aff>
      </contrib-group>
      <abstract>
        <p>Modern scale-out services are built on top of large datacenters composed of thousands of individual machines. These must be continuously monitored because unexpected failures can overload fail-over mechanism and cause large-scale outages. Such monitoring can be accomplished by periodically measuring hundreds of performance metrics and looking for outliers, often caused by miscon gurations, hardware failures or even software bugs. Previous work has shown that many failures are indeed preceded by such performance outliers, known as performance problems or latent faults. In this work we adapt an existing unsupervised statistical framework for latent fault detection to provide an online, communication- and computation-reduced version. The existing framework is e ective in predicting machine failures days before they happen, but requires each monitored machine to send all its periodic metric measurements, which is prohibitive in some settings and requires that the datacenter provide parallel storage and processing. Our adapted framework is able to reduce the amount of data sent and the processing cost at the central coordinator by processing the data in situ, making it usable in wider settings. We utilize techniques from the domain of stream processing, speci cally sketching and safe zones, to trade-o accuracy for communication and computation, without compromising its advantages. Like the original framework, our adapted framework is unsupervised, does not require domain knowledge, and provides statistical guarantees on the rate of false positives. Initial experiments show that scores yielded by the adapted framework match the original scores very well, while reducing communications by over 90%.</p>
      </abstract>
    </article-meta>
  </front>
  <body>
    <sec id="sec-1">
      <title>1. INTRODUCTION</title>
      <p>In recent years the demand for computing power and
storage has increased. Modern Web services and clouds rely on
large datacenters, often comprised of thousands of machines.
For such large services, it is unreasonable to assume that all
machines are working properly and are well con gured.</p>
      <p>
        Monitoring is essential in datacenters, since unnoticed
faults might accumulate to the point where redundancy and
fail-over mechanisms break. Yet the large number of
machines in datacenters makes manual monitoring
impractical. Instead machines are usually monitored by collecting
and analyzing performance counters [
        <xref ref-type="bibr" rid="ref11 ref3 ref5">3, 5, 11</xref>
        ]. Hundreds
of counters per machine are reported by the various service
layers, from service-speci c metrics (such as database query
statistics) to general metrics (such as CPU utilization).
      </p>
      <p>
        In this work we adapt an existing fault detection algorithm
[
        <xref ref-type="bibr" rid="ref9">9</xref>
        ] using sketching [
        <xref ref-type="bibr" rid="ref18 ref7 ref8">8, 18, 7</xref>
        ] and safe zones [
        <xref ref-type="bibr" rid="ref17 ref21">17, 21</xref>
        ] to reduce
communication and processing requirements by an order of
magnitude, while preserving its advantages.
      </p>
      <p>
        Many existing failure detectors are in exible [
        <xref ref-type="bibr" rid="ref9">9</xref>
        ], and most
require centralizing the data in some form. Rule-based
failure detectors de ne a set of watchdogs [
        <xref ref-type="bibr" rid="ref11">11</xref>
        ] that monitor
speci c counters and trigger an alert whenever a prede ned
threshold is crossed. However, maintaining these static rules
requires ongoing manual adjustments.
      </p>
      <p>
        More advanced methods model service behavior from
historical logs. Supervised machine learning approaches [
        <xref ref-type="bibr" rid="ref20 ref3 ref4 ref6">3, 6,
20, 4</xref>
        ] train detectors on historic annotated data. Others [
        <xref ref-type="bibr" rid="ref5">5</xref>
        ]
analyze logs from periods from when the service is
guaranteed to be healthy to extract model parameters. Such
approaches are sensitive to deviations in workloads and changes
in the monitored service itself [
        <xref ref-type="bibr" rid="ref10 ref23">23, 10</xref>
        ]. After such changes
the historical logs and the learned model are no longer
relevant. Approaches that require labeled data can be
expensive, since labels can be di cult to obtain, and re-labeling
may be needed after service changes.
      </p>
      <p>
        More exible, unsupervised approaches have been
proposed for high performance computing (HPC). Typical
approaches [
        <xref ref-type="bibr" rid="ref19 ref22">19, 22</xref>
        ] analyze textual console logs to detect
system or machine failures by examining frequency of log
messages. Console logs are impractical in high-volume services
for bandwidth and performance reasons: transactions are
very short, time-sensitive, and rapid.
      </p>
      <p>
        Finally, some approaches [
        <xref ref-type="bibr" rid="ref14 ref16">14, 16</xref>
        ] are unsupervised and
exible, but are not domain independent. They make use
of domain insights and knowledge of the monitored service,
for example in the domain of distributed le systems, and
are therefore limited to speci c systems.
      </p>
      <p>
        Recent approaches to the monitoring problem [
        <xref ref-type="bibr" rid="ref15 ref16 ref9">9, 16, 15</xref>
        ]
focus on early detection and handling of performance
problems, or latent faults. These are outliers { machine behaviors
that are indicative of a fault, or could eventually result in a
fault, yet y under the radar of monitoring systems because
they are not acute enough, or were not anticipated by the
monitoring system designers. Early detection of latent faults
can help prevent future failures and increase the reliability
of services.
      </p>
      <p>
        In previous work [
        <xref ref-type="bibr" rid="ref9">9</xref>
        ] we provided evidence that latent
faults are common, and we presented a novel, unsupervised
outlier detection framework for latent fault detection. In
experiments on a real-world production system comprised
of 4500 machines, we showed that over 20% of machine
failures were preceded by latent faults. Furthermore, we were
able to detect latent faults up to 14 days in advance of
actual machine or software failures with up to 70% precision
and 2% false positive rate { comparable to state of the art
supervised techniques in controlled settings [
        <xref ref-type="bibr" rid="ref4">4</xref>
        ]. We
demonstrated that our system is adaptable, requiring no domain
knowledge, no labeled examples, and no parameter tuning
in the face of workload changes and software updates.
Finally, our system has proven and demonstrated guarantees
on the false positive rates, it is non-intrusive, and it scales
to very large services.
      </p>
      <p>
        One drawback of previous work is the large
communication and processing costs, prohibitive in some settings.
Modern data centers are large, and consequently the
resultant counter logs are also large. It may be very di cult to
centralize and process such a large amount of data. In the
experiments described in [
        <xref ref-type="bibr" rid="ref9">9</xref>
        ], the log les were over 10TB
per day { too large to centralize and process in one location.
Instead we relied on a data-parallel infrastructure [
        <xref ref-type="bibr" rid="ref12">12</xref>
        ] built
into the data center. Parallel processing may not always
be feasible in all situations, however. Furthermore, some
large systems are not con ned to a single datacenter but are
geographically distributed.
      </p>
      <p>In this work we extend our latent fault detection using
techniques from the eld of stream processing to reduce the
size of the data by an order of magnitude, reducing
communication and processing requirements, and allowing
continuous online processing of distributed streams. The
resulting technique is essentially a distributed outlier detector for
multiple multivariate data streams, designed for monitoring
large-scale online services.</p>
    </sec>
    <sec id="sec-2">
      <title>2. SUMMARY OF PREVIOUS WORK</title>
      <p>
        In [
        <xref ref-type="bibr" rid="ref9">9</xref>
        ] we presented a statistical latent fault detection
framework with 3 derived tests. What follows is a short summary
of that work, with the sign test as example.
2.1
      </p>
    </sec>
    <sec id="sec-3">
      <title>Framework</title>
      <p>We begin with a reasonable assumption: in a large
cluster of machines doing the same job, most machines perform
well most of the time. Further, we expect similar machines
with similar hardware and software1 to exhibit roughly
similar behavior when measuring performance counters. We
therefore compare these machines to nd those whose
performance di ers notably.</p>
      <p>There are M machines, each reporting C performance
counters at every time t in a window of length T time points.
We denote by x(m; t) the vector of counter values for
machine m at time t. The hypothesis is that the inspected
machine is working properly and hence the statistical
process that generated this vector for machine m is the same
statistical process that generated the vector for any other
machine m0. However, if we see that the vector x(m; t) for
machine m is notably di erent from the vectors of other
machines, we reject the hypothesis and ag the machine m as
suspicious, meaning we suspect it manifests a latent fault.</p>
      <p>We now make explicit our assumptions on the behavior
of the monitored machines: a) the majority of machines are
working properly at any given point in time; b) the machines
are homogeneous, meaning they perform a similar task and
use similar hardware and software2; c) on average, the
workload is balanced across all machines; d) the counters are
ordinal and are reported at the same rate; and e) the counter
values are memoryless in the sense that they depend only on
the current time period (and are independent of the identity
of the machine).</p>
      <p>Formally, we assume that x(m; t) is a realization of a
random variable X(t) whenever machine m is working properly.
Since all machines perform the same task, and since the load
balancer attempts to split the load evenly between the
machines, the homogeneous assumption implies that we should
expect x(m; t) to show similar behavior. We do expect to
see changes over time, due to changes in the workload, for
example. However, we expect these changes to be similarly
re ected in all machines.</p>
      <p>At any time t, the input x(t) to a test S consists of the
vectors x(m; t) for all machines m. The test S(m; x(t))
analyzes the data and assigns a score (either a scalar or a vector)
to machine m at time t. Given a test S, and a signi cance
level &gt; 0, we can present the framework as follows:
1. Preprocess: select counters and scale to unit variance;</p>
      <sec id="sec-3-1">
        <title>2. Compute for every machine m the vector:</title>
        <p>vm = T1 Pt S(m; x(t)) (integration phase);
3. Compute the p-values (de ned below) p(m) from vm;</p>
      </sec>
      <sec id="sec-3-2">
        <title>4. Report every machine with p(m) &lt;</title>
        <p>as suspicious.</p>
        <p>Essentially, the scores for machine m are aggregated over
time, so that eventually the norm of the aggregated scores
converges, and is used to compute a p-value for m. The
longer the allowed time period for aggregating the scores
is, the more sensitive the test will be. At the same time,
aggregating over long periods of time creates latencies in the
detection process. In our previous work we aggregated data
over 24 hour intervals, as a compromise between sensitivity
and latency.</p>
        <p>The p-value for a machine m is a bound on the probability
that a random healthy machine would exhibit such aberrant
counter values. If the p-value falls below a prede ned
signi cance level , the null hypothesis is rejected, and the
machine is agged as suspicious.</p>
        <p>
          In [
          <xref ref-type="bibr" rid="ref9">9</xref>
          ] we derived and evaluated 3 di erent tests within
the framework (di erent S functions). The sign test
accumulates the average normalized direction from machine m
to the rest of the machines. The Tukey test measures the
average depth of x(m; t) compared to the vectors of other
machines at the same time. The LOF test similarly
compares the local density of points around x(m; t) to the local
density of its neighbors. What follows is a summary of the
sign test.
1These are reasonable assumptions in practice for many
services and datacenters [
          <xref ref-type="bibr" rid="ref14 ref19">14, 19</xref>
          ].
2If this is not the case, we can often split the collection of
machines to a few large homogeneous clusters.)
2.2
        </p>
      </sec>
    </sec>
    <sec id="sec-4">
      <title>The Sign Test</title>
      <p>The sign test extends the classic statistical sign test to
allow the simultaneous comparison of multiple machines. The
\sign" of a machine m at time t is the average direction of
its vector x(m; t) to all other machines' vectors, and its score
vm is the sum of all these directions, divided by T .</p>
      <p>The intuition is that healthy machines are similar on
average, and any di erences are random. Average directions
are therefore random and tend to cancel each other out
when added together, meaning vm will be a relatively short
vector for healthy machines. Conversely, if m has a latent
fault, then some of its metrics are consistently di erent from
healthy machines, and so the average directions are similar
in some dimensions. When summing up these average
directions, these similarities reinforce each other and therefore
vm tends to be a longer vector.</p>
      <p>Formally, let M denote the set of all machines in a test,
and M = jMj the number of machines. T are the time
points where counters are sampled during preprocessing (for
instance, every 5 minutes for 24 hours in our experiments),
t denote a speci c time point, and T = jT j. Let m and m0
be two machines and let x(m; t) and x(m0; t) be the vectors
of their reported and preprocessed counters at time t. We
use the test</p>
      <p>S (m; x(t)) =</p>
      <p>M
1</p>
      <p>X</p>
      <p>x(m; t)
1 m06=m kx(m; t)
x (m0; t)
x (m0; t)k
(1)
as a multivariate version of the sign function. If all the
machines are working properly, we expect this value to be
small. Therefore, the sum of several samples over time is
also expected not to grow far from zero.</p>
      <sec id="sec-4-1">
        <title>Algorithm 1: The sign test. foreach machine m do</title>
        <p>S (m; x(t))
vm</p>
        <p>1 P x(m;t) x(m0;t)</p>
        <p>M 1 m06=m kx(m;t) x(m0;t)k ;
1 Pt S (m; x(t));</p>
        <p>T
end
v^ M1 Pm kvmk;
foreach machine m do</p>
        <p>max (0; kvmk
p(m)
(M + 1) exp</p>
        <p>v^);
end
end
if p(m) then</p>
        <p>Report machine m as suspicious;</p>
        <p>T M 2
2(pM+2)2 ;</p>
        <p>If all machines are working properly, the norm of vm =
T1 Pt S(m; x(t)) should not be much larger than its
empirical mean. The p-value p(m) in Algorithm 1 controls this
statistic by guaranteeing a small number of false detections,
depending on the signi cance level .</p>
      </sec>
    </sec>
    <sec id="sec-5">
      <title>ONLINE DETECTOR WITH REDUCED</title>
    </sec>
    <sec id="sec-6">
      <title>COMMUNICATION</title>
      <p>We describe an online, communication-e cient version of
the latent fault detector summarized in Section 2.</p>
      <p>Detecting latent faults requires that each node must send
all performance counters measured at each time point: T
samples of C counters for each of the M machines. Beyond
bandwidth costs, processing so much data is di cult to do
on a single machine in a timely manner, due to the size and
high dimensionality of the data. We apply two techniques
to alleviate this issue.</p>
      <p>Sketching is used to reduce the amount of data sent from
each machine and processed by the coordinator. Instead of
sending all counters, each node calculates a sketch of the said
counters and sends only that. The coordinator (or
monitoring node) can then perform latent fault detection using the
sketches, rather than the original data. In addition to
reducing the communication load, this has the added bene t
of reducing the computational load, since the dimensionality
of the data is greatly reduced.</p>
      <p>
        The framework in Section 2.1 requires that counter values
be normalized during preprocessing (step 1), and this is true
as well for the sketched version3. We use the safe zone
approach [
        <xref ref-type="bibr" rid="ref17">17</xref>
        ] to monitor both the global mean and the global
variance of each counter so that they do not deviate too
much from their last known values. Each machine monitors
whether its data satis es a local constraint. If all local
constraints at all machines are satis ed, the global mean and
variance are known not to have deviated too far from their
last known values. These last known values are then used to
normalize the counter values at each node, before computing
the sketch. If there is any violation, the coordinator polls
each node for the current mean and variance, and distributes
the new global mean and variance to all nodes.
      </p>
      <p>The general pseudocode is shown in Algorithm 2 and
explained in detail below.
3.1</p>
    </sec>
    <sec id="sec-7">
      <title>Sketches</title>
      <p>
        Sketching [
        <xref ref-type="bibr" rid="ref18 ref8">18, 8</xref>
        ] is a common technique used to process
large, unpredictable data streams without having to send,
store and process all data. It reduces the size of the data,
while still enabling queries. See [
        <xref ref-type="bibr" rid="ref7">7</xref>
        ] for a recent survey of
sketched-based (and other) distributed monitoring.
      </p>
      <p>
        For our purposes, a sketch is a summary function that
takes a vector and transforms it to a smaller vector while
approximately preserving some desired property, for
example inner products [
        <xref ref-type="bibr" rid="ref1">1</xref>
        ]. We use sketches to modify our tests
to greatly reduce the amount of data that must be sent and
processed. For example, 200 counters could be reduced to 10
dimensions, achieving an immediate 95% reduction in size.
      </p>
      <p>Formally, rather than apply test S to the set of all local
counter vectors x(m; t), each machine m will rst apply a
sketching function f to its vectors, and send only the sketch
x^ = f (x(m; t)) for processing. The modi ed test S^ will
be applied to the sketches rather than the original vector:
vm = T1 Pt S^(m; x^(t)).</p>
      <p>
        One well-suited sketch is the AMS sketch [
        <xref ref-type="bibr" rid="ref1">1</xref>
        ], which
involves a random linear projection to k dimensions. In our
setting, each machine would project its counter vectors to k
dimensions using a specially constructed projection matrix:
x^(m; t) = f (x(m; t)) = Rx(m; t) where R is a random C k
matrix constructed as described in [
        <xref ref-type="bibr" rid="ref1">1</xref>
        ].
      </p>
      <p>
        The AMS sketch is general enough so that the same sketch
can be used as input to di erent tests. Because the sign test
relies on normalized directions, and since AMS sketches are
linear projections, the sign test can be applied directly to
the sketch. In other words, the sum of projected vectors is
3Automatic counter selection (part step 1) can be done in
advance, o ine, using the method described in [
        <xref ref-type="bibr" rid="ref9">9</xref>
        ].
      </p>
      <sec id="sec-7-1">
        <title>Algorithm 2: Online detection pseudocode.</title>
      </sec>
      <sec id="sec-7-2">
        <title>OFFLINE:</title>
        <p>Automatically select counters.</p>
      </sec>
      <sec id="sec-7-3">
        <title>INIT / COORDINATOR SYNC:</title>
        <p>foreach counter i in counters do</p>
        <p>Poll all nodes for mean and variance of counter i.</p>
        <p>Distribute new global mean, variance, safe zones.
end</p>
      </sec>
      <sec id="sec-7-4">
        <title>NODE at time point t:</title>
        <p>foreach counter i in counters do
if counter not in safe zone then</p>
        <p>Violation: send local mean, variance to
coordinator.</p>
        <p>Wait for new global mean and variance.
end
Let xi = value of counter i at time t .
Normalize xi with last known global mean and
variance.
end
Let x = vector of normalized counter values.
Compute sketch of x and send to coordinator.</p>
      </sec>
      <sec id="sec-7-5">
        <title>COORDINATOR at time point t:</title>
        <p>if violation for counter i then</p>
        <p>Run SYNC.
end
Receive sketches from all nodes.</p>
        <p>Compute test function S on received sketches.</p>
        <p>Add most recent test function result to vm.</p>
        <p>Subtract least recent test function result from vm.</p>
        <p>
          Calculate p-value for all machines and issue warnings.
the same as projecting the sum of the vectors. The
resulting vector is still small for healthy machines and large for
outliers. The Tukey test described in our previous work
already relies on a very similar technique, and has been shown
to be very e ective. The LOF test depends on the distance
of pairs of points. In this case, the Johnson-Lindenstrauss
lemma [
          <xref ref-type="bibr" rid="ref13">13</xref>
          ] guarantees that the projection to k = O log2M
preserves the distances within a factor of 1 . Since our
method averages T comparisons per day in the integration
phase, we can further expect that in practice the error will
be smaller.
3.1.1
        </p>
        <p>Sign Test on Linear Sketches</p>
        <p>The sign test function (1) from Section 2.2 depends only
on the normalized direction from x(m; t) to the other
vectors. Let B be the unit sphere in C dimensions. Given the
assumptions in Section 2.1, for healthy machines the
normalized directions to other machines tend to be distributed
spherically symmetric over B, resulting in the vector vm =
T1 Pt S (m; x(t)) being relatively short. Conversely, for
machines with consistently anomalous behavior, vm is a
relatively long vector.</p>
        <p>Given the sketched vectors x^(m; t) = Rx(m; t), the sign
test is still the sum of normalized directions from x(m; t),
after some transformation R. We now show that applying R to
the unit sphere B maintains this symmetrical distribution.
Let R = U DV T be the singular value decomposition of R.</p>
        <p>U and V T are unitary matrices, and D is a diagonal matrix
with positive elements. In geometrical terms, the
transformation R = U DV T is a composition of rotation, followed by
non-uniform scaling and dimensional reduction, and nally
another rotation { all of which preserve the symmetric
distribution around the origin. Therefore the transformation R
maps the unit sphere B (in C dimensions) to an ellipsoid B0
in k dimensions while preserving the symmetric distribution
around the origin.</p>
        <p>In summary, since the sign-test uses normalized directions
and R preserves their symmetry around x(m; t), we can
apply the sign test directly to the sketched vectors x^(m; t).
Moreover, the sign test p-value does not depend on the
dimensionality of the vectors, and so we can use it as is.</p>
        <p>Preliminary experiments on counter logs from a small
sample of 260 machines in a single day show that sign test
scores and p-values computed on sketched data match the
original very well. Figure 1 shows a comparison of sign test
scores based on AMS sketches to regular (centralized, or
parallel) sign test scores. The gure and linear regression
show that the scores match very well, with R2 = 0:966, very
close to 1. The sketch reduced the data size by 92% { from
123 counters to 10 dimensions. The p-values are similarly
close to the original values.
3.1.2</p>
        <p>Online Integration Using a Sliding Window</p>
        <p>The integration phase in stage 2 of the framework in
Section 2.1 computes vm = T1 Pt S (m; x(t)). Computing
S (m; x(t)) only requires the data from time t, and therefore
it is trivial to turn any test into an online test by keeping a
window of test function (S) outputs for the last T sketches
sent from the monitored machines. When new data arrives
at time t, the coordinator updates the current vm by
computing and adding T1 S (m; x^(t)), and subtracting the least
recent stored test result, T1 S (m; x^(t T 1)). The p-value
for each machine in the time window can then be computed
in the usual manner. Since the test function S need only be
computed for the most recent time, and since the sketches
are of low dimension k, processing and memory costs are
low. This allows the computation to be done on a single
coordinator machine on time, before the next round starts.
3.2</p>
      </sec>
    </sec>
    <sec id="sec-8">
      <title>Scaling By Monitoring Variance</title>
      <p>Our tests require the data to be standardized during
preprocessing: each counter should be globally centered to zero
mean and unit variance. In some settings we can assume
that a counter's mean and variance do not change much,
or that they have a daily cycle. However, we might wish to
avoid that assumption, and handle unpredictable workloads.</p>
      <p>
        We use the safe zones approach [
        <xref ref-type="bibr" rid="ref17 ref21">17, 21</xref>
        ] to monitor both
the global mean and the global variance of each counter. In
this approach, each monitored machine receives a local
constraint on its data x(m; t) from a coordinator machine, such
that if all local constraints are satis ed, the global monitored
value f (x(t)) for some function f of the global aggregate
is within a pre-de ned threshold. Violations of local
constraints are sent to the coordinator machine, which resolves
them and sends updated local constraints to participating
machines.
      </p>
      <p>
        Given the last known global mean and variance of the last
T samples, we de ne some lower and upper threshold, for
example 0.9 and 1.1 times the last known values. If there is
any violation, the coordinator polls each node for the current
mean and variance, and distributes the new global mean and
variance to all nodes. We can trade-o accuracy and
communication by adjusting the high and low thresholds when
monitoring. Violations are less likely if global mean and
variance are allowed to drift further from their last known
values { reducing communication but also decreasing
accuracy [
        <xref ref-type="bibr" rid="ref17">17</xref>
        ].
      </p>
      <p>
        We monitor each counter independently, so it is enough to
show how we monitor a single counter X. Further note that
all tests described in [
        <xref ref-type="bibr" rid="ref9">9</xref>
        ] are invariant to data translation,
and so we do not monitor the global mean explicitly.
3.2.1
      </p>
      <p>Notations</p>
      <p>The set of values of counter X over the last T times and
over M nodes (machines) is denoted by X(t). We denote
by Xi(t) the values of X at node i for the last T times
up to t. Thus E[Xi(t)] is the mean of the last T values at
node i in time t, while E[X(t)] is the global mean of the
last values at all nodes. Denote i(t) = E[Xi(t)] the local
means, and (t) = E[X(t)] the global mean. Similarly, we
denote i = E Xi(t)2 , the local mean of the squares, and
= E X(t)2 the global mean. Let V (t) = ( (t); (t)), and
Vi = ( i(t); i(t)), the global and local monitored vectors,
respectively.
3.2.2</p>
      <p>Monitoring</p>
      <p>We wish to monitor the global variance Var(X) at each
time t. Recall that:</p>
      <p>Var(X) = E X2
(E [X])2 =
2 .</p>
      <p>
        We therefore monitor the conditions L 2 H, for
some lower and upper variance thresholds L and H.
Figure 2 shows the admissible region (the region in which the
conditions hold), 0:5 2 1:5 . Following [
        <xref ref-type="bibr" rid="ref17">17</xref>
        ], we
aim to nd a convex safe zone G which is contained within
the admissible region. Since convex sets are closed under
averaging, when all local vectors are inside the safe zone,
the global mean is guaranteed to be inside as well.
      </p>
      <p>Let t = 0 be the last global synchronization time, and let
V (0) = ( (0); (0)) be the reference point, the last known
global mean and mean-of-squares, computed that time. For
each node i we de ne the local drift vector di(t) as the drift
of the current vector from the node's vector during the last
synchronization: di(t) = Vi(t) Vi(0).</p>
      <p>Since we wish to monitor that the global V (t) is within
some convex set G, we de ne equivalent local conditions on
the drift vectors. The current local vectors can be written
in terms of drift vector di: Vi(t) = Vi(0) + di(t). Note that
the global vector is the mean of the local vectors, and can
therefore be written as the mean of drifts and the reference
point:</p>
      <p>V (t) =</p>
      <p>Let Wi(t) = V (0) + di(t) be the local drift from the last
reference point. Note that V (t) = M1 Pi Wi, recall G is
convex, and from (2) we arrive at the local conditions: if
8i; Wi 2 G then V (t) 2 G.</p>
      <p>To monitor that the variance is between L and H, we
derive separate safe zones: one for variance above L and
another for variance below H. As long as the local
conditions for both safe zones are maintained in all nodes, we are
guaranteed that the variance is within the allowed range.
Variance Above Lower Threshold. We wish to de ne a
convex safe zone GL so that as long as V (t) 2 GL then
Var(X) L. This corresponds to monitoring that 2
L, which is already a convex set { the area above a parabola
{ and can be directly used as safe zone. Therefore the local
condition for each node i is trivial: Ii(t) 2 GL: Ii(t) =
V (0) + di(t) = (a; b) and monitor that b a2 L.
Variance Below Upper Threshold. We wish to de ne
a convex safe zone G so that as long as V (t) 2 G then
Var(X) H. This area is the area below a parabola, which
is not a convex set. However, we can nd a tangent
halfplane I below this parabola. This half-plane is a convex set,
and since I G, then as long as V (t) 2 I, V (t) 2 G and
therefore Var(X) H.</p>
      <p>We use the reference point V (0) to nd the optimal
hyperplane. The thresholds H and L are reset during
synchronization, so obviously V (0) 2 G. We can choose any
halfspace I such that V (0) 2 I, but to avoid future unnecessary
synchronization we choose I such that V (0) is far from the
boundary of G. Doing so ensures that drift has to be large
to cause a violation. Consequently, we choose I as the
tangent at point P , where P is the closest point to V (0) on the
parabola 2 = H, and the local condition is Wi 2 I. We
can nd P numerically, or by minimizing the distance from
the parabola to V (0). For example, if V (0) = (0:5; 1) and
H = 1:5, then the closest point on the parabola is 0:237.
This yields the point P = (0:237; 1:556), and nally the
induced safe zone I: the half-plane 0:474 &lt; 1:443 .
Figure 3(a) shows V (0), P and the resulting safe zone, and
Figure 3(b) shows the intersection with the safe zone for the
lower limit L = 0:5.
2.5
λ 1.5</p>
      <p>1
λ 1.5</p>
      <p>1</p>
      <p>P</p>
      <p>V(0)
V(0)
0
-3 -2 -1 0 1 2 3
μ
3.2.3</p>
      <p>Handling Violations</p>
      <p>If one of the local conditions Wj 2 G is violated, it may be
because Var(X) is no longer in the range, or due to a false
alarm. The simplest way to deal with a violation is to
perform a global synchronization: each node sends its current
Vi(t) to the coordinator. The coordinator \resets the time"
to t = 0, computes the new global reference point V (0), and
sends it to the nodes, where it is used for monitoring and
scaling.</p>
      <p>
        In terms of communication, our synchronizations are fairly
inexpensive. Each node sends only two numbers per counter
( and ), rather than the entire time window of T samples.
They also improve the accuracy of scaling, since nodes have
fresh global mean and variance. There are safe zone
techniques that allow partial synchronization for further
communication reduction, for example by balancing a node with
local violation with another node that has enough slack [
        <xref ref-type="bibr" rid="ref2">2</xref>
        ].
      </p>
    </sec>
    <sec id="sec-9">
      <title>FUTURE WORK</title>
      <p>
        This work uses sketching and safe zones to adapt the
latent fault detector in [
        <xref ref-type="bibr" rid="ref9">9</xref>
        ] to a streaming setting, resulting in
an online, communication-e cient outlier detector for
common scale-out systems. Preliminary results show that the
adapted detector obtains very similar results to those of
the original latent fault detector for the sign test. Future
work will concentrate on adapting additional tests,
evaluating the detector on real-world systems, and exploring the
communication-accuracy trade-o .
      </p>
    </sec>
    <sec id="sec-10">
      <title>ACKNOWLEDGMENTS</title>
      <p>The research leading to these results has received funding
from the European Union's Seventh Framework Programme
under grant agreement No 255951.</p>
    </sec>
  </body>
  <back>
    <ref-list>
      <ref id="ref1">
        <mixed-citation>
          [1]
          <string-name>
            <given-names>N.</given-names>
            <surname>Alon</surname>
          </string-name>
          ,
          <string-name>
            <given-names>Y.</given-names>
            <surname>Matias</surname>
          </string-name>
          , and
          <string-name>
            <given-names>M.</given-names>
            <surname>Szegedy</surname>
          </string-name>
          .
          <article-title>The space complexity of approximating the frequency moments</article-title>
          .
          <source>Journal of Computer and System Sciences</source>
          ,
          <year>1999</year>
          .
        </mixed-citation>
      </ref>
      <ref id="ref2">
        <mixed-citation>
          [2]
          <string-name>
            <given-names>D.</given-names>
            <surname>Ben-David</surname>
          </string-name>
          .
          <article-title>Violation resolution in distributed stream networks</article-title>
          .
          <source>Master's thesis</source>
          ,
          <string-name>
            <surname>Technion</surname>
            <given-names>I.I.T</given-names>
          </string-name>
          ,
          <year>2012</year>
          .
        </mixed-citation>
      </ref>
      <ref id="ref3">
        <mixed-citation>
          [3]
          <string-name>
            <given-names>P.</given-names>
            <surname>Bod k</surname>
          </string-name>
          , M. Goldszmidt,
          <string-name>
            <given-names>A.</given-names>
            <surname>Fox</surname>
          </string-name>
          ,
          <string-name>
            <given-names>D. B.</given-names>
            <surname>Woodard</surname>
          </string-name>
          , and
          <string-name>
            <given-names>H.</given-names>
            <surname>Andersen</surname>
          </string-name>
          .
          <article-title>Fingerprinting the datacenter: Automated classi cation of performance crises</article-title>
          .
          <source>In Proc. EuroSys</source>
          ,
          <year>2010</year>
          .
        </mixed-citation>
      </ref>
      <ref id="ref4">
        <mixed-citation>
          [4]
          <string-name>
            <given-names>G.</given-names>
            <surname>Bronevetsky</surname>
          </string-name>
          ,
          <string-name>
            <given-names>I.</given-names>
            <surname>Laguna</surname>
          </string-name>
          , B. De Supinski, and
          <string-name>
            <given-names>S.</given-names>
            <surname>Bagchi</surname>
          </string-name>
          .
          <article-title>Automatic fault characterization via abnormality-enhanced classi cation</article-title>
          .
          <source>In Proc. DSN</source>
          ,
          <year>2012</year>
          .
        </mixed-citation>
      </ref>
      <ref id="ref5">
        <mixed-citation>
          [5]
          <string-name>
            <given-names>H.</given-names>
            <surname>Chen</surname>
          </string-name>
          , G. Jiang, and
          <string-name>
            <given-names>K.</given-names>
            <surname>Yoshihira</surname>
          </string-name>
          .
          <article-title>Failure detection in large-scale internet services by principal subspace mapping</article-title>
          .
          <source>IEEE Trans. Knowl</source>
          . Data Eng.,
          <year>2007</year>
          .
        </mixed-citation>
      </ref>
      <ref id="ref6">
        <mixed-citation>
          [6]
          <string-name>
            <given-names>I.</given-names>
            <surname>Cohen</surname>
          </string-name>
          ,
          <string-name>
            <given-names>M.</given-names>
            <surname>Goldszmidt</surname>
          </string-name>
          ,
          <string-name>
            <given-names>T.</given-names>
            <surname>Kelly</surname>
          </string-name>
          , and
          <string-name>
            <given-names>J.</given-names>
            <surname>Symons</surname>
          </string-name>
          .
          <article-title>Correlating instrumentation data to system states: A building block for automated diagnosis and control</article-title>
          .
          <source>In Proc. OSDI</source>
          ,
          <year>2004</year>
          .
        </mixed-citation>
      </ref>
      <ref id="ref7">
        <mixed-citation>
          [7]
          <string-name>
            <surname>G. Cormode.</surname>
          </string-name>
          <article-title>The continuous distributed monitoring model</article-title>
          .
          <source>SIGMOD Rec</source>
          .,
          <year>2013</year>
          .
        </mixed-citation>
      </ref>
      <ref id="ref8">
        <mixed-citation>
          [8]
          <string-name>
            <given-names>G.</given-names>
            <surname>Cormode</surname>
          </string-name>
          and
          <string-name>
            <given-names>M.</given-names>
            <surname>Garofalakis</surname>
          </string-name>
          .
          <article-title>Sketching probabilistic data streams</article-title>
          .
          <source>In SIGMOD</source>
          ,
          <year>2007</year>
          .
        </mixed-citation>
      </ref>
      <ref id="ref9">
        <mixed-citation>
          [9]
          <string-name>
            <given-names>M.</given-names>
            <surname>Gabel</surname>
          </string-name>
          ,
          <string-name>
            <given-names>A.</given-names>
            <surname>Schuster</surname>
          </string-name>
          , R.-G. Bachrach, and
          <string-name>
            <given-names>N.</given-names>
            <surname>Bjorner</surname>
          </string-name>
          .
          <article-title>Latent fault detection in large scale services</article-title>
          .
          <source>In Proc. DSN</source>
          ,
          <year>2012</year>
          .
        </mixed-citation>
      </ref>
      <ref id="ref10">
        <mixed-citation>
          [10]
          <string-name>
            <given-names>C.</given-names>
            <surname>Huang</surname>
          </string-name>
          ,
          <string-name>
            <surname>I. Cohen</surname>
          </string-name>
          ,
          <string-name>
            <given-names>J.</given-names>
            <surname>Symons</surname>
          </string-name>
          , and
          <string-name>
            <given-names>T.</given-names>
            <surname>Abdelzaher</surname>
          </string-name>
          .
          <article-title>Achieving scalable automated diagnosis of distributed systems performance problems</article-title>
          .
          <source>Technical report, HP Labs</source>
          ,
          <year>2007</year>
          .
        </mixed-citation>
      </ref>
      <ref id="ref11">
        <mixed-citation>
          [11]
          <string-name>
            <given-names>M.</given-names>
            <surname>Isard</surname>
          </string-name>
          .
          <article-title>Autopilot: automatic data center management</article-title>
          .
          <source>SIGOPS Oper. Syst. Rev.</source>
          ,
          <year>2007</year>
          .
        </mixed-citation>
      </ref>
      <ref id="ref12">
        <mixed-citation>
          [12]
          <string-name>
            <given-names>M.</given-names>
            <surname>Isard</surname>
          </string-name>
          ,
          <string-name>
            <given-names>M.</given-names>
            <surname>Budiu</surname>
          </string-name>
          ,
          <string-name>
            <given-names>Y.</given-names>
            <surname>Yu</surname>
          </string-name>
          ,
          <string-name>
            <given-names>A.</given-names>
            <surname>Birrell</surname>
          </string-name>
          , and
          <string-name>
            <given-names>D.</given-names>
            <surname>Fetterly</surname>
          </string-name>
          .
          <article-title>Dryad: distributed data-parallel programs from sequential building blocks</article-title>
          .
          <source>In Proc. EuroSys</source>
          ,
          <year>2007</year>
          .
        </mixed-citation>
      </ref>
      <ref id="ref13">
        <mixed-citation>
          [13]
          <string-name>
            <given-names>W.</given-names>
            <surname>Johnson</surname>
          </string-name>
          and
          <string-name>
            <given-names>J.</given-names>
            <surname>Lindenstrauss</surname>
          </string-name>
          .
          <article-title>Extensions of Lipschitz mappings into a Hilbert space. In Conference in modern analysis and probability (New Haven, Conn</article-title>
          .,
          <year>1982</year>
          ),
          <source>Contemporary Mathematics</source>
          .
          <year>1984</year>
          .
        </mixed-citation>
      </ref>
      <ref id="ref14">
        <mixed-citation>
          [14]
          <string-name>
            <given-names>M. P.</given-names>
            <surname>Kasick</surname>
          </string-name>
          ,
          <string-name>
            <given-names>J.</given-names>
            <surname>Tan</surname>
          </string-name>
          ,
          <string-name>
            <given-names>R.</given-names>
            <surname>Gandhi</surname>
          </string-name>
          , and
          <string-name>
            <given-names>P.</given-names>
            <surname>Narasimhan</surname>
          </string-name>
          .
          <article-title>Black-box problem diagnosis in parallel le systems</article-title>
          .
          <source>In Proc. FAST</source>
          ,
          <year>2010</year>
          .
        </mixed-citation>
      </ref>
      <ref id="ref15">
        <mixed-citation>
          [15]
          <string-name>
            <given-names>S.</given-names>
            <surname>Kavulya</surname>
          </string-name>
          ,
          <string-name>
            <given-names>S.</given-names>
            <surname>Daniels</surname>
          </string-name>
          ,
          <string-name>
            <given-names>K.</given-names>
            <surname>Joshi</surname>
          </string-name>
          ,
          <string-name>
            <given-names>M.</given-names>
            <surname>Hiltunen</surname>
          </string-name>
          ,
          <string-name>
            <given-names>R.</given-names>
            <surname>Gandhi</surname>
          </string-name>
          , and
          <string-name>
            <given-names>P.</given-names>
            <surname>Narasimhan</surname>
          </string-name>
          . Draco:
          <article-title>Statistical diagnosis of chronic problems in large distributed systems</article-title>
          .
          <source>In Proc. DSN</source>
          ,
          <year>2012</year>
          .
        </mixed-citation>
      </ref>
      <ref id="ref16">
        <mixed-citation>
          [16]
          <string-name>
            <given-names>S.</given-names>
            <surname>Kavulya</surname>
          </string-name>
          ,
          <string-name>
            <given-names>R.</given-names>
            <surname>Gandhi</surname>
          </string-name>
          , and
          <string-name>
            <given-names>P.</given-names>
            <surname>Narasimhan</surname>
          </string-name>
          .
          <article-title>Gumshoe: Diagnosing performance problems in replicated le-systems</article-title>
          .
          <source>In Proc. SRDS</source>
          ,
          <year>2008</year>
          .
        </mixed-citation>
      </ref>
      <ref id="ref17">
        <mixed-citation>
          [17]
          <string-name>
            <given-names>D.</given-names>
            <surname>Keren</surname>
          </string-name>
          ,
          <string-name>
            <surname>I. Sharfman</surname>
          </string-name>
          ,
          <string-name>
            <given-names>A.</given-names>
            <surname>Schuster</surname>
          </string-name>
          ,
          <article-title>and</article-title>
          <string-name>
            <given-names>A.</given-names>
            <surname>Livne</surname>
          </string-name>
          .
          <article-title>Shape sensitive geometric monitoring. Knowledge and Data Engineering</article-title>
          , IEEE Transactions on,
          <year>2012</year>
          .
        </mixed-citation>
      </ref>
      <ref id="ref18">
        <mixed-citation>
          [18]
          <string-name>
            <given-names>S.</given-names>
            <surname>Muthukrishnan</surname>
          </string-name>
          .
          <article-title>Data streams: Algorithms and applications</article-title>
          .
          <source>Foundations and Trends in Theoretical Computer Science</source>
          ,
          <year>2005</year>
          .
        </mixed-citation>
      </ref>
      <ref id="ref19">
        <mixed-citation>
          [19]
          <string-name>
            <given-names>A. J.</given-names>
            <surname>Oliner</surname>
          </string-name>
          ,
          <string-name>
            <given-names>A.</given-names>
            <surname>Aiken</surname>
          </string-name>
          , and
          <string-name>
            <given-names>J.</given-names>
            <surname>Stearley</surname>
          </string-name>
          .
          <article-title>Alert detection in system logs</article-title>
          .
          <source>In Proc. ICDM</source>
          ,
          <year>2008</year>
          .
        </mixed-citation>
      </ref>
      <ref id="ref20">
        <mixed-citation>
          [20]
          <string-name>
            <given-names>D.</given-names>
            <surname>Pelleg</surname>
          </string-name>
          ,
          <string-name>
            <given-names>M.</given-names>
            <surname>Ben-Yehuda</surname>
          </string-name>
          ,
          <string-name>
            <given-names>R.</given-names>
            <surname>Harper</surname>
          </string-name>
          ,
          <string-name>
            <given-names>L.</given-names>
            <surname>Spainhower</surname>
          </string-name>
          , and
          <string-name>
            <given-names>T.</given-names>
            <surname>Adeshiyan</surname>
          </string-name>
          .
          <article-title>Vigilant: out-of-band detection of failures in virtual machines</article-title>
          .
          <source>SIGOPS Oper. Syst. Rev.</source>
          ,
          <year>2008</year>
          .
        </mixed-citation>
      </ref>
      <ref id="ref21">
        <mixed-citation>
          [21]
          <string-name>
            <given-names>I.</given-names>
            <surname>Sharfman</surname>
          </string-name>
          ,
          <string-name>
            <given-names>A.</given-names>
            <surname>Schuster</surname>
          </string-name>
          , and
          <string-name>
            <given-names>D.</given-names>
            <surname>Keren</surname>
          </string-name>
          .
          <article-title>A geometric approach to monitoring threshold functions over distributed data streams</article-title>
          .
          <source>TODS</source>
          ,
          <year>2007</year>
          .
        </mixed-citation>
      </ref>
      <ref id="ref22">
        <mixed-citation>
          [22]
          <string-name>
            <given-names>W.</given-names>
            <surname>Xu</surname>
          </string-name>
          ,
          <string-name>
            <given-names>L.</given-names>
            <surname>Huang</surname>
          </string-name>
          ,
          <string-name>
            <given-names>A.</given-names>
            <surname>Fox</surname>
          </string-name>
          ,
          <string-name>
            <given-names>D.</given-names>
            <surname>Patterson</surname>
          </string-name>
          , and
          <string-name>
            <given-names>M. I.</given-names>
            <surname>Jordan</surname>
          </string-name>
          .
          <article-title>Detecting large-scale system problems by mining console logs</article-title>
          .
          <source>In Proc. SOSP</source>
          ,
          <year>2009</year>
          .
        </mixed-citation>
      </ref>
      <ref id="ref23">
        <mixed-citation>
          [23]
          <string-name>
            <given-names>S.</given-names>
            <surname>Zhang</surname>
          </string-name>
          , I. Cohen,
          <string-name>
            <given-names>M.</given-names>
            <surname>Goldszmidt</surname>
          </string-name>
          ,
          <string-name>
            <given-names>J.</given-names>
            <surname>Symons</surname>
          </string-name>
          ,
          <article-title>and</article-title>
          <string-name>
            <given-names>A.</given-names>
            <surname>Fox</surname>
          </string-name>
          .
          <article-title>Ensembles of models for automated diagnosis of system performance problems</article-title>
          .
          <source>In Proc. DSN</source>
          ,
          <year>2005</year>
          .
        </mixed-citation>
      </ref>
    </ref-list>
  </back>
</article>