<!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>PeNeLoop: Parallelizing Federated SPARQL Queries in Presence of Replicated Fragments</article-title>
      </title-group>
      <contrib-group>
        <contrib contrib-type="author">
          <string-name>Thomas Minier</string-name>
          <email>fthomas.minier1@etu</email>
        </contrib>
        <contrib contrib-type="author">
          <string-name>Gabriela Montoya</string-name>
          <email>gmontoya@cs.aau.dk</email>
          <xref ref-type="aff" rid="aff1">1</xref>
        </contrib>
        <contrib contrib-type="author">
          <string-name>Hala Skaf-Molli</string-name>
        </contrib>
        <contrib contrib-type="author">
          <string-name>Pascal Molli</string-name>
          <email>pascal.molli@guniv-nantes.fr</email>
        </contrib>
        <aff id="aff0">
          <label>0</label>
          <institution>Aalborg University</institution>
          ,
          <country country="DK">Denmark</country>
        </aff>
        <aff id="aff1">
          <label>1</label>
          <institution>Department of Computer Science</institution>
        </aff>
      </contrib-group>
      <fpage>96</fpage>
      <lpage>109</lpage>
      <abstract>
        <p>Replicating data fragments in Linked Data improves data availability and performances of federated query engines. Existing replication aware federated query engines mainly focus on source selection and query decomposition in order to prune redundant sources and reduce intermediate results thanks to data locality. In this paper, we extend replication-aware federated query engines with a replication-aware parallel join operator: PeNeLoop. PeNeLoop exploits redundant sources to parallelize the join operator and reduce execution time. We implemented PeNeLoop in the federated query engine FedX with the replicatedaware source selection Fedra and we empirically evaluated the performance of FedX Fedra PeNeLoop. Experimental results suggest that FedX Fedra PeNeLoop outperforms FedX Fedra in terms of execution time while preserving answer completeness.</p>
      </abstract>
    </article-meta>
  </front>
  <body>
    <sec id="sec-1">
      <title>-</title>
      <p>
        Following the Linked Data principles, billions of RDF triples are made available
through SPARQL endpoints. Even if federated SPARQL query engines [
        <xref ref-type="bibr" rid="ref1">8,15,1</xref>
        ]
allow to execute SPARQL queries over multiple SPARQL endpoints, data-availability
and reliability of SPARQL endpoints is still an issue [
        <xref ref-type="bibr" rid="ref5">5</xref>
        ].
      </p>
      <p>Data replication is a common practice to overcome availability issues in
distributed databases [13]. However, data replication in Linked Data is more
challenging: the autonomy of data providers hosting SPARQL endpoints, and data
consumers running federated query engines, prevent data replication to be
designed. The fragmentation schema and the replication schema remain unknown
until a data consumer de nes a federation of SPARQL endpoints in a federated
query engine.</p>
      <p>Existing replication-aware [11,12] and duplicate-aware [14] federated query
engines focus on source selection and query decomposition in order to prune
redundant sources and use data-locality to reduce intermediate results. We point
out that replicated data can also be used to parallelize query processing, and
consequently reduce execution time.</p>
      <p>In this paper, we extend replication-aware federated query engines with
PeNeLoop, a replication-aware parallel join operator. More precisely, PeNeLoop
solves the parallel join problem with fragment replication (PJP-FR). Given a
SPARQL query and a set of data sources with replicated fragments, the
problem is to use all data sources to reduce query execution time while preserving
answer completeness and reducing data redundancy.</p>
      <p>
        In contrast to inter-operator parallelism proposed in the state-of-the-art
federated query engines [
        <xref ref-type="bibr" rid="ref1">1,15</xref>
        ], PeNeLoop introduces parallelization at the
operator level in order to preserve properties ensured by replicated-aware source
selection strategies [11] and replication-aware query decompositions [12].
      </p>
      <p>PeNeLoop is based on Bound Join operator implemented in FedX [15].
Bound joins were originally designed to reduce the number of requests sent
in a nested loop join [13]. PeNeLoop extends bound joins processing to use
all relevant endpoints with replicated fragments and distribute join processing
among them. The contributions of this work are as follows:
(i) We present PeNeLoop, a novel replication-aware parallel join operator
that uses replicated fragments to reduce query execution time. PeNeLoop
is the st attempt to use replicated fragments to parallelize query
processing in Linked Data
(ii) We extend the federated query engine FedX [15] and the source selection
strategy Fedra [11] with PeNeLoop.
(iii) We experiment FedX, FedX Fedra and FedX Fedra PeNeLoop in
di erent setups. We show that FedX Fedra PeNeLoop outperforms
FedX and FedX Fedra in terms of execution time while preserving
properties of Fedra in terms of reduced number of transferred tuples and
answer completeness. The improvements are signi cative for queries with
a large number of intermediate results.</p>
      <p>The paper is organized as follows: Section 2 provides background and
