Flavio Paiva Junqueira

dblp:27/2404 · also Flavio Junqueira · DBLP profile ↗
← Back
54ranked-venue papers
14as first author
1since 2021 · last 2023
0009-0000-6789-5505ORCID · corroborated

Domains — the database's venue-derived domains; a paper can count in several

Databases, data management, data science and information retrieval · 23 · 1 first-authorSystems, architecture and hardware · 19 · 8 first-authorArtificial intelligence and machine learning · 6Software engineering, systems software and programming languages · 6 · 1 first-author · 1 since 2021Security and privacy · 3 · 1 first-authorApplied, interdisciplinary, general and emerging computing · 2Computer networks · 1

Expertise — from the expertise taxonomy: the topics of the expert's papers under the CCF categories. A weight counts papers with recency: 1 for a paper about the topic, 0.3 when the topic is its context, halved every five years.

Computer architecture, parallel and distributed computing, and storage systems
19 papers
Distributed systems · 80% Parallel and multicore computing · 5% Cloud and datacenter computing · 5%
Databases, data mining, and information retrieval
14 papers
Information retrieval · 72% Transaction processing and concurrency control · 11% Data integration and cleaning · 5%

Topics — the 30 heaviest of 46, each with the papers that count most for it

TopicWeightPapersLastEvidence papers
Distributed systems
fault tolerance
1.092015
Scalable Error Isolation for Distributed Systems · NSDI 2015
Visigoth fault tolerance · EuroSys 2015
Practical Hardening of Crash-Tolerant Systems · USENIX ATC 2012
Distributed systems
consensus
0.752016
Filo: Consolidated Consensus as a Cloud Service · USENIX ATC 2016
Visigoth fault tolerance · EuroSys 2015
The life and times of a zookeeper · PODC 2009
Information retrieval › search engines
search engine caching
0.552012
Online result cache invalidation for real-time web search · SIGIR 2012
A refreshing perspective of search engine caching · WWW 2010
Caching search engine results over incremental indices · WWW 2010
Distributed systems › fault tolerance
byzantine fault tolerance
0.422015
Visigoth fault tolerance · EuroSys 2015
On the (limited) power of non-equivocation · PODC 2012
Distributed systems
replication
0.442012
Dynamic Reconfiguration of Primary/Backup Clusters · USENIX ATC 2012
ZooKeeper: Wait-free Coordination for Internet-scale Systems · USENIX ATC 2010
Replicating Nondeterministic Services on Grid Environments · HPDC 2006
Information retrieval › distributed information retrieval
distributed search
0.332012
Reactive index replication for distributed search engines · SIGIR 2012
Document assignment in multi-site search engines · WSDM 2011
Challenges on Distributed Web Retrieval · ICDE 2007
Distributed systems
distributed coordination
0.322015
Extensible distributed coordination · EuroSys 2015
ZooKeeper: Wait-free Coordination for Internet-scale Systems · USENIX ATC 2010
Distributed systems › replication
state machine replication
0.322015
Visigoth fault tolerance · EuroSys 2015
Mencius: Building Efficient Replicated State Machine for WANs · OSDI 2008
Information retrieval
web search
0.322013
Exploiting user clicks for automatic seed set generation for entity matching · KDD 2013
Energy-price-driven query processing in multi-center web search engines · SIGIR 2011
Information retrieval › search engines
search engine architecture
0.332010
A refreshing perspective of search engine caching · WWW 2010
Caching search engine results over incremental indices · WWW 2010
Challenges on Distributed Web Retrieval · ICDE 2007
Information retrieval › query processing
list intersection
0.222011
Posting list intersection on multicore architectures · SIGIR 2011
On efficient posting list intersection with multicore processors · SIGIR 2009
Information retrieval
query processing
0.222011
Posting list intersection on multicore architectures · SIGIR 2011
On efficient posting list intersection with multicore processors · SIGIR 2009
Distributed systems › distributed algorithms
distributed synchronization
0.212015
Extensible distributed coordination · EuroSys 2015
Hardware reliability and fault tolerance › fault containment
error confinement
0.212015
Scalable Error Isolation for Distributed Systems · NSDI 2015
Distributed systems › distributed coordination
coordination service
0.222010
ZooKeeper: Wait-free Coordination for Internet-scale Systems · USENIX ATC 2010
The life and times of a zookeeper · PODC 2009
Transaction processing and concurrency control
distributed transaction processing
0.212014
Omid: Lock-free transactional support for distributed data stores · ICDE 2014
Transaction processing and concurrency control › isolation levels
snapshot isolation
0.212014
Omid: Lock-free transactional support for distributed data stores · ICDE 2014
Data integration and cleaning
entity matching
0.212013
Exploiting user clicks for automatic seed set generation for entity matching · KDD 2013
Query processing and optimization
query result caching
0.222008
ResIn: a combination of results caching and index pruning for high-performance web search engines · SIGIR 2008
The impact of caching on search engines · SIGIR 2007
Distributed systems › fault tolerance › failure models
crash failures
0.112012
Practical Hardening of Crash-Tolerant Systems · USENIX ATC 2012
Reconfigurable computing and FPGAs
dynamic reconfiguration
0.112012
Dynamic Reconfiguration of Primary/Backup Clusters · USENIX ATC 2012
Indexing and storage engines › index maintenance
incremental indexing
0.122010
Caching search engine results over incremental indices · WWW 2010
Caching search engine results over incremental indices · SIGIR 2010
Information retrieval › distributed information retrieval
document allocation
0.112011
Document assignment in multi-site search engines · WSDM 2011
Cloud and datacenter computing › resource management
datacenter resource management
0.112011
Energy-price-driven query processing in multi-center web search engines · SIGIR 2011
Hardware accelerators and domain-specific architectures › query processing
energy-efficient query processing
0.112011
Energy-price-driven query processing in multi-center web search engines · SIGIR 2011
Parallel and multicore computing › parallelization strategies
fine-grained parallelization
0.112011
Posting list intersection on multicore architectures · SIGIR 2011
Parallel and multicore computing
parallel algorithms
0.112011
Posting list intersection on multicore architectures · SIGIR 2011
Database system architecture and tuning
cache invalidation
0.112010
Caching search engine results over incremental indices · SIGIR 2010
Distributed systems › replication › geo-replication
wide-area replication
0.112008
Mencius: Building Efficient Replicated State Machine for WANs · OSDI 2008
Information retrieval
search engines
0.122012
Reactive index replication for distributed search engines · SIGIR 2012
Posting list intersection on multicore architectures · SIGIR 2011

