=Paper= {{Paper |id=Vol-2083/paper-17 |storemode=property |title=A Distributed Online Learning Approach for Pattern Prediction over Movement Event Streams with Apache Flink |pdfUrl=https://ceur-ws.org/Vol-2083/paper-17.pdf |volume=Vol-2083 |authors=Ehab Qadah,Michael Mock,Elias Alevizos,Georg Fuchs |dblpUrl=https://dblp.org/rec/conf/edbt/QadahMAF18 }} ==A Distributed Online Learning Approach for Pattern Prediction over Movement Event Streams with Apache Flink== https://ceur-ws.org/Vol-2083/paper-17.pdf
       A Distributed Online Learning Approach for Pattern
    Prediction over Movement Event Streams with Apache Flink
                                 Ehab Qadah                                                                        Michael Mock
                           Fraunhofer IAIS                                                                     Fraunhofer IAIS
                       Sankt Augustin, Germany                                                            Sankt Augustin, Germany
                     Ehab.Qadah@iais.fraunhofer.de                                                     Michael.Mock@iais.fraunhofer.de

                               Elias Alevizos                                                                       Georg Fuchs
                           NCSR "Demokritos"                                                                   Fraunhofer IAIS
                             Athens, Greece                                                               Sankt Augustin, Germany
                     alevizos.elias@iit.demokritos.gr                                                   Georg.Fuchs@iais.fraunhofer.de

ABSTRACT                                                                                  data with low latency is challenging, since data streams are large
In this paper, we present a distributed online prediction system                          and distributed in nature and continuously arrive at a high rate.
for user-defined patterns over multiple massive streams of move-                             In this paper, we present the design and implementation of an
ment events, built using the general purpose stream processing                            online, distributed and scalable pattern prediction system over
framework Apache Flink. The proposed approach is based on                                 multiple, massive streams of events. More precisely, we consider
combining probabilistic event pattern prediction models on multi-                         event streams related to trajectories of moving objects (i.e., ves-
ple predictor nodes with a distributed online learning protocol in                        sels). The proposed approach is based on a novel method that
order to continuously learn the parameters of a global prediction                         combines a distributed online prediction protocol [8, 16] with
model and share them among the predictors in a communication-                             an event forecasting method based on Markov chains [2]. It is
efficient way. Our approach enables the collaborative learning                            implemented on top of the Big Data framework for stream pro-
between the predictors (i.e., "learn from each other"), thus the                          cessing Apache Flink [13]. We evaluate our proposed system over
learning rate is accelerated with less data for each predictor. The                       real-world data streams of moving vessels, which are provided
underlying model provides online predictions about when a pat-                            in the context of the datAcron project1 .
tern (i.e., a regular expression over the event types) is expected                           The rest of the paper is organized as follows. We discuss the
to be completed within each event stream. We describe the dis-                            related work and used frameworks in Section 2. In Section 3,
tributed architecture of the proposed system, its implementation                          we describe the problem of pattern prediction, our proposed
in Flink, and present experimental results over real-world event                          approach, and the architecture of our system. The implementa-
streams related to trajectories of moving vessels.                                        tion details on top of Flink are presented in Section 4 and the
                                                                                          experimental results in Section 5. We conclude in Section 6.

1    INTRODUCTION                                                                         2 RELATED WORK AND BACKGROUND
In recent years, technological advances have led to a growing                             2.1 Related work
availability of massive amounts of continuous streaming data
                                                                                              Pattern prediction over event streams. The task of forecasting
(i.e., data streams observing events) in many application domains
                                                                                          over time-evolving streams of data can be formulated in various
such as social networks [23], Internet of Things (IoT) [24], and
                                                                                          ways and with varying assumptions. One common way to for-
maritime surveillance [31]. The ability to detect and predict the
                                                                                          malize this task is to assume that the stream is a time-series of
full matches of a pattern of interest (e.g., a certain sequence of
                                                                                          numerical values, and the goal is to forecast at each time point
events), defined by a domain expert, is typically important for
                                                                                          n the values at some future points n + 1, n + 2, etc., (or even
operational decision making tasks in the respective domains.
                                                                                          the output of some function of future values). This is the task of
    An event stream is an unbounded collection of time-ordered
                                                                                          time-series forecasting [25]. Another way to formalize this task is
data observations in the form of a tuple of attributes that is
                                                                                          to view streams as sequences of events, i.e., tuples with multiple,
composed of a value from finite event types along with other
                                                                                          possibly categorical, attributes, like event type, timestamp, etc.,
categorical and numerical attributes. In this work, we deal with
                                                                                          and the goal is to predict future events or patterns of events. In
movement event streams. For instance, in the context of mar-
                                                                                          this paper, we focus on this latter definition of forecasting (event
itime surveillance the event stream of a moving vessel consists
                                                                                          pattern forecasting).
of spatio-temporal and kinematic information along with the
                                                                                              A substantial body of work on event forecasting comes from
vessel’s identification and its trajectory related events, based
                                                                                          the field of temporal pattern mining where events are defined as
on the automatic identification system (AIS) [28] messages that
                                                                                          2-tuples of the form (EventType, Timestamp). The ultimate goal
are continuously sent by the vessel. Therefore, leveraging event
                                                                                          is to extract patterns of events in the form either of association
patterns prediction over real-time streams of moving vessels is
                                                                                          rules [1] or frequent episode rules [22]. These methods have been
useful to alert maritime operation managers about suspicious
                                                                                          extended in order to be able to learn not only rules for detecting
activities (e.g., fast sailing vessels near ports, or illegal fishing)
                                                                                          event patterns but also rules for predicting events. For example,
before they happen. However, processing real-time streaming
                                                                                          in [32], a variant of association rule mining is where the goal is
© 2018 Copyright held by the owner/author(s). Published in the Workshop
                                                                                          to extract sets of event types that frequently lead to a rare, target
Proceedings of the EDBT/ICDT 2018 Joint Conference (March 26, 2018, Vienna,               event within a temporal window.
Austria) on CEUR-WS.org (ISSN 1613-0073). Distribution of this paper is permitted
under the terms of the Creative Commons license CC-by-nc-nd 4.0.                          1 http://www.datacron-project.eu/




                                                                                    109
    In [19], a probabilistic model is presented for calculating the         is treated as a special case of the streaming one (i.e., finite stream).
probability of the immediately next event in the stream. This               Flink’s software stack includes the DataStream and DataSet APIs
is achieved by using standard frequent episode discovery algo-              for processing infinite and finite data, respectively. These two
rithms and combining them with Hidden Markov Models and                     core APIs are built on top of Flink’s core dataflow engine and
mixture models. The framework of episode rules is employed                  provide operations on data streams or sets such as mapping,
in [9] as well. The output of the proposed algorithms is a set              filtering, grouping, etc.
of predictive rules whose antecedent is minimal (in number of                   The two main data abstractions of Flink are DataStream and
events) and temporally distant from the consequent. In [35] a               DataSet, they represent read-only collections of data elements.
set of algorithms is proposed that target batch online mining of            The list of elements is bounded (i.e., finite) in DataSet, while it is
sequential patterns, without maintaining exact frequency counts.            unbounded (i.e., infinite) in the case of DataStream. Flink’s core
As the stream is consumed, the learned patterns can be used to              is a distributed streaming dataflow engine. Each Flink program
test whether a prefix matches the last events seen in the stream,           is represented by a data-flow graph (i.e., directed acyclic graph -
indicating a possibility of occurrence for events that belong to            DAG) that gets executed by Flink’s dataflow engine [6]. The data
the suffix of the rule.                                                     flow graphs are composed of stateful operators and intermediate
    Event forecasting has also attracted some attention from the            data stream partitions. The execution of each operator is handled
filed of Complex Event Processing (see [7] for a review of Com-             by multiple parallel instances whose number is determined by
plex Event Processing). One such early approach is presented                the parallelism level. Each parallel operator instance is executed
in [26]. Complex event patterns are converted to automata, and              in an independent task slot on a machine within a cluster of
subsequently, Markov chains are used in order to estimate when              computers [13].
a pattern is expected to be fully matched. A similar approach
                                                                               Apache Kafka. Apache Kafka is a scalable, fault-tolerant, and
is presented in [2], where again automata and Markov chains
                                                                            distributed streaming framework/messaging system [11]. It al-
are employed in order to provide (future) time intervals during
                                                                            lows to publish and subscribe to arbitrary data streams, which
which a match is expected with a probability above a confidence
                                                                            are managed in different categories (i.e., topics) and partitioned
threshold.
                                                                            in the Kafka cluster. The Kafka Producer API provides the abil-
    Distributed Online Learning. In recent years, the problem of            ity to publish a stream of messages to a topic. These messages
distributed online learning has received increased attention and            can then be consumed by applications, using the Consumer API
has been studied in [8, 16, 18, 33, 34]. A distributed online mini-         that allows them to read the published data stream in the Kafka
batch prediction approach over multiple data streams has been               cluster. In addition, the streams of messages are distributed and
proposed in [8]. This approach is based on a static synchroniza-            load balanced between the multiple receivers within the same
tion method. The learners periodically communicate their local              consumer group for the sake of scalability.
models with a central coordinator unit after consuming a fixed
number of input samples/events (i.e., batch size b), in order to            3 SYSTEM OVERVIEW
create a global model and share it between all learners. This work          3.1 Pattern prediction on a single stream
has been extended in [16] by introducing a dynamic synchroniza-
                                                                            For our work presented in this paper, we use the approach pre-
tion scheme that reduces the required communication overhead.
                                                                            sented in [2]. For the sake of self-containment, we briefly describe
It can do so by making the local learners communicate their mod-
                                                                            this approach in the following, first assuming that only a single
els only if they diverge from a reference model. In this work, we
                                                                            stream is consumed and then adjusting for the case of multiple
employ this protocol with event patterns prediction models over
                                                                            streams. We follow the terminology of [3, 21, 35] to formalize
multiple event streams.
                                                                            the problem we tackle.

2.2 Technological Background                                                   3.1.1 Problem formulation. We define an input event and a
In the last years, many systems for large-scale and distributed             stream of input events as follows:
stream processing have been proposed, including Spark Stream-                    Definition 3.1. Each event is defined as a tuple of attributes
ing [12], Apache Storm [14] and Apache Flink [13]. These frame-             ei = (id, type, τ , a 1 , a 2 ....., an ), where type is the event type at-
works can ingest and process real-time data streams, published              tribute that takes a value from a set of finite event types/symbols
from different distributed message queuing platforms, such as               Σ, τ represents the time when the event tuple was created, the
Apache Kafka [11] or Amazon Kinesis [5]. In this work, we im-               a 1 , a 2 , ..., an are spatial or other contextual features (e.g., speed);
plemented our system in Flink and Kafka (see Section 4).                    these features are varying from one application domain to an-
   In the datAcron project, the Flink streaming processing engine           other. The attribute id is a unique identifier that connects the
has been chosen as a primary platform for supporting the stream-            event tuple to an associated domain object.
ing operations, based on an internal comparative evaluation of
several streaming platforms. Hence, we used it to implement our                Definition 3.2. A stream s = ⟨e 1 , e 3 , ..., et , ...⟩ is a time-ordered
system. A predecessor distributed online learning framework has             sequence of events.
already been implemented in the FERARI project [10] based on                   A user-defined pattern P is given in the form of a regular
Apache Storm.                                                               expression (i.e., using operators for sequence, disjunction, and
                                                                            iteration) over Σ (i.e., event types) [2]. More formally, a pattern
   Apache Flink. Apache Flink is an open source project that
                                                                            is given through the following grammar:
provides a large-scale, distributed, and stateful stream processing
platform [6]. Flink is one of the most recent and pioneering Big               Definition 3.3. P := E | P1 ; P2 |P1 ∨ P2 | P1∗ , where E ∈ Σ is
Data processing frameworks. It provides processing models for               a constant event type. ; stands for sequence, ∨ for disjunction
both streaming and batch data, where the batch processing model             and ∗ for Kleene − ∗. The pattern P := E is matched by reading




                                                                      110
an event ei iff ei .type = E. The other cases are matched as in
standard automata theory.                                                                                             a
                                                                                                                                                d                a
   The problem at hand may then be stated as follows: given                                                                              a
                                                                                                                           1            a                                     a
a stream s of low-level events and a pattern P, the goal is to                                            a                        b                b            5
                                                                                                                  b
estimate at each new event arrival the number of future events
that we will need to wait for until the pattern is satisfied (and                                                                           2           b             c
                                                                                                                  c           a
therefore a full match is detected).                                                  start           0                       c
                                                                                                                                                b
                                                                                                                                                        b             6
                                                                                                                                       bd
   3.1.2 Proposed approach. As a first step, event patterns are                                               c                                             d                     d
converted to deterministic finite automata (DFA) through stan-                                                            3
                                                                                                                                  c d               4
dard conversion algorithms [15]. As an example, see Figure 1a for
the DFA of the simple sequential pattern P = a; d; c and an alpha-
                                                                                                           d
bet Σ = {a, b, c, d} (note that the DFA has no dead states since                                                                   d                        c
we need to handle streams and not strings). The next step is to
                                                                                                                                  c
derive a Markov chain that will be able to provide a probabilistic                                                         (a) DFAΣ∗ ;P
description of the DFA’s run-time behavior.
   Towards this goal, we use Pattern Markov Chains, as was                                                                    P(d |a)
                                                                                          P(a|a)
proposed in [27]. Under the assumption that the input stream is
generated by an m-order Markov source, and after performing a                                                           P(a|d)
transformation step on the initial DFA to handle the mt h order                               1                   P(a|b)                                                  5
                                                                                                                         P(b |b)
(see [27] for more details). It can be shown that there is a direct                                                                             P(b|d)
mapping of the states of the DFA to states of a Markov chain and                                          P(b |a)
                                                                                                                                   2                            P(c |d)
the transitions of the DFA to transitions of the Markov chain.
The transition probabilities are then conditional probabilities on                                    P(a|c)
                                                                                                                          P(a|d)                                                  P(d |d)
the event types.                                                                                                                                                          6
                                                                                       P(c |a)                                                P(b |d)
   We call such a derived Markov chain a Pattern Markov Chain                                              P(c |b)                      P(d |b)
(PMC) of order m and denote by PMC m       , where P is the initial                                                                                                       1
                                         P                                                                     P(b |c)
pattern and m the assumed order. As an example, see Figure 1b,                                    3                                                               4
which depicts the PMC of order 1 for the generated DFA of
                                                                                                                           P(d |c)
Figure 1a.                                                                              P(c |c)                                                                       P(d |d)
                                                                                                                           P(c |d)
                                                                                                                                     1
                                                                                                                              (b) PMCP
    After constructing a PMC, we can use it to calculate the so-
called waiting-time distributions. Given a specific state of the                Figure 1: DFA and PMC for P = a; d; c with Σ = {a, b, c, d},
PMC, a waiting-time distribution gives us the probability of reach-             and order m = 1.
ing a set of absorbing states in n transition from now (absorbing
states are states with self-loops and probability equal to 1.0). By
mapping the final states of the initial DFA to absorbing states
of the PMC (see again Figure 1). Therefore, we can calculate the
probability of reaching a final state, or, in other words, of detect-
ing a full match of the original regular expression in n events
from now.
    In order to estimate the final forecasts, another step is required,
since our aim is not to provide a single future point with the high-
est probability but an interval. Predictions are given in the form
of intervals, as I = (start, end). The meaning of such an interval
is that the DFA is expected to reach a final state sometime in the
future between the start and end with probability at least some
constant threshold θ f c (provided by the user). These intervals are
estimated by a single-pass algorithm that scans a waiting-time
distribution and finds the smallest (in terms of length) interval                   (a) Waiting-time distribution.                                              (b) Prediction intervals.
that has probability exceeds this threshold (θ f c ). For example,
Figure 2a shows the waiting-time distributions for the non-final                Figure 2: Example of how prediction intervals are pro-
states of the DFA in Figure 1, and the computed prediction inter-               duced. P = a; d; c, Σ = {a, b, c, d}, m = 1, θ fc = 0.5.
vals are depicted in Figure 2b.
    The method described above assumes that we know the (pos-
sibly conditional) occurrence probabilities of the various event
types appearing in a stream (as would be the case with syntheti-                part of the stream to obtain the maximum-likelihood estimators
cally generated streams). However, this is not always the case in               for the transition probabilities [4]. If Π is the transition matrix of
real-world situations. Therefore, it is crucial for a system imple-             a Markov chain with a set of states Q, πi, j the transition probabil-
menting this method to have the capability to learn the values of               ity from state i to state j, ni, j the number of observed transitions
the PMC’s transition matrix. One way to do this is to use some                  from state i to state j, then the maximum likelihood estimator




                                                                          111
for πi, j is given by:                                                              computed from the already received models less or equal than
                                  ni, j    ni, j                                    the divergence threshold ∆. Finally, the aggregated model fˆ is
                      π̂i, j = Í         =                                          sent back to the predictor nodes that sent their models after the
                                       n
                                k ∈Q i,k    ni
                                                                                    violation or have been queried by the coordinator.
Executing this learning step on a single node might require a
vast amount of time until we arrive at a sufficiently good model.
                                                                                     Algorithm 1: Communication-efficient Distributed Online
In this paper, we present a distributed method for learning the
                                                                                     Learning Protocol
transition probability matrix.
                                                                                       Predictor node ni : at observing event e j
3.2 Pattern prediction on multiple streams                                                update the prediction model parameters fi and provide a
     3.2.1 Problem formulation. Let O = {o 1 , ..., ok } be a set of K                    prediction service ;
objects (i.e., moving objects) and S = {s 1 , ..., sk } a set of real-time                if j mod b = 0 and ∥ fi − fr ∥ 2 > ∆ then
streams of events, where si is generated by the object oi . Let P be                            send fi to the Coordinator (violation) ;
a user-defined pattern which we want to apply to every stream
si , i.e., each object will have its own DFA.                                          Coordinator:
     The setting that is considered in this work is then described                        receive local models with violation B = { fi }i=1 m ;
in the following: we have K input event streams S and a system                                                  1  Í             ˆ 2
                                                                                          while |B| , k and |B | fi ∈Π ∥ fi − f ∥ > ∆ do
consisting of K distributed predictor nodes n 1 , n 2 ..., nk , each of                       add other nodes that have not reported violation for
which consumes an input event stream si ∈ S. The goal is to                                   their models B ← { fl : fl < B and l ∈ [k]} ;
provide timely predictions and be able to do this at large-scale.                             receive models from nodes in B;
Each node ni handles a single event stream si associated with a
moving object oi ∈ O. In addition, it maintains a local prediction                        compute a new global model fˆ ;
model fi for the user-defined pattern P. The fi model provides                            send fˆ to all the predictors in B and set f 1 . . . fm = fˆ;
the online prediction about the future full match of the pattern                           if |B| = k then
P in si for each new arriving event tuple.                                                     set a new reference model fr ← fˆ ;
     In short, we have multiple running instances of an online
prediction algorithm on distributed nodes for multiple input
event streams. More specifically, the input to our system consists                      This protocol was introduced for linear models, and has been
of massive streams of events that describe trajectories of moving                   extended to handle kernelized online learning models [17]. We
vessels in the context of maritime surveillance, where there is                     also employ this protocol for the pattern prediction model, which
one predictor node for each vessel’s event stream.                                  is internally based on the PMC PMC m  P
                                                                                                                            . This allows the distributed
                                                                                    PMC m   predictors for  multiple event streams  to synchronize their
    3.2.2 The proposed approach. We designed and developed a                              P
                                                                                    models (i.e., the transition probability matrix of each predictor)
scalable and distributed pattern prediction system over a massive
                                                                                    within the system in a communication-efficient manner.
input event streams of moving objects. As the base prediction
                                                                                        We propose a synchronization operation for the parameters of
model, we use the PMC forecasting method [2]. Moreover, we pro-
                                                                                    the models (fi = Πi : i ∈ [k]) of the k distributed PMC predictors.
pose to enable the information exchange between the distributed
                                                                                    The operation is based on distributing the maximum-likelihood
predictors/learners of the input event streams, by adapting the
                                                                                    estimation [4] for the transition probabilities of the underlying
distributed online prediction protocol of [16] to synchronize the
                                                                                    PMC m   models described by:
prediction models, i.e., the transitions probabilities matrix of the                      P
                                                                                                                    Í
PMC predictors.                                                                                                       k ∈K nk,i, j
    Algorithm 1 presents the distributed online prediction protocol                                      π̂i, j = Í     Í
                                                                                                                   k ∈K l ∈L nk,i,l
by dynamic model synchronization on both the predictor nodes
                                                                                       Moreover, we measure the divergence of local models from
and the coordinator. We refer to the PMC’s transition matrix Πi
                                                                                    the reference model ∥ fk − fr ∥ 2 by calculating the sum of square
on predictor node ni by fi . That is, when a predictor ni : i ∈ [k]
                                                                                    difference between the transition probabilities Πi and Πr :
observes an event e j it revises its internal model state (i.e., fi ) and                                          Õ
provides a prediction report. Then it checks the local conditions                                  ∥ fk − fr ∥ 2 =     (π̂k i, j − π̂r i, j)2
(batch size b and local model divergence from a reference model                                                     i, j
fr ) to decide whether there is a need to synchronize its local                        In general, our approach relies on enabling the collaborative
model with the coordinator [or not]. fr is maintained in the                        learning among the distributed predictors. Each predictor node
predictor node as a copy of the last computed aggregated model                      receives a stream of events related to a distinct moving object,
fˆ from the previous full synchronization step, which is shared                     and the central coordinator is responsible for synchronizing their
between all local predictors/learners. By monitoring the local                      prediction models using the synchronization operation. Moreover,
condition ∥ fi − fr ∥ 2 > ∆ on all local predictors, we have a                      the predictors they only need to share the parameters of their
guarantee that if none of the local conditions is violated, the                     models, not the consumed event streams.
                                                      Í
divergence (i.e., variance of local models δ (f ) = k1 kj=1 ∥ fi − fˆ ∥ 2 )            We assume that the underlying event streams belong to the
does not exceed the threshold ∆ [16].                                               same distribution and share the same behavior (e.g., mobility
    On the other hand, the coordinator receives the prediction                      patterns). We claim this assumption is reasonable in many appli-
models from the predictor nodes that requested for model syn-                       cation domains: for instance, in the context of maritime surveil-
chronization (violation). Then it tries to keep incrementally query-                lance, vessels travel through standard routes, defined by the
ing other nodes for their local prediction models until reaching                    International Maritime Organization (IMO). Additionally, vessels
out all nodes, or the variance of the aggregated model fˆ that is                   have similar mobility patterns in specific areas such as moving




                                                                              112
with low speed and multiple turns near the ports [20, 29]. That al-                     4     IMPLEMENTATION DETAILS
lows our system to construct a coherent global prediction model                         In this section, we briefly describe in detail the implementation of
dynamically for all input event streams based on merging their                          our system on top of Apache Flink and Apache Kafka frameworks.
local prediction models.                                                                Each of the three sub-modules, described in Section 3.3, have
                                                                                        been implemented as Flink operations over the Kafka events
3.3 Distributed architecture                                                            stream.
Our system consumes as an input2 an aggregated stream of events                            Pre-processing and Prediction Operators. Listing 1 shows
coming from a large number of moving objects, which is con-                             how the main workflow of the system is implemented as Flink
tinuously collected and fed into the system. It allows users to                         data flow program.
register a pattern P to be monitored over each event stream of                             The system ingests the input events stream from a Kafka
a moving object. The output stream consists of original input                           cluster that is mapped to a DataStream of events, which is then
events and predictions of full matches of P, displayed to the end                       processed by an EventTuplesMapper to create tuples of (id, event),
users. Figure 3 presents the overview of our system architecture                        where the id is associated to the identifier of the moving object.
and its main components.                                                                To handle events coming in out of order in a certain margin, the
                                                                                        stream of event tuples is processed by a TimestampAssigner, it
                                                                                        assigns the timestamps for the input events based on the extracted
                                                                                        creation time. Afterwards, an ordered stream of event tuples is
                                                                                        generated using a process function EventSorter.

                                                                                            DataStream eventsStream =
                                                                                                 env.addSource(kafkaConsumer);
                                                                                            // Create tuples (id,event) and assign time stamps
                                                                                            DataStream> eventTuplesStream =
                                                                                            inputEventsStream.map(new EventTuplesMapper())
                                                                                            .assignTimestampsAndWatermarks(new
                                                                                                 EventTimeAssigner());
                                                                                            // Create the ordered keyed stream
                                                                                             orderedEventsStream =
                                                                                                  eventsStream.keyBy(0).process(new
                                                                                                  EventSorter()).keyBy(0);
                   Figure 3: System Architecture.                                           // Consume the events by the predictors
                                                                                            LocalPredictorNode predictorNode =new
                                                                                                 LocalPredictorNode(P);
                                                                                            DataStream processedEventsStream =
   The system is composed of three processing units: (i) pre-pro-
                                                                                            orderedEventsStream.map(predictorNode);
cessing operators that receive the input event stream and per-
form filtering and ordering operations, before partitioning the                             Listing 1: Flink pipeline for local predictors workflow
input event stream to multiple event streams based on the asso-
ciated moving object (ii) predictor nodes (learners), which are                            The ordered stream is then transformed to a keyedEventsStream
responsible for maintaining a prediction model for the input                            by partitioning it, based on the ids values, using a keyBy oper-
event streams. Each prediction node is configured to handle an                          ation. A local predictor node in a distributed environment is
event stream from the same moving object, in order to provide                           represented by a map function over the keyedEventsStream. Each
online predictions for a predefined pattern P (iii) a coordinator                       parallel instance of the map operator (predictor) always pro-
node that communicates through Kafka stream channels with                               cesses all events of the same moving object (i.e., equivalent id),
the predictors to realize the distributed online learning protocol.                     and maintains a bounded prediction model (i.e., PMC m     predictor)
                                                                                                                                               P
It builds a global prediction model, based on the received local                        using the Flink’s Keyed State 3 . The output streams of the moving
models, and then shares it among the predictors.                                        objects from the parallel instances of the predictor map functions
   Our distributed system consists of multiple pre-processing op-                       are sent to a new Kafka stream (i.e., same topic name). They then
erators, prediction nodes, and a central coordinator node. These                        can be processed by other components like visualization or users
units run concurrently and are arranged as a data processing                            notifier.
pipeline, depicted in Figure 3. We leverage Apache Kafka as a                              Moreover, the implementation of the predictor map function in-
messaging platform to ingest the input event streams and to                             cludes the communication with coordinator using Kafka streams.
publish the resulting streams. Also, it is used as the communica-                       At the beginning of the execution, it sends a registration request
tion channel between the predictor nodes and the coordinator.                           to the coordinator. Also at the run-time, it sends its local pre-
Apache Flink is employed to execute the system’s distributed                            diction model as synchronization request, or as a response for a
processing units over the input event streams: the pre-processing                       resolution request from the coordinator. These communication
operators, the prediction units, and the coordinator node. Our                          messages are published into different Kafka topics as depicted in
system architecture can be modeled as a logical network of pro-                         Table 1.
cessing nodes, organized in the form of a DAG, inspired by the                             Coordinator. It manages the distributed online learning pro-
Flink runtime dataflow programs [6].                                                    tocol operations, which is also implemented as Flink program.
                                                                                        The coordinator receives messages from the local predictors
2 In practice, the aggregated input events stream is composed of multiple event         3 Keyed State in Flink: https://ci.apache.org/projects/flink/flink-docs-release-
streams (partitions) from a set of moving objects.                                      1.3/dev/stream/state.html#kayed-state




                                                                                  113
        Table 1: Messages to Kafka topics mapping.                             and 32GB RAM. We used Apache Flink v1.3.2 and Apache Kafka
                                                                               v0.10.2.1 for our tests.
     Message                     Kafka Topic                                      Evaluation criteria. Our goal is to evaluate our distributed pat-
     RegisterNode,               LocalToCoordinatorTopicId                     tern prediction system, which enables the synchronization of
     RequestSync, and                                                          prediction models (i.e., PMC models) on the distributed predic-
     ResolutionAnswer                                                          tor nodes. Our proposed system can operate in three different
                                                                               modes of protocols/schemes of models synchronization: (i) static
     CoordinatorSync and         CoordinatorToLocalTopicId                     scheme based on synchronizing the prediction models periodi-
     RequestResolution                                                         cally every b of input events in each stream, (ii) continuous, full
                                                                               synchronization for each incoming event (hypothetical), (iii) dy-
                                                                               namic synchronization protocol based on making the predictors
through a Kafka Stream of a topic named "LocalToCoordinator-                   communicate their local prediction models periodically but only
TopicId". It is implemented as a single map function over the                  under condition that the divergence of the local models from a
messages stream, by setting the parallelism level of the Flink pro-            reference model exceeds a variance threshold ∆ (recommended).
gram to "1". Increasing the parallelism will scale up the number                  We compare our proposed system against the isolated predic-
of parallel coordinator instances, for example, in order to handle             tion mode, in which models are computed on single streams only,
different groupings of the input event streams. The map oper-                  and compare the predictive performance in terms of :
ator of the coordinator handles three message types from the                                      # of correct predictions
predictors: (i) RegisterNode that contains a registration request                 (i) Precision = # of total predictions is the fraction of the pro-
for a new predictor node, (ii) RequestSync to receive a local                         duced predictions that are correct. For each new event in
model after violation, (iii) ResolutionAnswer to receive a reso-                      the stream, the predictor provides a prediction interval
lution response from a local predictor node. In addition, it sends                    where the full match of the pattern might occur. Thus, the
CoordinatorSync messages for all predictors after creating a                          predictions are temporarily stored until a full match is
new global prediction model, or RequestResolution to a ask                            detected. At that point, all stored prediction intervals are
the local predictors for their prediction models.                                     evaluated by considering those intervals where the full
                                                                                      match occurred within as correct.
5    EMPIRICAL EVALUATION                                                        (ii) Spread= end(I ) − start(I ) is the width of the prediction
In this section, we evaluate our proposed system by analyzing                         interval I , which represents the number of events between
the predictive performance and communication complexity us-                           the start and the end of I .
ing real-world event streams provided by the datAcron project                       Moreover, we study the communication cost by measuring the
in the context of maritime monitoring. The used event streams                  cumulative communication that captures the number of messages,
describe critical points (i.e., synopses) of moving vessels trajecto-          which are required to perform the distributed online learning
ries, which are derived from raw AIS messages as described in                  modes to synchronize the prediction models. Next, we present
[30]. In particular, for our evaluation experiments we used a data             the experimental results for the patterns P1 = Sailinд with an
set of synopses that contains 4, 684, 444 critical points of 5055              order of m = 2, and P2 =changeInHeading; gapStart; gapEnd;
vessels sailing in the Atlantic Ocean during the period from 1                 changeInHeading with first order m = 1. All experiments are
October 2015 to 31 March 2016.                                                 performed with setting the batch size to 100 (b = 100), the vari-
   We used the synopses data set to generate a simulated stream                ance threshold of 2 (∆ = 2), 80% as PMC prediction threshold
of event tuples i.e., (id, timestamp, longitude, latitude, annotation,         (θ f c = 80%), and 200 for the maximum spread.
speed, heading), which are processed by the system to attach an ex-
                                                                                  Experimental results. Figure 4 depicts the average precision
tra attribute type that represents the event value, where type ∈ Σ,
                                                                               scores of predictions models (one prediction model per vessel)
and Σ = Σ1 ={VerySlow, Slow, Moving, Sailing, Stopping}, which is
                                                                               of all synchronization modes for the first pattern P1 = Sailinд,
based on a discretization of the speed values. That is, Σ1 includes
                                                                               namely, isolated without synchronization, continuous (full-sync),
a simple derived event types based on the speed value that can
                                                                               static, and our recommended approach based on the dynamic
be used over streams of raw AIS or critical points. Or Σ = Σ2 =
                                                                               synchronization scheme. It can be clearly seen that all methods
{stopStart, stopEnd, changeInSpeedStart, changeInSpeedEnd, slow-
                                                                               of distributed learning outperform the isolated prediction models.
MotionStart, slowMotionEnd, gapStart, gapEnd, changeInHeading},
                                                                               The hypothetical method of full continuous synchronization has
which is derived based on the values of the annotation attribute
                                                                               the highest precision rates, while the static and dynamic syn-
that encodes the extracted trajectory movement events [30]. Σ2
                                                                               chronization schemes have close precision scores. Consequently,
represents the set of possible mobility changes in the vessel’s
                                                                               dynamic synchronization is not much weaker than the static
trajectory [30], each critical point has at least one event. Where
                                                                               synchronization, but requires much less communication, as ex-
in the case of multiple values, we generate duplicate points each
                                                                               plained below.
of which corresponding to one event in the same order of Σ2 .
                                                                                  Figure 5 provides the amount of the accumulated communica-
   In our experiments, we monitor a pattern P1 = Sailinд with
                                                                               tion that is required by the three modes of the distributed online
Σ1 that detects when the vessel is underway (sailing). Likewise,
                                                                               learning, while the isolated approach does not require any com-
we test a second pattern P2 =changeInHeading; gapStart; gapEnd;
                                                                               munication between the predictors. These results are shown for
changeInHeading with Σ2 that describes a potential illegal fishing
                                                                               P1 . As expected, a larger amount of communication is required
activity [2].
                                                                               for the continuous synchronization comparing to the static and
   Experimental setup. We ran our experiments on single-node                   dynamic approaches. Also, it can be seen that we can reduce the
standalone Flink cluster deployed on an Ubuntu Server 17.04                    communication overhead by applying the dynamic synchroniza-
with Intel(R) Core(TM) i7-7700 CPU @ 3.60GHz X 8 processors                    tion protocol (a reduction by a factor of 100) comparing to the




                                                                         114
Figure 4: Precision scores with respect to the number of                              Figure 6: Average spread value for P1 .
input events over time for P1 .


static synchronization scheme, even with a small variance thresh-
old ∆ = 2. Furthermore, the dynamic protocol is still preserving
a close predictive performance to the static one (see Figure 4).
Therefore, we will only consider the dynamic synchronization
and the isolated approach in the evaluation of the second pattern.




                                                                           Figure 7: Precision scores of P2 for PLEASURE CRAFT ves-
                                                                           sels.


                                                                               In Figure 4, we also noted that the precision is going down in a
                                                                           first phase and stabilizes then. This seems to be counter-intuitive,
                                                                           as the models should improve when getting more data up to a
                                                                           certain point. For explanation, we have investigated the effect
                                                                           of the distributed synchronization of the prediction models on
                                                                           the average spread value, Figure 6 shows the spread results for
Figure 5: Cumulative communication with respect to the
                                                                           all approaches. It can be seen that the spread is higher for the
number of input events over time for P1 .
                                                                           distributed learning based methods comparing to the isolated
                                                                           approach. Furthermore, the average spread is decreasing over
                                                                           time until convergence, as result of confidence increase in the
                                                                           models. This may explain the drop in the precision scores from




                                                                     115
the beginning until reaching the convergence. We will investigate                            [7] Gianpaolo Cugola and Alessandro Margara. 2012. Processing Flows of Infor-
further in the interrelation between precision and spread in future                              mation: From Data Stream to Complex Event Processing. ACM Comput. Surv.
                                                                                                 44, 3, Article 15 (June 2012), 62 pages. https://doi.org/10.1145/2187671.2187677
work.                                                                                        [8] Ofer Dekel, Ran Gilad-Bachrach, Ohad Shamir, and Lin Xiao. 2012. Optimal
    For the second, more complex pattern (P2 ), we have found that                               distributed online prediction using mini-batches. Journal of Machine Learning
                                                                                                 Research 13, Jan (2012), 165–202.
the precision was worse for a distributed model generated over all                           [9] Lina Fahed, Armelle Brun, and Anne Boyer. 2014. Efficient Discovery of
vessels than in the model created for each vessel in isolation. This                             Episode Rules with a Minimal Antecedent and a Distant Consequent. In
indicates that there is no global model describing the behavior of                               Knowledge Discovery, Knowledge Engineering and Knowledge Management.
                                                                                                 Springer.
all models consistently. However, when looking at specific groups                           [10] Ioannis Flouris, Vasiliki Manikaki, Nikos Giatrakos, Antonios Deligiannakis,
of vessels, we achieved an improvement in terms of precision.                                    Minos Garofalakis, Michael Mock, Sebastian Bothe, Inna Skarbovsky, Fabiana
As initial experiment, we only enable the synchronization of the                                 Fournier, Marko Stajcer, et al. 2016. FERARI: A Prototype for Complex Event
                                                                                                 Processing over Streaming Multi-cloud Platforms. In Proceedings of the 2016
prediction models associated with vessels that belong to the same                                International Conference on Management of Data. ACM, 2093–2096.
vessel class. Currently, this change is technically performed by                            [11] The Apache Software Foundation. 2012. Apache Kafka. https://kafka.apache.
                                                                                                 org/. (2012).
an extra filter step that passes only one type of vessels, while                            [12] The Apache Software Foundation. 2013. Apache Spark Streaming. http:
multiple runs of the system are required for all vessel types. For                               //spark.apache.org/streaming/. (2013).
example, Figure 7 shows the precision scores for vessels of class                           [13] The Apache Software Foundation. 2014. Apache Flink. https://flink.apache.
                                                                                                 org/. (2014).
PLEASURE CRAFT. An interesting observation is that the dynamic                              [14] The Apache Software Foundation. 2014. Apache Storm. http://storm.apache.
synchronization approach still has higher precision scores than                                  org/. (2014).
the isolated approach. This case might seem to contradict of                                [15] John E Hopcroft, Rajeev Motwani, and Jeffrey D Ullman. 2006. Automata
                                                                                                 theory, languages, and computation. International Edition 24 (2006).
our assumption that the input event streams belong to the same                              [16] Michael Kamp, Mario Boley, Daniel Keren, Assaf Schuster, and Izchak Sharf-
distribution and share the same behavior, but it actually follows                                man. 2014. Communication-efficient distributed online prediction by dynamic
                                                                                                 model synchronization. In Joint European Conference on Machine Learning
the same assumption but between the predictors of vessels within                                 and Knowledge Discovery in Databases. Springer, 623–639.
the same type group. We will further investigate the effect of                              [17] Michael Kamp, Sebastian Bothe, Mario Boley, and Michael Mock. 2016.
groupings and more patterns in future work.                                                      Communication-Efficient Distributed Online Learning with Kernels. In Joint
                                                                                                 European Conference on Machine Learning and Knowledge Discovery in
                                                                                                 Databases. Springer, 805–819.
6    CONCLUSION                                                                             [18] John Langford, Alex J Smola, and Martin Zinkevich. 2009. Slow learners are
                                                                                                 fast. Advances in Neural Information Processing Systems 22 (2009), 2331–2339.
In this paper, we have presented a system that provides a dis-                              [19] Srivatsan Laxman, Vikram Tankasali, and Ryen W. White. 2008. Stream
tributed pattern prediction over multiple large-scale event streams                              Prediction Using a Generative Model Based on Frequent Episodes in Event
                                                                                                 Sequences. In ACM SIGKDD.
of moving objects (vessels). The system uses the event forecast-                            [20] Bo Liu, Erico N de Souza, Stan Matwin, and Marcin Sydow. 2014. Knowledge-
ing with Pattern Markov Chain (PMC) [2] as the base prediction                                   based clustering of ship trajectories using density-based approach. In Big Data
                                                                                                 (Big Data), 2014 IEEE International Conference on. IEEE, 603–608.
model on each event stream, and it applies the protocol for dis-                            [21] David Luckham. 2008. The power of events: An introduction to complex event
tributed online prediction [16] to exchange information between                                  processing in distributed enterprise systems. In International Workshop on
the prediction models over multiple input event streams. Our                                     Rules and Rule Markup Languages for the Semantic Web. Springer, 3–3.
                                                                                            [22] Heikki Mannila, Hannu Toivonen, and A. Inkeri Verkamo. 1997. Discovery of
proposed system has been implemented using Apache Flink and                                      Frequent Episodes in Event Sequences. Data Mining and Knowledge Discovery
Apache Kafka. In order to show the usefulness and effectiveness                                  (1997).
of our approach, we empirically tested it against large real-world                          [23] Michael Mathioudakis and Nick Koudas. 2010. Twittermonitor: trend detection
                                                                                                 over the twitter stream. In Proceedings of the 2010 ACM SIGMOD International
event streams related to trajectories of moving vessels.                                         Conference on Management of data. ACM, 1155–1158.
   As future work, we will address the open issues emerging                                 [24] Daniele Miorandi, Sabrina Sicari, Francesco De Pellegrini, and Imrich Chlam-
                                                                                                 tac. 2012. Internet of things: Vision, applications and research challenges. Ad
from the current findings. Firstly, we will study the interrelation                              Hoc Networks 10, 7 (2012), 1497–1516.
between precision and spread scores by validating the approach                              [25] Douglas C Montgomery, Cheryl L Jennings, and Murat Kulahci. 2015. Intro-
over synthetic event streams. Secondly, we will investigate the                                  duction to time series analysis and forecasting. Wiley.
                                                                                            [26] Vinod Muthusamy, Haifeng Liu, and Hans-Arno Jacobsen. 2010. Predictive
effect of grouping the input event streams on the predictive                                     Publish/Subscribe Matching. In DEBS. ACM.
performance of our method.                                                                  [27] Grégory Nuel. 2008. Pattern Markov Chains: Optimal Markov Chain Embed-
                                                                                                 ding through Deterministic Finite Automata. Journal of Applied Probability
                                                                                                 (2008).
ACKNOWLEDGMENTS                                                                             [28] International Maritime Organization. 2001. Automatic identification systems.
                                                                                                 http://www.imo.org/OurWork/Safety/Navigation/Pages/AIS.aspx. (2001).
This work was supported by EU Horizon 2020 datAcron project                                 [29] Giuliana Pallotta, Michele Vespe, and Karna Bryan. 2013. Vessel pattern
(grant agreement No 687591).                                                                     knowledge discovery from AIS data: A framework for anomaly detection and
                                                                                                 route prediction. Entropy 15, 6 (2013), 2218–2245.
                                                                                            [30] Kostas Patroumpas, Elias Alevizos, Alexander Artikis, Marios Vodas, Nikos
REFERENCES                                                                                       Pelekis, and Yannis Theodoridis. 2017. Online event recognition from moving
                                                                                                 vessel trajectories. GeoInformatica 21, 2 (2017), 389–427.
 [1] Rakesh Agrawal, Tomasz Imieliński, and Arun Swami. 1993. Mining Associa-               [31] Kostas Patroumpas, Alexander Artikis, Nikos Katzouris, Marios Vodas, Yan-
     tion Rules Between Sets of Items in Large Databases. In ACM SIGMOD.                         nis Theodoridis, and Nikos Pelekis. 2015. Event Recognition for Maritime
 [2] Elias Alevizos, Alexander Artikis, and George Paliouras. 2017. Event Forecast-              Surveillance.. In EDBT. 629–640.
     ing with Pattern Markov Chains. In Proceedings of the 11th ACM International           [32] R. Vilalta and Sheng Ma. 2002. Predicting rare events in temporal domains. In
     Conference on Distributed and Event-based Systems. ACM, 146–157.                            ICDM.
 [3] Elias Alevizos, Anastasios Skarlatidis, Alexander Artikis, and Georgios                [33] Lin Xiao. 2010. Dual averaging methods for regularized stochastic learning
     Paliouras. 2015. Complex event recognition under uncertainty: A short sur-                  and online optimization. Journal of Machine Learning Research 11, Oct (2010),
     vey. Event Processing, Forecasting and Decision-Making in the Big Data Era                  2543–2596.
     (EPForDM) (2015), 97–103.                                                              [34] Feng Yan, Shreyas Sundaram, SVN Vishwanathan, and Yuan Qi. 2013. Dis-
 [4] Theodore W Anderson and Leo A Goodman. 1957. Statistical inference about                    tributed autonomous online learning: Regrets and intrinsic privacy-preserving
     Markov chains. The Annals of Mathematical Statistics (1957), 89–110.                        properties. IEEE Transactions on Knowledge and Data Engineering 25, 11 (2013),
 [5] Amazon Web Services (AWS). 2013. Amazon Kinesis. https://aws.amazon.                        2483–2493.
     com/de/kinesis/. (2013).                                                               [35] Cheng Zhou, Boris Cule, and Bart Goethals. 2015. A pattern based predictor
 [6] Paris Carbone, Asterios Katsifodimos, Stephan Ewen, Volker Markl, Seif Haridi,              for event streams. Expert Systems with Applications (2015).
     and Kostas Tzoumas. 2015. Apache flink: Stream and batch processing in a
     single engine. Bulletin of the IEEE Computer Society Technical Committee on
     Data Engineering 36, 4 (2015).




                                                                                      116