motivations. Section 3 presents the PeNeLoop approach and algorithm. Section 4
presents our experimental setup and describes our results. Section 5 summarizes
related works. Finally, conclusions and future works are outlined in Section 6.
2</p>
    </sec>
    <sec id="sec-2">
      <title>Background and Motivations</title>
      <p>For replicating data, we follow the approach of replicated fragments introduced
in [11,12]. Data consumers replicate fragments composed of RDF triples that
satisfy a given triple pattern. Figure 1a shows a fragment from DBpedia which
contains RDF triples that match the triple pattern ?film dbo:director ?director.
Fragments are described using a 2-tuple fd that indicates the authoritative
source of the fragment, e.g. DBpedia, and the triple pattern met by the
fragment's triples.</p>
      <p>Figure 1b shows a federation with four SPARQL endpoints: E0, E1; E2 and
E3. These endpoints expose replicated fragments from DBpedia and
LinkedMDB. Figure 1c describes a federated SPARQL query Q1 executed against this
(a) Fragment description</p>
      <p>(b) Replicated fragments
triples(f): { dbr:A Knight’s Tale
dbo:director dbr:Brian Helgeland,
dbr:A Thousand Clowns
dbo:director dbr:Fred Coe,
dbr:Alfie (1966 film)
dbo:director dbr:Lewis Gilbert,
dbr:A Moody Christmas
dbo:director dbr:Trent O’Donnell,
dbr:A Movie dbo:director
dbr:Bruce Conner, · · · }
fd(f): &lt;dbpedia, ?film dbo:director ?director&gt;</p>
      <p>DBpedia
f1
E0
f2
f2</p>
      <p>E1 E2 E3
fd(f1): &lt;dbpedia, ?director dbo:nationality ?nat&gt;
fd(f2): &lt;dbpedia, ?film dbo:director ?director&gt;
fd(f3): &lt;linkedmdb, ?movie owl:sameAs ?film&gt;
fd(f4): &lt;linkedmdb, ?movie linkedmdb:genre ?genre&gt;
fd(f5): &lt;linkedmdb, ?genre linkedmdb:film genre name ?name&gt;
f4
f3,f5</p>
      <p>LinkedMDB
f4, f5
1
(c) Federated SPARQL query Q1 and its relevant fragments an d1 endpoints
s e l e c t d i s t i n c t
where f
? d i r e c t o r dbo : n a t i o n a l i t y ? nat . ( tp1 )
? f i l m db : d i r e c t o r ? d i r e c t o r . ( tp2 )
? movie owl : sameAs ? f i l m . ( tp3 )
? movie linkedmdb : g e n r e ? g e n r e . ( tp4 )
? g e n r e linkedmdb : f i l m g e n r e n a m e ?gname . ( tp5 )
g
Triple
pattern
tp1
tp2
tp3
tp4
tp5</p>
      <p>Relevant Relevant
fragment endpoint
f1 E0
f2 E1, E2
f3 E2
f4 E1, E3
f5 E2, E3
federation and its relevant fragments. For instance, the triple pattern tp4 has
relevant fragment f4 that has been replicated at E1 and E3.</p>
      <p>The logical plan of Q1 produced by FedX [15] is presented in Figure 2a. As
FedX is not replication-aware, i.e., it does not know that the evaluation of tp2
at E1 or E2 will produce the same results, query execution following this plan
will retrieve redundant data from endpoints and increase signi cantly the query
execution time.</p>
      <p>The Fedra [11] replication-aware source selection prunes redundant sources
in order to minimize intermediate results. Fedra selects E2 for tp2; tp3 and tp5,
E1 for tp4 and E0 for tp1. Next, Fedra lets FedX builds the logical plan of
Figure 2b that minimizes intermediate results.</p>
      <p>As pointed in Figure 2b, Fedra has removed E3 from selected sources of tp4.
However, it also removes an opportunity of parallelization. Indeed, it is possible
to use both endpoints to perform in parallel half of the join of '2 with E1 and
the other half with E3, as they mirror each other 3.</p>
      <p>Such parallelization can be obtained with a replication-aware query
