EDBT 2026 Demo / reviewers in the wild / expert
Flavio Paiva Junqueira
dblp:27/2404 · also Flavio Junqueira
· DBLP profile ↗
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
| Topic | Weight | Papers | Last | Evidence papers |
|---|---|---|---|---|
Distributed systems
fault tolerance |
1.0 | 9 | 2015 | 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.7 | 5 | 2016 | 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.5 | 5 | 2012 | 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.4 | 2 | 2015 | Visigoth fault tolerance · EuroSys 2015 On the (limited) power of non-equivocation · PODC 2012 |
Distributed systems
replication |
0.4 | 4 | 2012 | 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.3 | 3 | 2012 | 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.3 | 2 | 2015 | Extensible distributed coordination · EuroSys 2015 ZooKeeper: Wait-free Coordination for Internet-scale Systems · USENIX ATC 2010 |
Distributed systems › replication
state machine replication |
0.3 | 2 | 2015 | Visigoth fault tolerance · EuroSys 2015 Mencius: Building Efficient Replicated State Machine for WANs · OSDI 2008 |
Information retrieval
web search |
0.3 | 2 | 2013 | 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.3 | 3 | 2010 | 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.2 | 2 | 2011 | Posting list intersection on multicore architectures · SIGIR 2011 On efficient posting list intersection with multicore processors · SIGIR 2009 |
Information retrieval
query processing |
0.2 | 2 | 2011 | 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.2 | 1 | 2015 | Extensible distributed coordination · EuroSys 2015 |
Hardware reliability and fault tolerance › fault containment
error confinement |
0.2 | 1 | 2015 | Scalable Error Isolation for Distributed Systems · NSDI 2015 |
Distributed systems › distributed coordination
coordination service |
0.2 | 2 | 2010 | 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.2 | 1 | 2014 | Omid: Lock-free transactional support for distributed data stores · ICDE 2014 |
Transaction processing and concurrency control › isolation levels
snapshot isolation |
0.2 | 1 | 2014 | Omid: Lock-free transactional support for distributed data stores · ICDE 2014 |
Data integration and cleaning
entity matching |
0.2 | 1 | 2013 | Exploiting user clicks for automatic seed set generation for entity matching · KDD 2013 |
Query processing and optimization
query result caching |
0.2 | 2 | 2008 | 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.1 | 1 | 2012 | Practical Hardening of Crash-Tolerant Systems · USENIX ATC 2012 |
Reconfigurable computing and FPGAs
dynamic reconfiguration |
0.1 | 1 | 2012 | Dynamic Reconfiguration of Primary/Backup Clusters · USENIX ATC 2012 |
Indexing and storage engines › index maintenance
incremental indexing |
0.1 | 2 | 2010 | 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.1 | 1 | 2011 | Document assignment in multi-site search engines · WSDM 2011 |
Cloud and datacenter computing › resource management
datacenter resource management |
0.1 | 1 | 2011 | 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.1 | 1 | 2011 | Energy-price-driven query processing in multi-center web search engines · SIGIR 2011 |
Parallel and multicore computing › parallelization strategies
fine-grained parallelization |
0.1 | 1 | 2011 | Posting list intersection on multicore architectures · SIGIR 2011 |
Parallel and multicore computing
parallel algorithms |
0.1 | 1 | 2011 | Posting list intersection on multicore architectures · SIGIR 2011 |
Database system architecture and tuning
cache invalidation |
0.1 | 1 | 2010 | Caching search engine results over incremental indices · SIGIR 2010 |
Distributed systems › replication › geo-replication
wide-area replication |
0.1 | 1 | 2008 | Mencius: Building Efficient Replicated State Machine for WANs · OSDI 2008 |
Information retrieval
search engines |
0.1 | 2 | 2012 | 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
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2023 | Pravega: A Tiered Storage System for Data StreamsabstractThe 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 |
Middleware | 2 |
| 2018 | The Tortoise and the Hare: Characterizing Synchrony in Distributed Environments (Practical Experience Report)abstractThe 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 |
DSN | 3 |
| 2016 | Filo: Consolidated Consensus as a Cloud Service
Parisa Jalili Marandi, Christos Gkantsidis, Flavio Paiva Junqueira, Dushyanth Narayanan |
USENIX ATC | 3 |
| 2015 | Extensible distributed coordinationabstractMost 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 |
EuroSys | 5 |
| 2015 | Visigoth fault toleranceabstractWe 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 |
EuroSys | 6 |
| 2015 | Scalable Error Isolation for Distributed Systems
Diogo Behrens, Marco Serafini, Flavio Paiva Junqueira, Sergei Arnautov, Christof Fetzer |
NSDI | 3 |
| 2014 | Omid: Lock-free transactional support for distributed data storesabstractIn 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 |
ICDE | 2 |
| 2014 | Slider: incremental sliding window analyticsabstractSliding 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 |
Middleware | 3 |
| 2013 | Cache refreshing for online social news feedsabstractSeveral 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 |
CIKM | 2 |
| 2013 | Exploiting user clicks for automatic seed set generation for entity matchingabstractMatching 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 |
KDD | 2 |
| 2013 | DynaSoRe: Efficient In-Memory Store for Social Applications
Xiao Bai 0002, Arnaud Jégou, Flavio Paiva Junqueira, Vincent Leroy 0001 |
Middleware | 3 |
| 2013 | On Barriers and the Gap between Active and Passive Replication
Flavio Paiva Junqueira, Marco Serafini |
DISC | 1 |
| 2013 | Piggybacking on Social NetworksabstractThe 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 replicationabstractDeferred 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 |
DSN | 3 |
| 2012 | On the (limited) power of non-equivocationabstractIn 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 |
PODC | 2 |
| 2012 | Online result cache invalidation for real-time web searchabstractCaches 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 |
SIGIR | 2 |
| 2012 | Reactive index replication for distributed search enginesabstractDistributed 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 |
SIGIR | 1 |
| 2012 | Practical Hardening of Crash-Tolerant Systems
Miguel Correia 0001, Daniel Gómez Ferro, Flavio Paiva Junqueira, Marco Serafini |
USENIX ATC | 3 |
| 2012 | Dynamic Reconfiguration of Primary/Backup Clusters
Alexander Shraer, Benjamin C. Reed, Dahlia Malkhi, Flavio Paiva Junqueira |
USENIX ATC | 4 |
| 2012 | Brief Announcement: Consensus and Efficient Passive Replication
Flavio Paiva Junqueira, Marco Serafini |
DISC | 1 |
| 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 feedbackabstractSearch 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 |
CIKM | 3 |
| 2011 | Assigning documents to master sites in distributed searchabstractAn 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 |
CIKM | 3 |
| 2011 | Zab: High-performance broadcast for primary-backup systemsabstractZab 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 |
DSN | 1 |
| 2011 | Leader Election for Replicated Services Using Application Scores
Diogo Becker, Flavio Paiva Junqueira, Marco Serafini |
Middleware | 2 |
| 2011 | Energy-price-driven query processing in multi-center web search enginesabstractConcurrently 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 |
SIGIR | 4 |
| 2011 | Posting list intersection on multicore architecturesabstractIn 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 |
SIGIR | 3 |
| 2011 | Document assignment in multi-site search enginesabstractAssigning 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 |
WSDM | 3 |
| 2010 | Caching search engine results over incremental indicesabstractA 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 |
SIGIR | 3 |
| 2010 | ZooKeeper: Wait-free Coordination for Internet-scale Systems
Patrick Hunt, Mahadev Konar, Flavio Paiva Junqueira, Benjamin C. Reed |
USENIX ATC | 3 |
| 2010 | Caching search engine results over incremental indicesabstractA 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 |
WWW | 3 |
| 2010 | A refreshing perspective of search engine cachingabstractCommercial 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 |
WWW | 2 |
| 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 |
CIDR | 4 |
| 2009 | On the feasibility of multi-site web search enginesabstractWeb 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 |
CIKM | 3 |
| 2009 | Introduction
Dejan Kostic, Guillaume Pierre, Flavio Paiva Junqueira, Peter R. Pietzuch |
Euro-Par | 3 |
| 2009 | The life and times of a zookeeperabstractA 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 |
PODC | 1 |
| 2009 | On efficient posting list intersection with multicore processorsabstractNo abstract available. Shirish Tatikonda, Flavio Paiva Junqueira, Berkant Barla Cambazoglu, Vassilis Plachouras |
SIGIR | 2 |
| 2009 | The life and times of a zookeeperabstractA 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 |
SPAA | 1 |
| 2009 | Brief Announcement Zab: A Practical Totally Ordered Broadcast Protocol
Flavio Paiva Junqueira, Benjamin C. Reed |
DISC | 1 |
| 2008 | Mencius: Building Efficient Replicated State Machine for WANs
Yanhua Mao, Flavio Paiva Junqueira, Keith Marzullo |
OSDI | 2 |
| 2008 | ResIn: a combination of results caching and index pruning for high-performance web search enginesabstractResults 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 |
SIGIR | 2 |
| 2008 | Optimizing Threshold Protocols in Adversarial Structures
Maurice Herlihy, Flavio Paiva Junqueira, Keith Marzullo, Lucia Draque Penso |
DISC | 2 |
| 2008 | Design trade-offs for search engine cachingabstractIn 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. Web | 3 |
| 2007 | Challenges on Distributed Web RetrievalabstractIn 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 |
ICDE | 3 |
| 2007 | The impact of caching on search enginesabstractIn 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 |
SIGIR | 3 |
| 2007 | Admission Policies for Caches of Search Engine Results
Ricardo Baeza-Yates, Flavio Paiva Junqueira, Vassilis Plachouras, Hans Friedrich Witschel |
SPIRE | 2 |
| 2007 | A framework for the design of dependent-failure algorithmsabstractAbstract 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 EnvironmentsabstractReplication 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 |
HPDC | 2 |
| 2005 | Replication Predicates for Dependent-Failure Algorithms
Flavio Paiva Junqueira, Keith Marzullo |
Euro-Par | 1 |
| 2005 | Surviving Internet Catastrophes
Flavio Paiva Junqueira, Ranjita Bhagwan, Alejandro Hevia, Keith Marzullo, Geoffrey M. Voelker |
USENIX ATC, General Track | 1 |
| 2005 | Coterie Availability in Sites
Flavio Paiva Junqueira, Keith Marzullo |
DISC | 1 |
| 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 |
HotOS | 1 |
| 2003 | Synchronous Consensus for Dependent Process FailureabstractWe 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 |
ICDCS | 1 |