Methods — techniques the papers use, named apart from their topics

snapshot isolation · 0.4centralized transaction metadata · 0.4heuristic · 0.3approximation algorithm · 0.3consensus as a service · 0.2threshold-based timing model · 0.2dynamic extension · 0.2access control · 0.2random walk with restart · 0.2co-clustering · 0.2classifier · 0.2online replication · 0.1hardening · 0.1fault injection · 0.1document popularity · 0.1workload shifting · 0.1query log analysis · 0.1machine-learned assignment · 0.1
YearPublicationVenuePosition
2023 Pravega: A Tiered Storage System for Data Streams
abstract
The growing popularity of the data stream abstraction entails new challenging requirements when it comes to data ingestion and storage. Many organizations expect to retain data streams for extended periods of time and to store such stream data in a cost-effective manner. It is also crucial to reconcile apparently opposite properties, like data durability and consistency, along with high performance. Furthermore, data streams should not only deal with a high degree of parallelism, but also adapt to fluctuating workloads with little or no admin intervention. To our knowledge, no storage system for data streams fully copes with all these requirements.
Raúl Gracia Tinedo, Flavio Paiva Junqueira, Tom Kaitchuck, Sachin Joshi
Middleware2
2018 The Tortoise and the Hare: Characterizing Synchrony in Distributed Environments (Practical Experience Report)
abstract
The design of distributed protocols that run in data centers and enterprise clusters is heavily dependent on synchrony assumptions regarding the timing behavior of the participating nodes and the network. However, little is known about the actual synchrony of real distributed systems, and how it varies across deployments. To better understand this timing behavior and how it impacts the design and implementation of distributed protocols, we conduct an extensive measurement study of the latency for transmitting and processing messages between nodes in four different environments. Our study determines how protocol characteristics affect the latency behavior. We also determine how different environmental factors can affect the measured latency and whether high latency events manifest globally or locally. Our results suggest several directions for reducing latency, and for leveraging recent distributed computing models in a more judicious way.
Daniel Porto 0002, João Leitão 0001, Flavio Paiva Junqueira, Rodrigo Rodrigues 0001
DSN3
2016 Filo: Consolidated Consensus as a Cloud Service
Parisa Jalili Marandi, Christos Gkantsidis, Flavio Paiva Junqueira, Dushyanth Narayanan
USENIX ATC3
2015 Extensible distributed coordination
abstract
Most services inside a data center are distributed systems requiring coordination and synchronization in the form of primitives like distributed locks and message queues. We argue that extensibility is a crucial feature of the coordination infrastructures used in these systems. Without the ability to extend the functionality of coordination services, applications might end up using sub-optimal coordination algorithms, possibly leading to low performance. Adding extensibility, however, requires mechanisms that constrain extensions to be able to make reasonable security and performance guarantees. We propose a scheme that enables extensions to be introduced and removed dynamically in a secure way. To avoid performance overheads due to poorly designed extensions, it constrains the access of extensions to resources. Evaluation results for extensible versions of ZooKeeper and DepSpace show that it is possible to increase the throughput of a distributed queue by more than an order of magnitude (17x for ZooKeeper, 24x for DepSpace) while keeping the underlying coordination kernel small.
Tobias Distler, Christopher Bahn, Alysson Neves Bessani, Frank Fischer 0004, Flavio Paiva Junqueira
EuroSys5
2015 Visigoth fault tolerance
abstract
We present a new technique for designing distributed protocols for building reliable stateful services called Visigoth Fault Tolerance (VFT). VFT introduces the Visigoth model, which makes it possible to calibrate the timing assumptions of a system using a threshold of slow processes or messages, and also to distinguish between non-malicious arbitrary faults and correlated attack scenarios. This enables solutions that leverage the characteristics of data center systems, namely their secure environment and predictable performance, in order to allow replicated systems to be more efficient with respect to the utilization of resources than those designed under asynchrony and Byzantine assumptions, while avoiding the need to make a system synchronous, or to restrict failure modes to silent crashes. We implemented a VFT protocol for a state machine replication library, and ran several benchmarks. Our evaluation shows that VFT has comparable performance to existing schemes and brings significant benefits in terms of the throughput per dollar, i.e., the server cost for sustaining a certain level of request execution.
Daniel Porto 0002, João Leitão 0001, Cheng Li 0001, Allen Clement, Aniket Kate, Flavio Paiva Junqueira, Rodrigo Rodrigues 0001
EuroSys6
2015 Scalable Error Isolation for Distributed Systems
Diogo Behrens, Marco Serafini, Flavio Paiva Junqueira, Sergei Arnautov, Christof Fetzer
NSDI3
2014 Omid: Lock-free transactional support for distributed data stores
abstract
In this paper, we introduce Omid, a tool for lock-free transactional support in large data stores such as HBase. Omid uses a centralized scheme and implements snapshot isolation, a property that guarantees that all read operations of a transaction are performed on a consistent snapshot of the data. In a lock-based approach, the unreleased, distributed locks that are held by a failed or slow client block others. By using a centralized scheme for Omid, we are able to implement a lock-free commit algorithm, which does not suffer from this problem. Moreover, Omid lightly replicates a read-only copy of the transaction metadata into the clients where they can locally service a large part of queries on metadata. Thanks to this technique, Omid does not require modifying either the source code of the data store or the tables' schema, and the overhead on data servers is also negligible. The experimental results show that our implementation on a simple dual-core machine can service up to a thousand of client machines. While the added latency is limited to only 10 ms, Omid scales up to 124K write transactions per second. Since this capacity is multiple times larger than the maximum reported traffic in similar systems, we do not expect the centralized scheme of Omid to be a bottleneck even for current large data stores.
Daniel Gómez Ferro, Flavio Paiva Junqueira, Ivan Kelly, Benjamin C. Reed, Maysam Yabandeh
ICDE2
2014 Slider: incremental sliding window analytics
abstract
Sliding window analytics is often used in distributed data-parallel computing for analyzing large streams of continuously arriving data. When pairs of consecutive windows overlap, there is a potential to update the output incrementally, more efficiently than recomputing from scratch. However, in most systems, realizing this potential requires programmers to explicitly manage the intermediate state for overlapping windows, and devise an application-specific algorithm to incrementally update the output.
Pramod Bhatotia, Umut A. Acar, Flavio Paiva Junqueira, Rodrigo Rodrigues 0001
Middleware3
2013 Cache refreshing for online social news feeds
abstract
Several social networking applications enable users to view the events generated by other users, typically friends in the social network, in the form of ``news feeds''. Friends and events are typically maintained per user and cached in memory to enable efficient generation of news feeds. Caching user friends and events, however, raises concerns about the freshness of news feeds as users may not observe the most recent events when cache content becomes stale. Mechanisms to keep cache content fresh are thus critical for user satisfaction while computing news feeds efficiently through caching.
Xiao Bai 0002, Flavio Paiva Junqueira, Adam Silberstein
CIKM2
2013 Exploiting user clicks for automatic seed set generation for entity matching
abstract
Matching entities from different information sources is a very important problem in data analysis and data integration. It is, however, challenging due to the number and diversity of information sources involved, and the significant editorial efforts required to collect sufficient training data. In this paper, we present an approach that leverages user clicks during Web search to automatically generate training data for entity matching. The key insight of our approach is that Web pages clicked for a given query are likely to be about the same entity. We use random walk with restart to reduce data sparseness, rely on co-clustering to group queries and Web pages, and exploit page similarity to improve matching precision. Experimental results show that: (i) With 360K pages from 6 major travel websites, we obtain 84K matchings (of 179K pages) that refer to the same entities, with an average precision of 0.826; (ii) The quality of matching obtained from a classifier trained on the resulted seed data is promising: the performance matches that of editorial data at small size and improves with size.
Xiao Bai 0002, Flavio Paiva Junqueira, Srinivasan H. Sengamedu
KDD2
2013 DynaSoRe: Efficient In-Memory Store for Social Applications
Xiao Bai 0002, Arnaud Jégou, Flavio Paiva Junqueira, Vincent Leroy 0001
Middleware3
2013 On Barriers and the Gap between Active and Passive Replication
Flavio Paiva Junqueira, Marco Serafini
DISC1
2013 Piggybacking on Social Networks
abstract
The popularity of social-networking sites has increased rapidly over the last decade. A basic functionalities of social-networking sites is to present users with streams of events shared by their friends. At a systems level, materialized per-user views are a common way to assemble and deliver such event streams on-line and with low latency. Access to the data stores, which keep the user views, is a major bottleneck of social-networking systems. We propose to improve the throughput of these systems by using social piggybacking, which consists of processing the requests of two friends by querying and updating the view of a third common friend. By using one such hub view, the system can serve requests of the first friend without querying or updating the view of the second. We show that, given a social graph, social piggybacking can minimize the overall number of requests, but computing the optimal set of hubs is an NP-hard problem. We propose anO(logn) approximation algorithm and a heuristic to solve the problem, and evaluate them using the full Twitter and Flickr social graphs, which have up to billions of edges. Compared to existing approaches, using social piggybacking results in similar throughput in systems with few servers, but enables substantial throughput improvements as the size of the system grows, reaching up to a 2-factor increase. We also evaluate our algorithms on a real social networking system prototype and we show that the actual increase in throughput corresponds nicely to the gain anticipated by our cost function.
Aristides Gionis, Flavio Paiva Junqueira, Vincent Leroy 0001, Marco Serafini, Ingmar Weber
Proc. VLDB Endow.2
2012 Scalable deferred update replication
abstract
Deferred update replication is a well-known approach to building data management systems as it provides both high availability and high performance. High availability comes from the fact that any replica can execute client transactions; the crash of one or more replicas does not interrupt the system. High performance comes from the fact that only one replica executes a transaction; the others must only apply its updates. Since replicas execute transactions concurrently, transaction execution is distributed across the system. The main drawback of deferred update replication is that update transactions scale poorly with the number of replicas, although read-only transactions scale well. This paper proposes an extension to the technique that improves the scalability of update transactions. In addition to presenting a novel protocol, we detail its implementation and provide an extensive analysis of its performance.
Daniele Sciascia, Fernando Pedone, Flavio Paiva Junqueira
DSN3
2012 On the (limited) power of non-equivocation
abstract
In recent years, there have been a few proposals to add a small amount of trusted hardware at each replica in a Byzantine fault tolerant system to cut back replication factors. These trusted components eliminate the ability for a Byzantine node to perform equivocation, which intuitively means making conflicting statements to different processes.
Allen Clement, Flavio Paiva Junqueira, Aniket Kate, Rodrigo Rodrigues 0001
PODC2
2012 Online result cache invalidation for real-time web search
abstract
Caches of results are critical components of modern Web search engines, since they enable lower response time to frequent queries and reduce the load to the search engine backend. Results in long-lived cache entries may become stale, however, as search engines continuously update their index to incorporate changes to the Web. Consequently, it is important to provide mechanisms that control the degree of staleness of cached results, ideally enabling the search engine to always return fresh results. In this paper, we present a new mechanism that identifies and invalidates query results that have become stale in the cache online. The basic idea is to evaluate at query time and against recent changes if cache hits have had their results have changed. For enhancing invalidation efficiency, the generation time of cached queries and their chronological order with respect to the latest index update are used to early prune unaffected queries. We evaluate the proposed approach using documents that change over time and query logs of the Yahoo! search engine. We show that the proposed approach ensures good query results (50% fewer stale results) and high invalidation accuracy (90% fewer unnecessary invalidations) compared to a baseline approach that makes invalidation decisions off-line. More importantly, the proposed approach induces less processing overhead, ensuring an average throughput 73% higher than that of the baseline approach.
Xiao Bai 0002, Flavio Paiva Junqueira
SIGIR2
2012 Reactive index replication for distributed search engines
abstract
Distributed search engines comprise multiple sites deployed across geographically distant regions, each site being specialized to serve the queries of local users. When a search site cannot accurately compute the results of a query, it must forward the query to other sites. This paper considers the problem of selecting the documents indexed by each site focusing on replication to increase the fraction of queries processed locally. We propose RIP, an algorithm for replicating documents and posting lists that is practical and has two important features. RIP evaluates user interests in an online fashion and uses only local data of a site. Being an online approach simplifies the operational complexity, while locality enables higher performance when processing queries and documents. The decision procedure, on top of being online and local, incorporates document popularity and user queries, which is critical when assuming a replication budget for each site. Having a replication budget reflects the hardware constraints of any given site. We evaluate RIP against the approach of replicating popular documents statically, and show that we achieve significant gains, while having the additional benefit of supporting incremental indexes.
Flavio Paiva Junqueira, Vincent Leroy 0001, Matthieu Morel
SIGIR1
2012 Practical Hardening of Crash-Tolerant Systems
Miguel Correia 0001, Daniel Gómez Ferro, Flavio Paiva Junqueira, Marco Serafini
USENIX ATC3
2012 Dynamic Reconfiguration of Primary/Backup Clusters
Alexander Shraer, Benjamin C. Reed, Dahlia Malkhi, Flavio Paiva Junqueira
USENIX ATC4
2012 Brief Announcement: Consensus and Efficient Passive Replication
Flavio Paiva Junqueira, Marco Serafini
DISC1
2012 A five-level static cache architecture for web search engines
Rifat Ozcan, Ismail Sengör Altingövde, Berkant Barla Cambazoglu, Flavio Paiva Junqueira, Özgür Ulusoy
Inf. Process. Manag.4
2011 Discovering URLs through user feedback
abstract
Search engines rely upon crawling to build their Web page collections. A Web crawler typically discovers new URLs by following the link structure induced by links on Web pages. As the number of documents on the Web is large, discovering newly created URLs may take arbitrarily long, and depending on how a given page is connected to others, such a crawler may miss the pages altogether. In this paper, we evaluate the benefits of integrating a passive URL discovery mechanism into a Web crawler. This mechanism is passive in the sense that it does not require the crawler to actively fetch documents from the Web to discover URLs. We focus here on a mechanism that uses toolbar data as a representative source for new URL discovery. We use the toolbar logs of Yahoo! to characterize the URLs that are accessed by users via their browsers, but not discovered by Yahoo! Web crawler. We show that a high fraction of URLs that appear in toolbar logs are not discovered by the crawler. We also reveal that a certain fraction of URLs are discovered by the crawler later than the time they are first accessed by users. One important conclusion of our work is that web search engines can highly benefit from user feedback in the form of toolbar logs for passive URL discovery.
Xiao Bai 0002, Berkant Barla Cambazoglu, Flavio Paiva Junqueira
CIKM3
2011 Assigning documents to master sites in distributed search
abstract
An appealing solution to scale Web search with the growth of the Internet is the use of distributed architectures. Distributed search engines rely on multiple sites deployed in distant regions across the world, where each site is specialized to serve queries issued by the users of its region. This paper investigates the problem of assigning each document to a master site. We show that by leveraging similarities between a document and the activity of the users, we can accurately detect which site is the most relevant to place a document. We conduct various experiments using two document assignment approaches, showing performance improvements of up to 20.8% over a baseline technique which assigns the documents to search sites based on their language.
Roi Blanco, Berkant Barla Cambazoglu, Flavio Paiva Junqueira, Ivan Kelly, Vincent Leroy 0001
CIKM3
2011 Zab: High-performance broadcast for primary-backup systems
abstract
Zab is a crash-recovery atomic broadcast algorithm we designed for the ZooKeeper coordination service. ZooKeeper implements a primary-backup scheme in which a primary process executes clients operations and uses Zab to propagate the corresponding incremental state changes to backup processes. Due the dependence of an incremental state change on the sequence of changes previously generated, Zab must guarantee that if it delivers a given state change, then all other changes it depends upon must be delivered first. Since primaries may crash, Zab must satisfy this requirement despite crashes of primaries.
Flavio Paiva Junqueira, Benjamin C. Reed, Marco Serafini
DSN1
2011 Leader Election for Replicated Services Using Application Scores
Diogo Becker, Flavio Paiva Junqueira, Marco Serafini
Middleware2
2011 Energy-price-driven query processing in multi-center web search engines
abstract
Concurrently processing thousands of web queries, each with a response time under a fraction of a second, necessitates maintaining and operating massive data centers. For large-scale web search engines, this translates into high energy consumption and a huge electric bill. This work takes the challenge to reduce the electric bill of commercial web search engines operating on data centers that are geographically far apart. Based on the observation that energy prices and query workloads show high spatio-temporal variation, we propose a technique that dynamically shifts the query workload of a search engine between its data centers to reduce the electric bill. Experiments on real-life query workloads obtained from a commercial search engine show that significant financial savings can be achieved by this technique.
Enver Kayaaslan, Berkant Barla Cambazoglu, Roi Blanco, Flavio Paiva Junqueira, Cevdet Aykanat
SIGIR4
2011 Posting list intersection on multicore architectures
abstract
In current commercial Web search engines, queries are processed in the conjunctive mode, which requires the search engine to compute the intersection of a number of posting lists to determine the documents matching all query terms. In practice, the intersection operation takes a significant fraction of the query processing time, for some queries dominating the total query latency. Hence, efficient posting list intersection is critical for achieving short query latencies. In this work, we focus on improving the performance of posting list intersection by leveraging the compute capabilities of recent multicore systems. To this end, we consider various coarse-grained and fine-grained parallelization models for list intersection. Specifically, we present an algorithm that partitions the work associated with a given query into a number of small and independent tasks that are subsequently processed in parallel. Through a detailed empirical analysis of these alternative models, we demonstrate that exploiting parallelism at the finest-level of granularity is critical to achieve the best performance on multicore systems. On an eight-core system, the fine-grained parallelization method is able to achieve more than five times reduction in average query processing time while still exploiting the parallelism for high query throughput.
Shirish Tatikonda, Berkant Barla Cambazoglu, Flavio Paiva Junqueira
SIGIR3
2011 Document assignment in multi-site search engines
abstract
Assigning documents accurately to sites is critical for the performance of multi-site Web search engines. In such settings, sites crawl only documents they index and forward queries to obtain best-matching documents from other sites. Inaccurate assignments may lead to inefficiencies when crawling Web pages or processing user queries. In this work, we propose a machine-learned document assignment strategy that uses the locality of document views in search results to decide upon assignments. We evaluate the performance of our strategy using various document features extracted from a large Web collection. Our experimental setup uses query logs from a number of search front-ends spread across different geographic locations and uses these logs to learn the document access patterns. We compare our technique against baselines such as region- and language-based document assignment and observe that our technique achieves substantial performance improvements with respect to recall. With our technique, we are able to obtain a small query forwarding rate (0.04) requiring roughly 45% less replication of documents compared to replicating all documents across all sites.
Ulf Brefeld, Berkant Barla Cambazoglu, Flavio Paiva Junqueira
WSDM3
2010 Caching search engine results over incremental indices
abstract
A Web search engine must update its index periodically to incorporate changes to the Web. We argue in this paper that index updates fundamentally impact the design of search engine result caches, a performance-critical component of modern search engines. Index updates lead to the problem of cache invalidation: invalidating cached entries of queries whose results have changed. Naive approaches, such as flushing the entire cache upon every index update, lead to poor performance and in fact, render caching futile when the frequency of updates is high. Solving the invalidation problem efficiently corresponds to predicting accurately which queries will produce different results if re-evaluated, given the actual changes to the index.
Roi Blanco, Edward Bortnikov, Flavio Paiva Junqueira, Ronny Lempel, Luca Telloli, Hugo Zaragoza
SIGIR3
2010 ZooKeeper: Wait-free Coordination for Internet-scale Systems
Patrick Hunt, Mahadev Konar, Flavio Paiva Junqueira, Benjamin C. Reed
USENIX ATC3
2010 Caching search engine results over incremental indices
abstract
A Web search engine must update its index periodically to incorporate changes to the Web, and we argue in this work that index updates fundamentally impact the design of search engine result caches. Index updates lead to the problem of cache invalidation: invalidating cached entries of queries whose results have changed. To enable efficient invalidation of cached results, we propose a framework for developing invalidation predictors and some concrete predictors. Evaluation using Wikipedia documents and a query log from Yahoo! shows that selective invalidation of cached search results can lower the number of query re-evaluations by as much as 30% compared to a baseline time-to-live scheme, while returning results of similar freshness.
Roi Blanco, Edward Bortnikov, Flavio Paiva Junqueira, Ronny Lempel, Luca Telloli, Hugo Zaragoza
WWW3
2010 A refreshing perspective of search engine caching
abstract
Commercial Web search engines have to process user queries over huge Web indexes under tight latency constraints. In practice, to achieve low latency, large result caches are employed and a portion of the query traffic is served using previously computed results. Moreover, search engines need to update their indexes frequently to incorporate changes to the Web. After every index update, however, the content of cache entries may become stale, thus decreasing the freshness of served results. In this work, we first argue that the real problem in today's caching for large-scale search engines is not eviction policies, but the ability to cope with changes to the index, i.e., cache freshness. We then introduce a novel algorithm that uses a time-to-live value to set cache entries to expire and selectively refreshes cached results by issuing refresh queries to back-end search clusters. The algorithm prioritizes the entries to refresh according to a heuristic that combines the frequency of access with the age of an entry in the cache. In addition, for setting the rate at which refresh queries are issued, we present a mechanism that takes into account idle cycles of back-end servers. Evaluation using a real workload shows that our algorithm can achieve hit rate improvements as well as reduction in average hit ages. An implementation of this algorithm is currently in production use at Yahoo!.
Berkant Barla Cambazoglu, Flavio Paiva Junqueira, Vassilis Plachouras, Scott A. Banachowski, Baoqiu Cui, Swee Lim, Bill Bridge
WWW2
2010 Threshold protocols in survivor set systems
Flavio Paiva Junqueira, Keith Marzullo, Maurice Herlihy, Lucia Draque Penso
Distributed Comput.1
2009 Interactive Analysis of Web-Scale Data
Christopher Olston, Edward Bortnikov, Khaled Elmeleegy, Flavio Paiva Junqueira, Benjamin C. Reed
CIDR4
2009 On the feasibility of multi-site web search engines
abstract
Web search engines are often implemented as centralized systems. Designing and implementing a Web search engine in a distributed environment is a challenging engineering task that encompasses many interesting research questions. However, distributing a search engine across multiple sites has several advantages, such as utilizing less compute resources and exploiting data locality. In this paper we investigate the cost-effectiveness of building a distributed Web search engine. We propose a model for assessing the total cost of a distributed Web search engine that includes the computational costs and the communication cost among all distributed sites. We then present a query-processing algorithm that maximizes the amount of queries answered locally, without sacrificing the quality of the results compared to a centralized search engine. We simulate the algorithm on real document collections and query workloads to measure the actual parameters needed for our cost model, and we show that a distributed search engine can be competitive compared to a centralized architecture with respect to real cost.
Ricardo Baeza-Yates, Aristides Gionis, Flavio Paiva Junqueira, Vassilis Plachouras, Luca Telloli
CIKM3
2009 Introduction
Dejan Kostic, Guillaume Pierre, Flavio Paiva Junqueira, Peter R. Pietzuch
Euro-Par3
2009 The life and times of a zookeeper
abstract
A distributed application may comprise a range of processors, from two to thousands. For the processors of an application comprising many different machines to work together, they must often coordinate. Coordination may be as simple as agreeing upon a basic configuration, e.g., a list of server addresses and ports for example. In the context of distributed systems, even agreeing upon a basic configuration can be complicated: servers may be down during configuration changes, servers may need to adapt quickly to dynamic configuration changes, and during changes servers may have conflicting configurations. Most distributed applications also require more sophisticated coordination primitives such as leader election, group membership, and rendezvous.At Yahoo!, we noticed that many distributed applications were re-implementing coordination primitives in their distributed applications. In theory, this is completely reasonable: coordination protocols are well known, and usually the coordination logic is intermingled with the application logic. In practice, the story is much different; these protocols have sometimes subtle requirements that can easily be overlooked when implementing them. Further, an application developer is usually much more interested in working on application logic than the coordination protocol that the logic depends on. We found many cases of applications whose coordination primitives were buggy, a single point of failure, poorly performing, or oversimplified; in some cases the applications suffered from all of the above.We started ZooKeeper as a system that could address all these problems in a general way so that all of our applications could use it for coordination and application developers could focus on developing their applications. By providing a general system that all applications could use, we could devote the time to making it robust, fault-tolerant, and with good enough performance to be used extensively by applications. We also needed to balance two possibly conflicting goals: ZooKeeper needed to be general enough to address our coordination needs and simple enough to implement a correct high performance service. We found that we were able to achieve our goals by trading strong synchronization for strong ordering guarantees and a wait-free interface.Inside and outside of Yahoo! developers of distributed applications have enthusiastically embraced ZooKeeper. Developers who have previously dealt with the difficulties of implementing distributed applications quickly see the benefits of ZooKeeper, and we have many applications inside of Yahoo! that make use of it. It was developed quickly by a small group of developers in a small amount of time.The general techniques we use in ZooKeeper are not fundamentally new. For instance, ZooKeeper uses a leader-based atomic broadcast protocol that may seem at a first glance not novel. However, a deeper inspection reveals that it has some key propertiesthat are crucial to guarantee the properties that developers of large-scale applications require. To the best of our knowledge, no previous algorithm presented all necessary properties out-of-the-box.In this talk, we share previous work that influenced ZooKeeper. In particular, we compare at the protocol level with Paxos and Viewstamped Replication, and at the service level with services like Chubby and Boxwood. We give a high level overview of the service itself and describe its implementation. Along the way we point out some of the design and implementation details that were key in achieving the correctness and performance that we need. Finally, we show how ZooKeeper performance has evolved over time. We have increased the performance of ZooKeeper an order of magnitude from our first implementation, and none of that increase has come from protocol changes. Indeed, our protocol has remained constant from the very beginning.
Flavio Paiva Junqueira, Benjamin C. Reed
PODC1
2009 On efficient posting list intersection with multicore processors
abstract
No abstract available.
Shirish Tatikonda, Flavio Paiva Junqueira, Berkant Barla Cambazoglu, Vassilis Plachouras
SIGIR2
2009 The life and times of a zookeeper
abstract
A distributed application may comprise a range of processors, from two to thousands. For the processors of an application comprising many different machines to work together, they must often coordinate. Coordination may be as simple as agreeing upon a basic configuration, e.g., a list of server addresses and ports for example. In the context of distributed systems, even agreeing upon a basic configuration can be complicated: servers may be down during configuration changes, servers may need to adapt quickly to dynamic configuration changes, and during changes servers may have conflicting configurations. Most distributed applications also require more sophisticated coordination primitives such as leader election, group membership, and rendezvous.
Flavio Paiva Junqueira, Benjamin C. Reed
SPAA1
2009 Brief Announcement Zab: A Practical Totally Ordered Broadcast Protocol
Flavio Paiva Junqueira, Benjamin C. Reed
DISC1
2008 Mencius: Building Efficient Replicated State Machine for WANs
Yanhua Mao, Flavio Paiva Junqueira, Keith Marzullo
OSDI2
2008 ResIn: a combination of results caching and index pruning for high-performance web search engines
abstract
Results caching is an efficient technique for reducing the query processing load, hence it is commonly used in real search engines. This technique, however, bounds the maximum hit rate due to the large fraction of singleton queries, which is an important limitation. In this paper we propose ResIn - an architecture that uses a combination of results caching and index pruning to overcome this limitation.
Gleb Skobeltsyn, Flavio Paiva Junqueira, Vassilis Plachouras, Ricardo Baeza-Yates
SIGIR2
2008 Optimizing Threshold Protocols in Adversarial Structures
Maurice Herlihy, Flavio Paiva Junqueira, Keith Marzullo, Lucia Draque Penso
DISC2
2008 Design trade-offs for search engine caching
abstract
In this article we study the trade-offs in designing efficient caching systems for Web search engines. We explore the impact of different approaches, such as static vs. dynamic caching, and caching query results vs. caching posting lists. Using a query log spanning a whole year, we explore the limitations of caching and we demonstrate that caching posting lists can achieve higher hit rates than caching query answers. We propose a new algorithm for static caching of posting lists, which outperforms previous methods. We also study the problem of finding the optimal way to split the static cache between answers and posting lists. Finally, we measure how the changes in the query log influence the effectiveness of static caching, given our observation that the distribution of the queries changes slowly over time. Our results and observations are applicable to different levels of the data-access hierarchy, for instance, for a memory/disk layer or a broker/remote server layer.
Ricardo Baeza-Yates, Aristides Gionis, Flavio Paiva Junqueira, Vanessa Murdock 0001, Vassilis Plachouras, Fabrizio Silvestri
ACM Trans. Web3
2007 Challenges on Distributed Web Retrieval
abstract
In the ocean of Web data, Web search engines are the primary way to access content. As the data is on the order of petabytes, current search engines are very large centralized systems based on replicated clusters. Web data, however, is always evolving. The number of Web sites continues to grow rapidly and there are currently more than 20 billion indexed pages. In the near future, centralized systems are likely to become ineffective against such a load, thus suggesting the need of fully distributed search engines. Such engines need to achieve the following goals: high quality answers, fast response time, high query throughput, and scalability. In this paper we survey and organize recent research results, outlining the main challenges of designing a distributed Web retrieval system.
Ricardo Baeza-Yates, Carlos Castillo 0001, Flavio Paiva Junqueira, Vassilis Plachouras, Fabrizio Silvestri
ICDE3
2007 The impact of caching on search engines
abstract
In this paper we study the trade-offs in designing efficient caching systems for Web search engines. We explore the impact of different approaches, such as static vs. dynamic caching, and caching query results vs.caching posting lists. Using a query log spanning a whole year we explore the limitations of caching and we demonstrate that caching posting lists can achieve higher hit rates than caching query answers. We propose a new algorithm for static caching of posting lists, which outperforms previous methods. We also study the problem of finding the optimal way to split the static cache between answers and posting lists. Finally, we measure how the changes in the query log affect the effectiveness of static caching, given our observation that the distribution of the queries changes slowly over time. Our results and observations are applicable to different levels of the data-access hierarchy, for instance, for a memory/disk layer or a broker/remote server layer.
Ricardo Baeza-Yates, Aristides Gionis, Flavio Paiva Junqueira, Vanessa Murdock 0001, Vassilis Plachouras, Fabrizio Silvestri
SIGIR3
2007 Admission Policies for Caches of Search Engine Results
Ricardo Baeza-Yates, Flavio Paiva Junqueira, Vassilis Plachouras, Hans Friedrich Witschel
SPIRE2
2007 A framework for the design of dependent-failure algorithms
abstract
Abstract Dependent failures constitute a real problem in distributed systems. In this paper, we present a framework for the design of distributed algorithms for systems in which failures of processes are not necessarily independent or identically distributed. To derive this framework, we revisit the traditional way of designing distributed algorithms: assuming a threshold t on the number of process failures and determining constraints on process replication of the form All mathematics should have $n > k\cdot t$ . Under this model, there are several important results in the literature, and we use observations from these results to derive a framework that enables more expressive characterizations of failures, but still captures the essence of previous results. Our framework then has two parts: a characterization of the subsets of processes that can fail in an execution and properties that express constraints on process replication, called replication predicates. After presenting our model to characterize failures, we first consider a class of replication predicates that represent most well‐known problems in distributed computing. Second, we extend this set to also include some predicates for problems that have unusual bounds. We argue that, although being unusual, they have important practical implications, such as using fewer replicas. Copyright © 2007 John Wiley & Sons, Ltd.
Flavio Paiva Junqueira, Keith Marzullo
Concurr. Comput. Pract. Exp.1
2006 Replicating Nondeterministic Services on Grid Environments
abstract
Replication is a technique commonly used to increase the availability of services in distributed systems, including grid and Web services. While replication is relatively easy for services with fully deterministic behavior, grid and Web services often include nondeterministic operations. The traditional way to replicate such nondeterministic services is to use the primary-backup approach. While this is straightforward in synchronous systems with perfect failure detection, typical grid environments are not usually considered to be synchronous systems. This paper addresses the problem of replicating nondeterministic services by designing a protocol based on Paxos and proposing two performance optimizations suitable for replicated grid services. The first improves the performance in the case where some service operations do not change the service state, while the second optimizes grid service requests that use transactions. Evaluations done both on a local cluster and on Planet-Lab demonstrate that these optimizations significantly reduce the service response time and increase the throughput of replicated services
Xianan Zhang, Flavio Paiva Junqueira, Matti A. Hiltunen, Keith Marzullo, Richard D. Schlichting
HPDC2
2005 Replication Predicates for Dependent-Failure Algorithms
Flavio Paiva Junqueira, Keith Marzullo
Euro-Par1
2005 Surviving Internet Catastrophes
Flavio Paiva Junqueira, Ranjita Bhagwan, Alejandro Hevia, Keith Marzullo, Geoffrey M. Voelker
USENIX ATC, General Track1
2005 Coterie Availability in Sites
Flavio Paiva Junqueira, Keith Marzullo
DISC1
2003 The Phoenix Recovery System: Rebuilding from the Ashes of an Internet Catastrophe
Flavio Paiva Junqueira, Ranjita Bhagwan, Keith Marzullo, Stefan Savage, Geoffrey M. Voelker
HotOS1
2003 Synchronous Consensus for Dependent Process Failure
abstract
We present a new abstraction to replace the t of n assumption used in designing fault-tolerant algorithms. This abstraction models dependent process failures yet it is as simple to use as the t of n assumption. To illustrate this abstraction, we consider Consensus for synchronous systems with both crash and arbitrary process failures. By considering failure correlations, we are able to reduce latency and enable the solution of Consensus for system configurations in which it is not possible when forced to use algorithms designed under the t of n assumption. We show that, in general, the number of rounds required in the worst case when assuming crash failures is different from the number of rounds required when assuming arbitrary failures. This is in contrast with the traditional result under the t of n assumption.
Flavio Paiva Junqueira, Keith Marzullo
ICDCS1