decomposer or with intra-operator [13] parallelism. In this paper, we focus on the
second approach because it can be easily embedded in current federated query
3 Note that joins '1 and '3 cannot be parallelized in this way, because '1 is a local
join performed at E2, and tp1 has only one relevant source.
(a) FedX Left-Linear plan for Q1
(b) FedX</p>
      <p>Fedra Left-Linear plan for Q1
Π
./4
./3
./2
./1
tp2
./1
tp3
Π
./3
tp5
./2
Fig. 2: Logical plans generated by FedX and FedX
1</p>
      <p>Fedra for Q1
1
engines. Consequently, the challenge is to build replication-aware parallel
operators to speed-up query execution.</p>
      <sec id="sec-2-1">
        <title>Parallel Join Problem with Fragment Replication (PJP-FR)</title>
        <p>Given S1 and S2 two disjoint sets of replicated data sources. A set of
replicated data sources is a set of endpoints that replicate the same fragments. Given
a join 'i between O1 and O2 with relevant sources respectively, S1 and S2. The
parallel join problem with fragment replication is to distribute the execution of
join 'i among endpoints of S1 and S2 in order to minimize the execution time
while guaranteeing complete query answers.
3</p>
        <p>PeNeLoop : A Replication-Aware Nested Loop Join
Operator
PeNeLoop is a solution for parallel join problem with fragment replication with
the following assumptions: (i) we focus on nested loop join (NLJ), (ii) we do not
consider the load of di erent endpoints, (iii) we consider that replicated
fragments are synchronized, (iv) replicated sources are determined by a
replicationaware source selection algorithm as Fedra before pruning.
3.1</p>
      </sec>
      <sec id="sec-2-2">
        <title>NLJ Processing</title>
        <p>During a NLJ processing, the query engine iteratively evaluates each triple
pattern, starting with a single pattern and substituting the set of mappings produced
by the pattern's execution in the next evaluation step. Even if a NLJ is more
e cient when the rst evaluated triple patten is more selective than the others,
it still produces many remote requests in a distributed setting. In FedX [15],
the Bound Join (BJ) operator is proposed to minimize the number of join steps
and the number of requests sent in nested loop joins. A BJ consists of a nested
loop join where sets of mappings are grouped in blocks, i.e., as a single subquery
using SPARQL UNION constructs. The subquery is then sent to the relevant
endpoint in a single remote request. This technique acts as a distributed semijoin
and allows to reduce the number of requests by a factor equivalent to the size of
the block.</p>
        <p>PeNeLoop proposes to parallelize the BJ operator itself. Instead of sending
all blocks to the same endpoint, PeNeLoop uses the knowledge about replicated
sources to further parallelize the bound join operator. When processing a join in a
basic graph pattern (BGP), if the current triple pattern has N relevant sources
that replicate the same fragment, PeNeLoop sends each block to a di erent
endpoint in a Round Robin fashion, i.e., the block bi is sent to the endpoint Ek,
k i mod N . Therefore, PeNeLoop does not increase the number of remote
calls while increasing the parallelization during join processing.
3.2</p>
      </sec>
      <sec id="sec-2-3">
        <title>PeNeLoop Algorithm</title>
      </sec>
      <sec id="sec-2-4">
        <title>Algorithm 1: PeNeLoop</title>
        <p>Input: tp  s; p; o¡: a triple pattern, E tE0; : : : ; Em 1u: relevant endpoints
of tp, N extOp: next operator in the pipeline, b: maximum number of
mappings per block
Data: Mi: a set of mappings produced by the previous operator in the pipeline,</p>
        <p>B tM1; : : : ; Mnu: block of sets of mappings waiting to be sent</p>
        <p>Init: B tu, k 0
1 SendBlock(block, tp):
2 Q GroupedSubquery(block, tp)
3 SendQuery(Q) to Ek
4 B = tu
5 k = pk 1q mod Size(E)
6 onMappings(Mi):
7 B = B Y tMiu
8 if Size(B) ¥ b then
9 SendBlock(B, tp)
10 end</p>
        <p>PeNeLoop is de ned as part of a pipelining approach allowing for
intermediate results to be processed by the next operator as soon as they are ready,
providing higher throughput than a blocking model.</p>
        <p>Algorithm 1 describes the PeNeLoop algorithm using an event driven paradigm.
Sets of mappings Mi are produced by the previous operator in the pipeline and
sent in continuous to PeNeLoop operator. When a set Mi arrives (Line 6), it
is stored in the next block B. When B reaches its maximum size b (Line 8),
PeNeLoop generates a subquery in a Bound Join fashion using B and tp
Start</p>
        <p>local join ./ 1
tp2.tp3.tp5</p>
        <p>M6
./i PeNeLoop Join
./i Parallel Bound Join
./ 2
M5
B</p>
        <p>k =⇒
{M1, M2}
{M3, M4}
./ 3
Π</p>
        <p>End
(Line 2). Then, the subquery is sent to the endpoint Ek (Line 3), B is cleared
and the next endpoint is selected using our Round Robin approach (Line 5).</p>
        <p>When results, i.e., new sets of mappings, arrive from the requested endpoints
(Line 11), they are sent to the next operator in the pipeline. Finally, when
the previous operator has completed its work and will not produce any more
data (Line 13), PeNeLoop sends the last non-empty block and then close the
operator.</p>
        <p>In the following, we illustrate PeNeLoop processing for the query Q1
(Figure 1c) using the query plan generated by FedX Fedra (Figure 2b). For
simplicity, we x b 2.</p>
        <p>Figure 3 illustrates a snapshot of the pipeline during the evaluation of the
triple pattern tp4 of the query Q1. We focus on processing of join '2, performed
using PeNeLoop. Two blocks tM1; M2u and tM3; M4u have been already sent
to E1 and E3, respectively. A set of mappings M5 arrived from the join '1 and
was placed in the next block. When another set of mappings M6 arrives, the
block will be full and sent to the next endpoint E1. Join '2 ends when no more
mappings are produced by join '1.
4</p>
      </sec>
    </sec>
    <sec id="sec-3">
      <title>Experimental Study</title>
      <p>The goal of the experimental study is to evaluate the execution time reduction
obtained with the parallelization enabled by PeNeLoop. Moreover, such
reduction is obtained without degrading the reduced number of transferred tuples
and the answer completeness granted by Fedra. We compare the performance
of the federated query engine FedX alone, FedX with the addition of
Fedra (FedX Fedra) and FedX with both Fedra and PeNeLoop (FedX
Fedra PeNeLoop).</p>
      <p>We expect to see that FedX Fedra PeNeLoop exhibits lower query
execution time than FedX and FedX Fedra, while maintaining the same
number of transferred tuples and answer completeness.</p>
      <p>
        Dataset and Queries: We use one instance of the Waterloo SPARQL
Diversity Test Suite (WatDiv) synthetic dataset [
        <xref ref-type="bibr" rid="ref2 ref3">2,3</xref>
        ] with 105 triples. We
generate 50,000 queries from 500 templates. Next, we unbound subjects and objects
of each query. 100 queries with at least one join are then randomly picked to
be executed against our federations. Generated queries are STAR, PATH and
SNOWFLAKE shaped queries, we use the DISTINCT modi er.
      </p>
      <p>Federations: We setup three federations with respectively 10, 20 and 30
SPARQL endpoints, and generate three versions of each of these federations
by randomizing the fragmentation schema. Every schema is distinct from the
others. Fragments are created from the 100 random queries and are replicated
exactly three times to provide opportunities of parallelization.</p>
      <p>To measure the number of transferred tuples, the federated query engine
accesses SPARQL endpoints through a proxy. All the federation endpoints are
deployed on the same machine, and to simulate the network latency, the proxies
were con gured to add a delay of 30ms to each request.</p>
      <p>Hardware con guration: One machine with Intel Xeon E5-2680 v2 2.80GHz
and 128GB of RAM hosts the SPARQL endpoints and performs the queries. Each
SPARQL endpoint is deployed using Jena Fuseki 1.1.14. Fuseki is con gured to
handle incoming queries on only one executing thread to increase the stress load
and study the e ect of the parallelization done by the engine. Endpoints have
no limitations in term of memory used.</p>
      <p>Implementations: FedX Fedra implementation5 (in Java) has been
modi ed to preserve the multiple sources that provide the same relevant
fragments. Additionally, FedX join processing has been modi ed to remove some
redundant synchronization barriers imposed by FedX on the rst join of a plan,
i.e., the right operand can start execution before the left one has nished its
evaluation, and to use PeNeLoop operator when possible6. Every con guration
of this experimental study has received the same modi cations. Proxies used to
measure results are implemented in Java 1.7, using the Apache HttpComponents
Client library 4.3.57.</p>
      <p>Evaluation Metrics: i) Execution Time (ET): is the elapsed time since the
query is posed until the complete answer is produced. We used a timeout of
1800 seconds. ii) Number of parallelized queries (NPQ): is the number of queries
where at least one join has been parallelized by PeNeLoop. This metric is
only used in FedX Fedra PeNeLoop. Queries marked as improved have
a lower execution time (ET ) with FedX Fedra PeNeLoop than with
FedX Fedra. iii) Number of Transferred Tuples (NTT): is the number of
transferred tuples from all the endpoints to the query engine during a query
evaluation. iv) Completeness (C): is the ratio between the answers produced by
the query execution engine and the answers produced by the evaluation of the
4 http://jena.apache.org/, January 2015.
5 https://github.com/gmontoya/fedra, June 2016.
6 Implementation available at: https://github.com/Callidon/peneloop-fedx
7 https://hc.apache.org/, October 2014.
query over the set of all triples available in the federation; values range between
0.0 and 1.0.</p>
      <p>Results presented for ET, NTT and C correspond to the average over the
three versions generated for each size of federation. Queries that failed to deliver
an answer due to a query engine internal error are excluded from the nal results.</p>
      <p>Statistical Analysis: The Wilcoxon signed rank test [17] for paired
nonuniform data is used to study the signi cance of the improvements on
performance obtained when the join execution bene ts from replicated fragments.8
4.1</p>
      <sec id="sec-3-1">
        <title>Execution time</title>
        <p>Figure 4 summarizes the execution time (ET ) for the three federations.
Execution time (ET ) with FedX Fedra PeNeLoop is better for all federations
than with FedX and FedX Fedra. As queries have unbounded subjects
and unbounded objects, they generated more intermediate results during joins,
which allow PeNeLoop to distribute more bindings between relevant sources.
Figure 5 presents the execution time for queries with a large number of
intermediate results (at least 1000 tuples). This represents 562 queries out of 865 for all
federations. PeNeLoop is even more e cient for queries with a large number of
intermediate results. This is an important result because generally the number
of the intermediate results impacts negatively the query execution time.
8 The Wilcoxon signed rank test was computed using the R project (http://www.
r-project.org/)</p>
        <p>Both FedX Fedra and FedX Fedra PeNeLoop bene t from the
reduction of transferred tuples granted by Fedra, which reduce the number of
mappings that PeNeLoop can distribute.</p>
        <p>To con rm that PeNeLoop reduces the execution time of FedX Fedra,
a Wilcoxon signed rank test was run for results of Figure 4 with the hypotheses:
H0 : PeNeLoop does not change the engine query execution time.
H1 : PeNeLoop reduces FedX Fedra's query execution time.</p>
        <p>We obtain p-values no greater than 1:639 10 4 for each federation. These
low p-values allow for rejecting the null hypothesis that the execution time of
FedX Fedra and FedX Fedra PeNeLoop are the same. Additionally,
it supports the acceptance of the alternative hypothesis that FedX Fedra
PeNeLoop has a lower execution time.
4.2</p>
      </sec>
      <sec id="sec-3-2">
        <title>Number of Parallelized Queries</title>
        <p>Figure 6 presents the number of parallelized queries (NPQ ) in FedX Fedra
PeNeLoop for the three versions of each federation. PeNeLoop increases query
parallelization during join processing, especially in larger federations where
fragments are more scattered across endpoints. In most cases, queries parallelized by
PeNeLoop are improved, i.e., they exhibit a lower execution time compared to
FedX Fedra. Parallelized queries with unimproved execution time are those
that do not have a large number of intermediate results. Parallelization of such
10v1 10v2 10v3 20v1 20v2 20v3 30v1 30v2 30v3</p>
        <p>Number of endpoints in federation by version
timeout
unparallelized
parallelized + unimproved
parallelized + improved
queries does not improve query performance, as their joins were not originally
costly to evaluate.</p>
        <p>As pointed in Figure 6, the number of parallelized queries is not constant
within di erent versions the same federation, because the replication schema
directly in uences query parallelization. When this schema is not designed, as
in Linked Open Data, PeNeLoop creates parallelization where locality cannot
be used by Fedra to optimize the query execution plan.
4.3</p>
      </sec>
      <sec id="sec-3-3">
        <title>Number of transferred tuples</title>
        <p>Figure 7 summarizes the number of transferred tuples (NTT ) in di erent
federations. FedX Fedra PeNeLoop transfers the same amount of tuples as
FedX Fedra. This demonstrates that PeNeLoop does not deteriorate the
reduction of transferred tuples provided by Fedra. Moreover, modi cations
performed on FedX to remove some synchronisation barriers do not introduce any
di erence between FedX Fedra and FedX Fedra PeNeLoop in terms
of number of transferred tuples and do not impact FedX Fedra performance.
4.4</p>
      </sec>
      <sec id="sec-3-4">
        <title>Completeness</title>
        <p>Figure 8 presents results concerning answer completeness (C ) for the di erent
federations. In all cases, FedX Fedra PeNeLoop is able to produce the
same answers as FedX Fedra for all queries.</p>
        <p>As observed with the number of transferred tuples (NTT ), our modi cation
for FedX does not reduce the completeness of FedX and FedX Fedra, which
support our claim that this modi cation does not impact negatively FedX
Fedra.</p>
        <p>F
Experimental study results con rm that PeNeLoop can further increase the
performance of join processing in presence of replicated fragments. Execution
time in average is lower with FedX Fedra PeNeLoop than with FedX or
FedX Fedra, and the reduced number of transferred tuples granted by
Fedra is maintained. Answer completeness is not degraded. PeNeLoop is able
to parallelize a signi cant number of queries in presence of replicated fragments
and shows to be more e cient on larger federations. Query performance are
signi cantly improved for queries with a large number of intermediate results, and
the time to evaluate joins is reduced by taking advantage of parallel processing.
5</p>
      </sec>
    </sec>
    <sec id="sec-4">
      <title>Related Work</title>
      <p>Fedra [11] is a replication-aware source selection that uses data locality
produced by replicated fragments to enhance federated query engines performances.
Fedra uses Union and BGP reductions to prune data sources and nds as
many sub-queries that can be executed against the same endpoint as
possible, leading to evaluation of local joins and a reduced number of transferred
tuples. PeNeLoop uses replicated fragments di erently. As seen in Section 2,
Fedra prunes redundant endpoints that cannot be used to creates localities,
whereas PeNeLoop uses these endpoints to create more opportunities of
parallelization.</p>
      <p>LILAC [12] is a replication-aware decomposer. Compared to Fedra, LILAC
is able to reduce intermediate results by allocating a triple pattern to several
endpoints. As for Fedra, PeNeLoop can reuse source selection performed by
LILAC to introduce intra-operator parallelism.</p>
      <p>Other existing sources selection techniques reduce the number of selected
sources by a federated SPARQL query engine. BBQ [9] and DAW [14] use
sketches to estimate the overlapping among sources, but they only operate on
duplicated sources and not on replication itself. They do not provide information
about replicated fragments that allow PeNeLoop to e ciently parallelize join
processing.</p>
      <p>
        Parallel join processing in distributed database systems has been the subject
of signi cant investigation. Parallel nested loop algorithms have been
investigated in [
        <xref ref-type="bibr" rid="ref4 ref6">4,6</xref>
        ], but they do not use replication for parallelization. Instead,
replication is mostly used for fault tolerance and to locate data closer to their access
points [10,13], improving query performance by reducing communication time.
PeNeLoop does not use localities created by data redundancy, but
opportunities of parallelization created by this redundancy.
      </p>
      <p>
        Parallel join processing has been also studied in federated query engines.
For instance, [
        <xref ref-type="bibr" rid="ref1">1,15,7</xref>
        ] propose parallel architectures for executing queries
concurrently at di erent data sources. Anapsid [
        <xref ref-type="bibr" rid="ref1">1</xref>
        ] takes advantage of bushy query
execution plans to create inter-operator parallelism. FedX [15] implements bound
joins in a distributed and highly parallelized environment where di erent
subqueries can be executed at the endpoints concurrently. PeNeLoop creates
intraoperator parallelism and proposes a more advanced parallel join processing using
replication. Similar to FedX, subqueries are executed concurrently, but they are
distributed between endpoints, increasing parallelization.
      </p>
      <p>To our knowledge, none of existing federated query engines propose to take
advantage of replicated data for join processing or propose a replication-aware
parallel join operator.
6</p>
    </sec>
    <sec id="sec-5">
      <title>Conclusions and Future Works</title>
      <p>In this paper, we extended a replication-aware federated query engine with a new
replication-aware parallel join operator PeNeLoop. PeNeLoop provides
intraoperator parallelism relying on replicated data. In this way, PeNeLoop
preserves properties of source-selection and query decomposition replication-aware
federated query engines. We implemented PeNeLoop in FedX. Evaluation
results demonstrates that PeNeLoop improves signi cantly query performance.</p>
      <p>PeNeLoop is the rst attempt to use replicated data to parallelize query
processing in Linked Open Data and opens several perspectives. First, we made
the assumption that the load of the endpoints is uniform during query execution.
We can leverage this hypothesis by making PeNeLoop adaptive to the
performances of endpoints. Second, we focused on a Nested Loop Join operator, we
can also parallelize others operators such as Symmetric Hash-Join [18] used in
Anapsid. Finally, we focused on SPARQL endpoints, and we think that parallel
query processing in presence of replicated fragments can also be applied to the
Triple Pattern Fragment approach [16].</p>
      <p>Acknowledgments. This work is partially supported through the FaBuLA
project, part of the AtlanSTIC 2020 program.
7. Gorlitz, O., Staab, S.: Splendid: Sparql endpoint federation exploiting void
descriptions. In: Proceedings of the Second International Conference on Consuming
Linked Data - Volume 782. pp. 13{24. COLD'11, CEUR-WS.org, Aachen,
Germany, Germany (2010), http://dl.acm.org/citation.cfm?id=2887352.2887354
8. Gorlitz, O., Staab, S.: Federated Data Management and Query Optimization for</p>
      <p>Linked Open Data, vol. 331, pp. 109{137. Springer, Heidelberg (2011)
9. Hose, K., Schenkel, R.: Towards bene t-based rdf source selection for sparql
queries. In: Proceedings of the 4th International Workshop on Semantic Web
Information Management. p. 2. ACM (2012)
10. Kossmann, D.: The state of the art in distributed query processing. ACM
Computing Surveys (CSUR) 32(4), 422{469 (2000)
11. Montoya, G., Skaf-Molli, H., Molli, P., Vidal, M.E.: Federated sparql queries
processing with replicated fragments. In: International Semantic Web Conference. pp.
36{51. Springer International Publishing (2015)
12. Montoya, G., Skaf-Molli, H., Molli, P., Vidal, M.E.: Decomposing federated queries
in presence of replicated fragments. Web Semantics: Science, Services and Agents
on the World Wide Web 42, 1 { 18 (2017), //www.sciencedirect.com/science/
article/pii/S1570826816300580
13. Ozsu, M.T., Valduriez, P.: Principles of distributed database systems. Springer</p>
      <p>Science &amp; Business Media (2011)
14. Saleem, M., Ngomo, A.C.N., Parreira, J.X., Deus, H.F., Hauswirth, M.: Daw:
Duplicate-aware federated query processing over the web of data. In: International
Semantic Web Conference. pp. 574{590. Springer (2013)
15. Schwarte, A., Haase, P., Hose, K., Schenkel, R., Schmidt, M.: Fedx: Optimization
techniques for federated query processing on linked data. In: International Semantic
Web Conference. pp. 601{616. Springer (2011)
16. Verborgh, R., Vander Sande, M., Hartig, O., Van Herwegen, J., De Vocht, L.,
De Meester, B., Haesendonck, G., Colpaert, P.: Triple pattern fragments: A
lowcost knowledge graph interface for the web. Web Semantics: Science, Services and
Agents on the World Wide Web 37, 184{206 (2016)
17. Wilcoxon, F.: Individual comparisons by ranking methods. In: Breakthroughs in</p>
      <p>Statistics, pp. 196{202. Springer (1992)
18. Wilschut, A.N., Apers, P.M.: Data ow query execution in a parallel main-memory
environment. Distributed and Parallel Databases 1(1), 103{128 (1993)</p>
    </sec>
  </body>
  <back>
    <ref-list>
      <ref id="ref1">
        <mixed-citation>
          1.
          <string-name>
            <surname>Acosta</surname>
            ,
            <given-names>M.</given-names>
          </string-name>
          ,
          <string-name>
            <surname>Vidal</surname>
            ,
            <given-names>M.E.</given-names>
          </string-name>
          ,
          <string-name>
            <surname>Lampo</surname>
            ,
            <given-names>T.</given-names>
          </string-name>
          ,
          <string-name>
            <surname>Castillo</surname>
            ,
            <given-names>J.</given-names>
          </string-name>
          ,
          <string-name>
            <surname>Ruckhaus</surname>
          </string-name>
          , E.:
          <article-title>Anapsid: an adaptive query processing engine for sparql endpoints</article-title>
          .
          <source>In: International Semantic Web Conference</source>
          . pp.
          <volume>18</volume>
          {
          <fpage>34</fpage>
          . Springer (
          <year>2011</year>
          )
        </mixed-citation>
      </ref>
      <ref id="ref2">
        <mixed-citation>
          2.
          <string-name>
            <surname>Aluc</surname>
            ,
            <given-names>G.</given-names>
          </string-name>
          ,
          <string-name>
            <surname>Hartig</surname>
            ,
            <given-names>O.</given-names>
          </string-name>
          ,
          <string-name>
            <given-names>O</given-names>
            zsu, M.T.,
            <surname>Daudjee</surname>
          </string-name>
          ,
          <string-name>
            <surname>K.</surname>
          </string-name>
          :
          <article-title>Diversi ed stress testing of rdf data management systems</article-title>
          . In: International Semantic Web Conference. pp.
          <volume>197</volume>
          {
          <fpage>212</fpage>
          . Springer (
          <year>2014</year>
          )
        </mixed-citation>
      </ref>
      <ref id="ref3">
        <mixed-citation>
          3.
          <string-name>
            <surname>Aluc</surname>
            ,
            <given-names>G.</given-names>
          </string-name>
          ,
          <string-name>
            <surname>Ozsu</surname>
            ,
            <given-names>M.</given-names>
          </string-name>
          ,
          <string-name>
            <surname>Daudjee</surname>
            ,
            <given-names>K.</given-names>
          </string-name>
          ,
          <string-name>
            <surname>Hartig</surname>
          </string-name>
          , O.:
          <article-title>chameleon-db: a workload-aware robust rdf data management system</article-title>
          . university of waterloo.
          <source>Tech. rep., Tech. Rep. CS-2013-10</source>
          (
          <year>2013</year>
          )
        </mixed-citation>
      </ref>
      <ref id="ref4">
        <mixed-citation>
          4.
          <string-name>
            <surname>Bitton</surname>
            ,
            <given-names>D.</given-names>
          </string-name>
          ,
          <string-name>
            <surname>Boral</surname>
            ,
            <given-names>H.</given-names>
          </string-name>
          ,
          <string-name>
            <surname>DeWitt</surname>
          </string-name>
          , D.J.,
          <string-name>
            <surname>Wilkinson</surname>
            ,
            <given-names>W.K.</given-names>
          </string-name>
          :
          <article-title>Parallel algorithms for the execution of relational database operations</article-title>
          .
          <source>ACM Transactions on Database Systems (TODS) 8</source>
          (
          <issue>3</issue>
          ),
          <volume>324</volume>
          {
          <fpage>353</fpage>
          (
          <year>1983</year>
          )
        </mixed-citation>
      </ref>
      <ref id="ref5">
        <mixed-citation>
          5.
          <string-name>
            <surname>Buil-Aranda</surname>
            ,
            <given-names>C.</given-names>
          </string-name>
          ,
          <string-name>
            <surname>Hogan</surname>
            ,
            <given-names>A.</given-names>
          </string-name>
          ,
          <string-name>
            <surname>Umbrich</surname>
            ,
            <given-names>J.</given-names>
          </string-name>
          ,
          <string-name>
            <surname>Vandenbussche</surname>
          </string-name>
          , P.Y.:
          <article-title>Sparql webquerying infrastructure: Ready for action</article-title>
          ? In: International Semantic Web Conference. pp.
          <volume>277</volume>
          {
          <fpage>293</fpage>
          . Springer (
          <year>2013</year>
          )
        </mixed-citation>
      </ref>
      <ref id="ref6">
        <mixed-citation>
          6.
          <string-name>
            <surname>DeWitt</surname>
          </string-name>
          , D.J.,
          <string-name>
            <surname>Naughton</surname>
            ,
            <given-names>J.F.</given-names>
          </string-name>
          ,
          <string-name>
            <surname>Burger</surname>
          </string-name>
          , J.:
          <article-title>Nested loops revisited</article-title>
          .
          <source>In: Parallel and Distributed Information Systems</source>
          ,
          <year>1993</year>
          ., Proceedings of the Second International Conference on. pp.
          <volume>230</volume>
          {
          <fpage>242</fpage>
          .
          <string-name>
            <surname>IEEE</surname>
          </string-name>
          (
          <year>1993</year>
          )
        </mixed-citation>
      </ref>
    </ref-list>
  </back>
</article>