EDBT 2026 Demo / reviewers in the wild / expert
Alan L. Cox
dblp:c/AlanLCox
· DBLP profile ↗
73ranked-venue papers
4as first author
3since 2021 · last 2023
0009-0005-4904-9600ORCID · verified
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 44 · 3 first-author · 3 since 2021Software engineering, systems software and programming languages · 21 · 3 first-authorComputer networks · 11Databases, data management, data science and information retrieval · 4Security and privacy · 2Applied, interdisciplinary, general and emerging computing · 2
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2023 | The Impact of Page Size and Microarchitecture on Instruction Address Translation OverheadabstractAs the volume of data processed by applications has increased, considerable attention has been paid to data address translation overheads, leading to the widespread use of larger page sizes (“superpages”) and multi-level translation lookaside buffers (TLBs). However, far less attention has been paid to instruction address translation and its relation to TLB and pipeline structure. In prior work, we quantified the impact of using code superpages on a variety of widely used applications, ranging from compilers to web user-interface frameworks, and the impact of sharing page table pages for executables and shared libraries. Within this article, we augment those results by first uncovering the effects that microarchitectural differences between Intel Skylake and AMD Zen+, particularly their different TLB organizations, have on instruction address translation overhead. This analysis provides some key insights into the microarchitectural design decisions that impact the cost of instruction address translation. First, a lower-level (level 2) TLB that has both instruction and data mappings competing for space within the same structure allows better overall performance and utilization when using code superpages. Code superpages not only reduce instruction address translation overhead but also indirectly reduce data address translation overhead. In fact, for a few applications, the use of just a few code superpages has a larger impact on overall performance than the use of a much larger number of data superpages. Second, a level 1 (L1) TLB with separate structures for different page sizes may require careful tuning of the superpage promotion policy for code, and a correspondingly suboptimal utilization of the level 2 TLB. In particular, increasing the number of superpages when the size of the L1 superpage structure is small may result in more L1 TLB misses for some applications. Moreover, on some microarchitectures, the cost of these misses can be highly variable, because replacement is delayed until all of the in-flight instructions mapped by the victim entry are retired. Hence, more superpage promotions can result in a performance regression. Finally, our findings also make a case for first-class OS support for superpages on ordinary files containing executables and shared libraries, as well as a more aggressive superpage policy for code. Alan L. Cox, Sandhya Dwarkadas, Xiaowan Dong |
ACM Trans. Archit. Code Optim. | 2 |
| 2023 | An FPGA Accelerator for Genome Variant CallingabstractIn genome analysis, it is often important to identify variants from a reference genome. However, identifying variants that occur with low frequency can be challenging, as it is computationally intensive to do so accurately. LoFreq is a widely used program that is adept at identifying low-frequency variants. This article presents a design framework for an FPGA-based accelerator for LoFreq. In particular, this accelerator is targeted at virus analysis, which is particularly challenging, compared to human genome analysis, as the characteristics of the data to be analyzed are fundamentally different. Across the design space, this accelerator can achieve up to 120× speedups on the core computation of LoFreq and speedups of up to 51.7× across the entire program. Tiancheng Xu, Scott Rixner, Alan L. Cox |
ACM Trans. Reconfigurable Technol. Syst. | 3 |
| 2022 | An FPGA Accelerator for Genome Variant CallingabstractIn genome analysis, it is often important to identify variants from a reference genome. However, identifying variants that occur with low frequency can be challenging, as it is computationally intensive to do so accurately. LoFreq is a widely used program that is adept at identifying low frequency variants. This paper presents an FPGA-based accelerator for LoFreq. In particular, this accelerator is targeted at virus analysis, which is particularly challenging, compared to human genome analysis, as the characteristics of the data to be analyzed are fundamentally different. This accelerator can achieve up to 120× speedups on the core computation of LoFreq and speedups of up to 32.4× across the entire program. Tiancheng Xu, Scott Rixner, Alan L. Cox |
FCCM | 3 |
| 2020 | A Comprehensive Analysis of Superpage Management Mechanisms and Policies
Weixi Zhu, Alan L. Cox, Scott Rixner |
USENIX ATC | 2 |
| 2019 | On the Impact of Instruction Address Translation OverheadabstractEven on modern processors with their ever larger instruction translation lookaside buffers (TLBs), we find that a variety of widely used applications, ranging from compilers to web user-interface frameworks, suffer from high instruction address translation overheads. In this paper, we explore the efficacy of different operating system-level approaches to automatically reducing this instruction address translation overhead. Specifically, we evaluate the use of automatic superpage promotion and page table sharing as well as a transparent padding mechanism that enables small code regions to be mapped using superpages. Overall, we find that the combined effects of these different approaches can reduce an application's total execution cycles by up to 18%. Surprisingly, we find that improving address translation performance in the first-level instruction TLB can significantly reduce the address translation overhead for data accesses. The overall reduction in execution cycles is more than double the instruction address translation overhead on stock FreeBSD, demonstrating that data address translation and access synergistically benefit from less contention in the caches and TLBs that might be shared across instruction and data. Xiaowan Dong, Alan L. Cox, Sandhya Dwarkadas |
ISPASS | 3 |
| 2018 | Shielding Software From Privileged Side-Channel Attacks
Xiaowan Dong, Zhuojia Shen, John Criswell, Alan L. Cox, Sandhya Dwarkadas |
USENIX Security Symposium | 4 |
| 2016 | TPC: Target-Driven Parallelism Combining Prediction and Correction to Reduce Tail Latency in Interactive ServicesabstractIn interactive services such as web search, recommendations, games and finance, reducing the tail latency is crucial to provide fast response to every user. Using web search as a driving example, we systematically characterize interactive workload to identify the opportunities and challenges for reducing tail latency. We find that the workload consists of mainly short requests that do not benefit from parallelism, and a few long requests which significantly impact the tail but exhibit high parallelism speedup. This motivates estimating request execution time, using a predictor, to identify long requests and to parallelize them. Prediction, however, is not perfect; a long request mispredicted as short is likely to contribute to the server tail latency, setting a ceiling on the achievable tail latency. We propose TPC, an approach that combines prediction information judiciously with dynamic correction for inaccurate prediction. Dynamic correction increases parallelism to accelerate a long request that is mispredicted as short. TPC carefully selects the appropriate target latencies based on system load and parallelism efficiency to reduce tail latency. Myeongjae Jeon, Yuxiong He, Hwanju Kim, Sameh Elnikety, Scott Rixner, Alan L. Cox |
ASPLOS | 6 |
| 2016 | Shared address translation revisitedabstractModern operating systems avoid duplication of code and data when they are mapped by multiple processes by sharing physical memory through mechanisms like copy-on-write. Nonetheless, a separate copy of the virtual address translation structures, such as page tables, are still maintained for each process, even if they are identical. This duplication can lead to inefficiencies in the address translation process and interference within the memory hierarchy. In this paper, we show that on Android platforms, sharing address translation structures, specifically, page tables and TLB entries, for shared libraries can improve performance. For example, at a low level, sharing address translation structures reduces the cost of fork by more than half by reducing page table construction overheads. At a higher level, application launch and IPC are faster due to page fault elimination coupled with better cache and TLB performance when context switching. Xiaowan Dong, Sandhya Dwarkadas, Alan L. Cox |
EuroSys | 3 |
| 2016 | Deadlock-free local fast failover for arbitrary data center networksabstractToday, given data center networks' sizes and bursty workloads, it is likely that at any moment there is packet loss due to some type of failure in the network. This paper focuses on solving the two most common types of data center network failures: congestion and routing failures. Recently, there has been demand for lossless Ethernet (DCB) in data center networks as a solution to congestion failures. However, DCB complicates fault tolerance by introducing a new type of failure, deadlock. If DCB is enabled, then all routing must be deadlock free. To the best of our knowledge, this paper describes the first ever deadlock-free approaches to local fast failover that can be combined with DCB, DF-FI and DF-EDST resilience. Moreover, in the evaluation, this paper shows that DF-EDST resilience, which is the paper's main contribution, can improve fault tolerance without adversely impacting performance when compared to a state-of-the-art approach to deadlock-free routing. If, however, a small reduction in aggregate throughput is acceptable, then it is possible to build routes such that only 0.00001% of the total flows in the network are likely to fail given 16 edge failures on networks with 1K-4K hosts. Brent E. Stephens, Alan L. Cox |
INFOCOM | 2 |
| 2015 | GD-Wheel: a cost-aware replacement policy for key-value storesabstractMemory-based key-value stores, such as Memcached and Redis, are often used to speed up web applications. Specifically, they are used to cache the results of computations, such as database queries and dynamically generated web pages, so that a future request to the web application may not have to repeat the same computation. Currently, when memory-based key-value stores reach their capacity limits, they use replacement policies, like LRU and random, that are oblivious to differences among the cached results in their recomputation costs. However, this paper shows that if the costs of recomputing cached results vary significantly, as in the RUBiS and TPC-W benchmarks, then a cost-aware replacement policy will not only reduce the web application's total recomputation cost but also reduce its average response time. Conglong Li, Alan L. Cox |
EuroSys | 2 |
| 2014 | Practical DCB for improved data center networksabstractStorage area networking is driving commodity data center switches to support lossless Ethernet (DCB). Unfortunately, to enable DCB for all traffic on arbitrary network topologies, we must address several problems that can arise in lossless networks, e.g., large buffering delays, unfairness, head of line blocking, and deadlock. We propose TCP-Bolt, a TCP variant that not only addresses the first three problems but reduces flow completion times by as much as 70%. We also introduce a simple, practical deadlock-free routing scheme that eliminates deadlock while achieving aggregate network throughput within 15% of ECMP routing. This small compromise in potential routing capacity is well worth the gains in flow completion time. We note that our results on deadlock-free routing are also of independent interest to the storage area networking community. Further, as our hardware testbed illustrates, these gains are achievable today, without hardware changes to switches or NICs. Brent E. Stephens, Alan L. Cox, Ankit Singla, John B. Carter, Colin Dixon, Wes Felter |
INFOCOM | 2 |
| 2014 | Predictive parallelization: taming tail latencies in web searchabstractWeb search engines are optimized to reduce the high-percentile response time to consistently provide fast responses to almost all user queries. This is a challenging task because the query workload exhibits large variability, consisting of many short-running queries and a few long-running queries that significantly impact the high-percentile response time. With modern multicore servers, parallelizing the processing of an individual query is a promising solution to reduce query execution time, but it gives limited benefits compared to sequential execution since most queries see little or no speedup when parallelized. The root of this problem is that short-running queries, which dominate the workload, do not benefit from parallelization. They incur a large parallelization overhead, taking scarce resources from long-running queries. On the other hand, parallelization substantially reduces the execution time of long-running queries with low overhead and high parallelization efficiency. Motivated by these observations, we propose a predictive parallelization framework with two parts: (1) predicting long-running queries, and (2) selectively parallelizing them. For the first part, prediction should be accurate and efficient. For accuracy, we study a comprehensive feature set covering both term features (reflecting dynamic pruning efficiency) and query features (reflecting query complexity). For efficiency, to keep overhead low, we avoid expensive features that have excessive requirements such as large memory footprints. For the second part, we use the predicted query execution time to parallelize long-running queries and process short-running queries sequentially. We implement and evaluate the predictive parallelization framework in Microsoft Bing search. Our measurements show that under moderate to heavy load, the predictive strategy reduces the 99th-percentile response time by 50% (from 200 ms to 100 ms) compared with prior approaches that parallelize all queries. Myeongjae Jeon, Saehoon Kim, Seung-won Hwang, Yuxiong He, Sameh Elnikety, Alan L. Cox, Scott Rixner |
SIGIR | 6 |
| 2013 | Adaptive parallelism for web searchabstractA web search query made to Microsoft Bing is currently parallelized by distributing the query processing across many servers. Within each of these servers, the query is, however, processed sequentially. Although each server may be processing multiple queries concurrently, with modern multicore servers, parallelizing the processing of an individual query within the server may nonetheless improve the user's experience by reducing the response time. In this paper, we describe the issues that make the parallelization of an individual query within a server challenging, and we present a parallelization approach that effectively addresses these challenges. Since each server may be processing multiple queries concurrently, we also present a adaptive resource management algorithm that chooses the degree of parallelism at run-time for each query, taking into account system load and parallelization efficiency. As a result, the servers now execute queries with a high degree of parallelism at low loads, gracefully reduce the degree of parallelism with increased load, and choose sequential execution under high load. We have implemented our parallelization approach and adaptive resource management algorithm in Bing servers and evaluated them experimentally with production workloads. The experimental results show that the mean and 95th-percentile response times for queries are reduced by more than 50% under light or moderate load. Moreover, under high load where parallelization adversely degrades the system performance, the response times are kept the same as when queries are executed sequentially. In all cases, we observe no degradation in the relevance of the search results. Myeongjae Jeon, Yuxiong He, Sameh Elnikety, Alan L. Cox, Scott Rixner |
EuroSys | 4 |
| 2013 | Plinko: building provably resilient forwarding tablesabstractThis paper introduces Plinko, a network architecture that uses a novel forwarding model and routing algorithm to build networks with forwarding paths that, assuming arbitrarily large forwarding tables, are provably resilient against t link failures, ∀t ∈ N. However, in practice, there are clearly limits on the size of forwarding tables. Nonetheless, when constrained to hardware comparable to modern top-of-rack (TOR) switches, Plinko scales with high resilience to networks with up to ten thousand hosts. Thus, as long as t or fewer links have failed, the only reason packets of any flow in a Plinko network will be dropped are congestion, packet corruption, and a partitioning of the network topology, and, even after t + 1 failures, most, if not all, flows may be unaffected. In addition, Plinko is topology independent, supports arbitrary paths for routing, provably bounds stretch, and does not require any additional computation during forwarding. To the best of our knowledge, Plinko is the first network to have all of these properties. Brent E. Stephens, Alan L. Cox, Scott Rixner |
HotNets | 2 |
| 2013 | Hyper-Switch: A Scalable Software Virtual Switching Architecture
Kaushik Kumar Ram, Alan L. Cox, Mehul Chadha, Scott Rixner |
USENIX ATC | 2 |
| 2013 | Reducing DRAM row activations with eager read/write clusteringabstractThis article describes and evaluates a new approach to optimizing DRAM performance and energy consumption that is based on eagerly writing dirty cache lines to DRAM. Under this approach, many dirty cache lines are written to DRAM before they are evicted. In particular, dirty cache lines that have not been recently accessed are eagerly written to DRAM when the corresponding row has been activated by an ordinary, noneager access, such as a read. This approach enables clustering of reads and writes that target the same row, resulting in a significant reduction in row activations. Specifically, for a variety of applications, it reduces the number of DRAM row activations by an average of 42% and a maximum of 82%. Moreover, the results from a full-system simulator show compelling performance improvements and energy consumption reductions. Out of 23 applications, 6 have overall performance improvements between 10% and 20%, and 3 have improvements in excess of 20%. Furthermore, 12 consume between 10% and 20% less DRAM energy, and 7 have energy consumption reductions in excess of 20%. Myeongjae Jeon, Conglong Li, Alan L. Cox, Scott Rixner |
ACM Trans. Archit. Code Optim. | 3 |
| 2012 | PAST: scalable ethernet for data centersabstractWe present PAST, a novel network architecture for data center Ethernet networks that implements a Per-Address Spanning Tree routing algorithm. PAST preserves Ethernet's self-configuration and mobility support while increasing its scalability and usable bandwidth. PAST is explicitly designed to accommodate unmodified commodity hosts and Ethernet switch chips. Surprisingly, we find that PAST can achieve performance comparable to or greater than Equal-Cost Multipath (ECMP) forwarding, which is currently limited to layer-3 IP networks, without any multipath hardware support. In other words, the hardware and firmware changes proposed by emerging standards like TRILL are not required for high-performance, scalable Ethernet networks. We evaluate PAST on Fat Tree, HyperX, and Jellyfish topologies, and show that it is able to capitalize on the advantages each offers. We also describe an OpenFlow-based implementation of PAST in detail. Brent E. Stephens, Alan L. Cox, Wes Felter, Colin Dixon, John B. Carter |
CoNEXT | 2 |
| 2011 | A Scalability Study of Enterprise Network ArchitecturesabstractThe largest enterprise networks already contain hundreds of thousands of hosts. Enterprise networks are composed of Ethernet subnets interconnected by IP routers. These routers require expensive configuration and maintenance. If the Ethernet subnets are made more scalable, the high cost of the IP routers can be eliminated. Unfortunately, it has been widely acknowledged that Ethernet does not scale well because it relies on broadcast, which wastes bandwidth, and a cycle-free topology, which poorly distributes load and forwarding state. There are many recent proposals to replace Ethernet, each with its own set of architectural mechanisms. These mechanisms include eliminating broadcasts, using source routing, and restricting routing paths. Although there are many different proposed designs, there is little data available that allows for comparisons between designs. This study performs simulations to evaluate all of the factors that affect the scalability of Ethernet together, which has not been done in any of the proposals. The simulations demonstrate that, in a realistic environment, source routing reduces the maximum state requirements of the network by over an order of magnitude. About the same level of traffic engineering achieved by load-balancing all the flows at the TCP/UDP flow granularity is possible by routing only the heavy flows at the TCP/UDP granularity. Additionally, requiring routing restrictions, such as deadlock-freedom or minimum-hop routing, can significantly reduce the network's ability to perform traffic engineering across the links. Brent E. Stephens, Alan L. Cox, Scott Rixner, T. S. Eugene Ng |
ANCS | 2 |
| 2011 | SpecTLB: a mechanism for speculative address translationabstractData-intensive computing applications are using more and more memory and are placing an increasing load on the virtual memory system. While the use of large pages can help alleviate the overhead of address translation, they limit the control the operating system has over memory allocation and protection. We present a novel device, the SpecTLB, that exploits the predictable behavior of reservation-based physical memory allocators to interpolate address translations. Thomas W. Barr, Alan L. Cox, Scott Rixner |
ISCA | 2 |
| 2010 | sNICh: efficient last hop networking in the data centerabstractVirtualization has fundamentally changed the data center network. The last hop of the network is no longer handled by a physical network switch, but rather is typically performed in software inside the server to switch among virtual machines hosted by that server. Kaushik Kumar Ram, Jayaram Mudigonda, Alan L. Cox, Scott Rixner, Parthasarathy Ranganathan, Jose Renato Santos |
ANCS | 3 |
| 2010 | Axon: a flexible substrate for source-routed ethernetabstractThis paper introduces the Axon, an Ethernet-compatible device for creating large-scale datacenter networks. Axons are inexpensive, practical devices that are demonstrated using prototype hardware. Functionally, Axons replace Ethernet switches and maintain full compatibility with existing Ethernet hosts. Between themselves, however, Axons transparently use source-routed Ethernet. This unlocks many benefits, such as improved network scalability, performance, and flexibility.In an Axon network, all state required to route a host's packets is placed in the local Axon---the Axon to which the host is directly connected. Therefore, regardless of the scale of the network, the route computation and storage needs of a single Axon device only need to scale with the demands of its locally-connected hosts. This is in stark contrast to conventional switched Ethernet, which requires routing resources proportional to the traffic that flows through the device. Scalability is also increased by eliminating the use of packet flooding for automatic location and address discovery. Further, source-routed Ethernet increases network flexibility by supporting different route selection strategies. For example, shortest-path routing could be employed, or longer paths selected to minimize congestion by balancing traffic across redundant links. Jeffrey Shafer, Brent E. Stephens, Michael Foss, Scott Rixner, Alan L. Cox |
ANCS | 5 |
| 2010 | CONTRACT: Incorporating Coordination into the IP Network Control PlaneabstractThis paper presents the CONTRACT framework to address a fundamental deficiency of the IP network control plane, namely the lack of coordination between an IGP and other control functions involved in achieving a high level objective. For example, an IGP's default automatic reaction to a network failure may result in an SLA violation, even if the IGP link weights have been carefully chosen. This is because an IGP blindly routes traffic along the shortest paths based on link weights, and it is completely oblivious to the interactions between SLA compliance, load balancing and traffic policing objectives in a network. The CONTRACT framework makes it possible to coordinate these objectives. Under this framework, routers continue to operate autonomously, but they also coordinate their actions with a centralized network controller, which evaluates the impact of routing changes, decides whether the changes are SLA compliant, and performs load rebalancing and/or packet filter reconfiguration as necessary. The key contribution of CONTRACT is a set of coordination algorithms. We show that CONTRACT can effectively coordinate the actions of routing, load balancing and traffic policing to improve a network's SLA compliance. Zheng Cai, Florin Dinu, Alan L. Cox, T. S. Eugene Ng |
ICDCS | 4 |
| 2010 | Translation caching: skip, don't walk (the page table)abstractThis paper explores the design space of MMU caches that accelerate virtual-to-physical address translation in processor architectures, such as x86-64, that use a radix tree page table. In particular, these caches accelerate the page table walk that occurs after a miss in the Translation Lookaside Buffer. This paper shows that the most effective MMU caches are translation caches, which store partial translations and allow the page walk hardware to skip one or more levels of the page table. Thomas W. Barr, Alan L. Cox, Scott Rixner |
ISCA | 2 |
| 2010 | The Hadoop distributed filesystem: Balancing portability and performanceabstractHadoop is a popular open-source implementation of MapReduce for the analysis of large datasets. To manage storage resources across the cluster, Hadoop uses a distributed user-level filesystem. This filesystem - HDFS - is written in Java and designed for portability across heterogeneous hardware and software platforms. This paper analyzes the performance of HDFS and uncovers several performance issues. First, architectural bottlenecks exist in the Hadoop implementation that result in inefficient HDFS usage due to delays in scheduling new MapReduce tasks. Second, portability limitations prevent the Java implementation from exploiting features of the native platform. Third, HDFS implicitly makes portability assumptions about how the native platform manages storage resources, even though native filesystems and I/O schedulers vary widely in design and behavior. This paper investigates the root causes of these performance bottlenecks in order to evaluate tradeoffs between portability and performance in the Hadoop distributed filesystem. Jeffrey Shafer, Scott Rixner, Alan L. Cox |
ISPASS | 3 |
| 2009 | EtherProxy: Scaling Ethernet By Suppressing Broadcast TrafficabstractEthernet is the dominant technology for local area networks. This is mainly because of its autoconfiguration capability and its cost effectiveness. Unfortunately, a single Ethernet network can not scale to span a large enterprise network. A main reason for this is broadcast traffic resulting from many protocols running on top of Ethernet. This paper addresses Ethernet's scalability limits due to broadcast traffic. We studied and characterized broadcast traffic in Ethernet networks using traces collected from real networks. We found that broadcast is mainly used in Ethernet for service and resource discovery. For example, the address resolution protocol (ARP) uses broadcast to discover a MAC address that corresponds to an IP address. To avoid broadcast for service and resource discovery, we propose a new device, the EtherProxy. An EtherProxy uses caching to suppress broadcast traffic. EtherProxy is backward compatible and requires no changes to existing hardware, software, or protocols. Moreover, it requires no configuration. In our evaluation, we used real and synthetic workloads. Using both workloads, we experimentally demonstrate the effectiveness of the EtherProxy. Khaled Elmeleegy, Alan L. Cox |
INFOCOM | 2 |
| 2009 | Achieving 10 Gb/s using safe and transparent network interface virtualizationabstractThis paper presents mechanisms and optimizations to reduce the overhead of network interface virtualization when using the driver domain I/O virtualization model. The driver domain model provides benefits such as support for legacy device drivers and fault isolation. However, the processing overheads incurred in the driver domain to achieve these benefits limit overall I/O performance. This paper demonstrates the effectiveness of two approaches to reduce driver domain overheads. First, Xen is modified to support multi-queue network interfaces to eliminate the software overheads of packet demultiplexing and copying. Second, a grant reuse mechanism is developed to reduce memory protection overheads. These mechanisms shift the bottleneck from the driver domain to the guest domains, improving scalability and enabling significantly higher data rates. This paper also presents and evaluates a series of optimizations that substantially reduce the I/O virtualization overheads in the guest domain. In combination, these mechanisms and optimizations increase the maximum throughput achieved by guest domains from 2.9Gb/s to full 10 Gigabit Ethernet link rates. Kaushik Kumar Ram, Jose Renato Santos, Yoshio Turner, Alan L. Cox, Scott Rixner |
VEE | 4 |
| 2009 | Understanding and mitigating the effects of count to infinity in Ethernet networks
Khaled Elmeleegy, Alan L. Cox, T. S. Eugene Ng |
IEEE/ACM Trans. Netw. | 2 |
| 2008 | Investigating the TLB Behavior of High-end Scientific Applications on Commodity MicroprocessorsabstractThe floating point portion of the SPEC CPU suite and the HPC Challenge suite are widely recognized and utilized as benchmarks that represent scientific application behavior. In this work we show that while these benchmark suites may be representative of the cache behavior of production scientific applications, they do not accurately represent the TLB behavior of these applications. Furthermore, we demonstrate that the difference can have a significant impact on performance. In the first part of the paper we present results from implementation-independent trace-based simulations which demonstrate that benchmarks exhibit significantly different TLB behavior for a range of page sizes than a representative set of production applications. In the second part we validate these results on the AMD Opteron implementation of the x86 architecture, showing that false conclusions about choice of page size, drawn from benchmark performance, can result in performance degradations of up to nearly 50% for the production applications we investigated.. Collin McCurdy, Alan L. Cox, Jeffrey S. Vetter |
ISPASS | 2 |
| 2008 | Explaining the Impact of Network Transport Protocols on SIP Proxy PerformanceabstractThis paper characterizes the impact that the use of UDP versus TCP has on the performance and scalability of the OpenSER SIP proxy server. The session initiation protocol (SIP) is an application-layer signaling protocol that is widely used for establishing voice-over-IP (VoIP) phone calls. SIP can utilize a variety of transport protocols, including UDP and TCP. Despite the advantages of TCP, such as reliable delivery and congestion control, the common practice is to use UDP. This is a result of the belief that UDP's lower processor and network overhead results in improved performance and scalability of SIP services. This paper argues against this conventional wisdom. This paper shows that the principal reasons for OpenSER's poor performance using TCP are caused by the server's design, and not the low-level performance of UDP versus TCP. Specifically, OpenSER's architecture for handling concurrent calls is responsible for most of the difference. Moreover, once these issues are addressed, OpenSER's performance using TCP is much more competitive with its performance using UDP. Kaushik Kumar Ram, Ian C. Fedeli, Alan L. Cox, Scott Rixner |
ISPASS | 3 |
| 2008 | Protection Strategies for Direct Access to Virtualized I/O Devices
Paul Willmann, Scott Rixner, Alan L. Cox |
USENIX ATC | 3 |
| 2008 | Scheduling I/O in virtual machine monitorsabstractThis paper explores the relationship between domain scheduling in avirtual machine monitor (VMM) and I/O performance. Traditionally, VMM schedulers have focused on fairly sharing the processor resources among domains while leaving the scheduling of I/O resources as asecondary concern. However, this can resultin poor and/or unpredictable application performance, making virtualization less desirable for applications that require efficient and consistent I/O behavior. Diego Ongaro, Alan L. Cox, Scott Rixner |
VEE | 2 |
| 2007 | Whodunit: transactional profiling for multi-tier applicationsabstractThis paper is concerned with performance debugging of multi-tier applications, such as commonly found in servers and dynamic-content web sites. Existing tools and techniques for profiling such applications are not general enough to track and profile transactions in a generic multi-tier application. We propose transactional profiling that provides a general solution to this problem. We provide novel algorithms and techniques to track and profile transactions that flow through shared memory, events, stages or via interprocess communication using messages. We also measure interference among concurrent transactions. Anupam Chanda, Alan L. Cox, Willy Zwaenepoel |
EuroSys | 2 |
| 2007 | Concurrent Direct Network Access for Virtual Machine MonitorsabstractThis paper presents hardware and software mechanisms to enable concurrent direct network access (CDNA) by operating systems running within a virtual machine monitor. In a conventional virtual machine monitor, each operating system running within a virtual machine must access the network through a software-virtualized network interface. These virtual network interfaces are multiplexed in software onto a physical network interface, incurring significant performance overheads. The CDNA architecture improves networking efficiency and performance by dividing the tasks of traffic multiplexing, interrupt delivery, and memory protection between hardware and software in a novel way. The virtual machine monitor delivers interrupts and provides protection between virtual machines, while the network interface performs multiplexing of the network data. In effect, the CDNA architecture provides the abstraction that each virtual machine is connected directly to its own network interface. Through the use of CDNA, many of the bottlenecks imposed by software multiplexing can be eliminated without sacrificing protection, producing substantial efficiency improvements Jeffrey Shafer, David Carr, Aravind Menon, Scott Rixner, Alan L. Cox, Willy Zwaenepoel, Paul Willmann |
HPCA | 5 |
| 2007 | Etherfuse: an ethernet watchdogabstractEthernet is pervasive. This is due in part to its ease of use. Equipment can be added to an Ethernet network with little or no manual configuration. Furthermore, Ethernet is self-healing in the event of equipment failure or removal. However, there are scenarios where a local event can lead to network-wide packet loss and duplication due to slow or faulty reconfiguration of the spanning tree. Moreover, in some cases the packet loss and duplication may persist indefinitely. Khaled Elmeleegy, Alan L. Cox, T. S. Eugene Ng |
SIGCOMM | 2 |
| 2006 | On Count-to-Infinity Induced Forwarding Loops Ethernet NetworksabstractEthernet's high performance, low cost and ubiquity have made it the dominant networking technology for many application domains. Unfortunately, its distributed forwarding topology computation protocol - the Rapid Spanning Tree Proto- col (RSTP) - can suffer from a classic count-to-infinity problem that may lead to a forwarding loop under certain network failures. The consequences are serious. During the period of count-to-infinity, which can last tens of seconds even in a small network, the network can become highly congested by packets that persist in cycles in the network, even packet forwarding can fail as the forwarding tables are polluted. In this paper, we explain the origin of this problem in detail and study its behavior. We find that simply tuning RSTP's parameter settings cannot adequately address the fundamental problem with count-to- infinity. We propose a simple and effective solution called RSTP with Epochs. This approach uses epochs of sequence numbers in protocol messages to eliminate stale protocol information in the network and allows the forwarding topology to recover in merely one round-trip time across the network. Khaled Elmeleegy, Alan L. Cox, T. S. Eugene Ng |
INFOCOM | 2 |
| 2006 | Caching Dynamic Web Content: Designing and Analysing an Aspect-Oriented Solution
Sara Bouchenak, Alan L. Cox, Steven G. Dropsho, Sumit Mittal, Willy Zwaenepoel |
Middleware | 2 |
| 2006 | Optimizing Network Virtualization in Xen (awarded best paper)
Aravind Menon, Alan L. Cox, Willy Zwaenepoel |
USENIX ATC, General Track | 2 |
| 2006 | An Evaluation of Network Stack Parallelization Strategies in Modern Operating Systems
Paul Willmann, Scott Rixner, Alan L. Cox |
USENIX ATC, General Track | 3 |
| 2005 | Causeway: Operating System Support for Controlling and Analyzing the Execution of Distributed Programs
Anupam Chanda, Khaled Elmeleegy, Alan L. Cox, Willy Zwaenepoel |
HotOS | 3 |
| 2005 | A Comparative Evaluation of Transparent Scaling Techniques for Dynamic Content ServersabstractWe study several transparent techniques for scaling dynamic content Web sites, and we evaluate their relative impact when used in combination. Full transparency implies strong data consistency as perceived by the user, no modifications to existing dynamic content site tiers and no additional programming effort from the user or site administrator upon deployment. We study strategies for scheduling and load balancing queries on a cluster of replicated database back-ends. We also investigate transparent query caching as a means of enhancing database replication. Our work shows that, on an experimental platform with up to 8 database replicas, the various techniques work in synergy to improve overall scaling for the e-commerce TPC-W benchmark. We rank the techniques necessary for high performance in order of impact as follows. Key among the strategies are scheduling strategies, such as conflict-aware scheduling, that minimize consistency maintenance overheads. The choice of load balancing strategy is less important. Transparent query result caching increases performance significantly at any given cluster size for a mostly-read workload. Its benefits are limited for write-intensive workloads, where content-aware scheduling is the only scaling option. Cristiana Amza, Alan L. Cox, Willy Zwaenepoel |
ICDE | 2 |
| 2005 | Causeway: Support for Controlling and Analyzing the Execution of Multi-tier Applications
Anupam Chanda, Khaled Elmeleegy, Alan L. Cox, Willy Zwaenepoel |
Middleware | 3 |
| 2005 | A Portable Kernel Abstraction for Low-Overhead Ephemeral Mapping Management
Khaled Elmeleegy, Anupam Chanda, Alan L. Cox, Willy Zwaenepoel |
USENIX ATC, General Track | 3 |
| 2004 | Lazy Asynchronous I/O for Event-Driven Servers
Khaled Elmeleegy, Anupam Chanda, Alan L. Cox, Willy Zwaenepoel |
USENIX ATC, General Track | 3 |
| 2003 | Using Performance Reflection in Systems Software
Robert J. Fowler, Alan L. Cox, Sameh Elnikety, Willy Zwaenepoel |
HotOS | 2 |
| 2003 | Distributed Versioning: Consistent Replication for Scaling Back-End Databases of Dynamic Content Web Sites
Cristiana Amza, Alan L. Cox, Willy Zwaenepoel |
Middleware | 2 |
| 2003 | Run-time support for distributed sharing in safe languagesabstractWe present a new run-time system that supports object sharing in a distributed system. The key insight in this system is that a handle-based implementation of such a system enables efficient and transparent sharing of data with both fine- and coarse-grained access patterns. In addition, it supports efficient execution of garbage-collected programs. In contrast, conventional distributed shared memory (DSM) systems are limited to providing only one granularity with good performance, and have experienced difficulty in efficiently supporting garbage collection. A safe language, in which no pointer arithmetic is allowed, can transparently be compiled into a handle-based system and constitutes its preferred mode of use. A programmer can also directly use a handle-based programming model that avoids pointer arithmetic on the handles, and achieve the same performance but without the programming benefits of a safe programming language. This new run-time system, DOSA (Distributed Object Sharing Architecture), provides a shared object space abstraction rather than a shared address space abstraction. The key to its efficiency is the observation that a handle-based distributed implementation permits VM-based access and modification detection without suffering false sharing for fine-grained access patterns. We compare DOSA to TreadMarks, a conventional DSM system that is efficient at handling coarse-grained sharing. The performance of fine-grained applications and garbage-collected applications is considerably better than in TreadMarks, and the performance of coarse-grained applications is nearly as good as in TreadMarks. Inasmuch as the performance of such applications is already good in TreadMarks, we consider this an acceptable performance penalty. Y. Charlie Hu, Weimin Yu, Alan L. Cox, Dan S. Wallach, Willy Zwaenepoel |
ACM Trans. Comput. Syst. | 3 |
| 2002 | Practical, Transparent Operating System Support for Superpages
Juan Navarro, Sitaram Iyer, Peter Druschel, Alan L. Cox |
OSDI | 4 |
| 2001 | Contention elimination by replication of sequential sections in distributed shared memory programsabstractIn shared memory programs contention often occurs at the transition between a sequential and a parallel section of the code. As all threads start executing the parallel section, they often access data just modified by the thread that executed the sequential section, causing a flurry of data requests to converge on that processor. Honghui Lu, Alan L. Cox, Willy Zwaenepoel |
PPoPP | 2 |
| 2000 | Data Replication Strategies for Fault Tolerance and Availability on Commodity ClustersabstractRecent work has shown the advantages of using persistent memory transaction processing. In particular the Vista transaction system uses recoverable memory to avoid disk I/O, thus improving performance by several orders of magnitude. In such a system, however the data is safe when a node fails, but unavailable until it recovers, because the data is kept in only one memory. In contrast, our work uses data replication to provide both reliability and data availability while still maintaining very high transaction throughput. We investigate four possible designs for a primary-backup system, using a cluster of commodity servers connected by a write-through capable system area network (SAN). We show that logging approaches outperform mirroring approaches, even when communicating more data, because of their better locality. Finally, we show that the best logging approach also scales well to small shared-memory multiprocessors. Cristiana Amza, Alan L. Cox, Willy Zwaenepoel |
DSN | 2 |
| 2000 | Improving Fine-Grained Irregular Shared-Memory Benchmarks by Data ReorderingabstractWe demonstrate that data reordering can substantially improve the performance of fine-grained irregular shared-memory benchmarks, on both hardware and software shared-memory systems. In particular, we evaluate two distinct data reordering techniques that seek to co-locate in memory objects that are in close proximity in the physical system modeled by the computation. The effects of these techniques are increased spatial locality and reduced false sharing. We evaluate the effectiveness of the data reordering techniques on a set of five irregular applications from SPLASH-2 and Chaos. We implement both techniques in a small library, allowing us to enable them in an application by adding less than 10 lines of code. Our results on one hardware and two software shared-memory systems show that, with data reordering during initialization, the performance of these applications is improved by 12%-99% on the Origin 2000, 30%-366% on TreadMarks, and 14%-269% on HLRC. Y. Charlie Hu, Alan L. Cox, Willy Zwaenepoel |
SC | 2 |
| 2000 | OpenMP for Networks of SMPs
Y. Charlie Hu, Honghui Lu, Alan L. Cox, Willy Zwaenepoel |
J. Parallel Distributed Comput. | 3 |
| 1999 | A Performance Comparison of Homeless and Home-Based Lazy Release Consistency Protocols in Software Shared MemoryabstractIn this paper, we compare the performance of two multiple-writer protocols based on lazy release consistency. In particular, we compare the performance of Princeton's home-based protocol and TreadMarks' protocol on a 32-processor platform. We found that the performance difference between the two protocols was less than 4% for four out of seven applications. For the three applications on which performance differed by more than 4%, the TreadMarks protocol performed better for two because most of their data were migratory, while the home-based protocol performed better for one. For this one application, the explicit control over the location of data provided by the home-based protocol resulted in a better distribution of communication load across the processors. These results differ from those of a previous comparison of the two protocols. We attribute this difference to (1) a different ratio of memory to network bandwidth on our platform and (2) lazy diffing and request overlapping, two optimizations used by TreadMarks that were not used in the previous study. Alan L. Cox, Eyal de Lara, Y. Charlie Hu, Willy Zwaenepoel |
HPCA | 1 |
| 1999 | Efficient Mining for Association Rules with Relational Database SystemsabstractWith the tremendous growth of large scale data repositories, a need for integrating the exploratory techniques of data mining with the capabilities of relational systems to efficiently handle large volumes of data has now risen. We look at the performance of the most prevalent association rule mining algorithm-Apriori with IBM's DB2 Universal Database system. We show that a multi-column (MC) data model is preferable over the commonly used single column (SC) data model for association rule mining. We obtain factors of 4.8 to 6 improvement in performance for the MC data model over commercial implementations for the SC data model. We provide a new relational operator called Combinations, for efficient SQL implementation of Apriori in the database engine-this results in trivial parallelizability, reliability, and portability for the mining application. Karthick Rajamani, Alan L. Cox, Balakrishna R. Iyer, Atul Chadha |
IDEAS | 2 |
| 1999 | Extending the Applicability of Association Rules
Karthick Rajamani, Sam Yuan Sung, Alan L. Cox |
PAKDD | 3 |
| 1999 | Adaptive protocols for software distributed shared memoryabstractWe demonstrate the benefits of software shared memory protocols that adapt at run time to the memory access patterns observed in the applications. This adaptation is automatic-no user annotations are required-and does not rely on compiler support or special hardware. We investigate adaptation between singleand multiple-writer protocols, dynamic aggregation of pages into a larger transfer unit, and adaptation between invalidate and update. Our results indicate that adaptation between single- and multiple-writer and dynamic page aggregation are clearly beneficial. The results for the adaptation between invalidate and update are less compelling, showing at best gains similar to the dynamic aggregation adaptation and at worst serious performance deterioration. Cristiana Amza, Alan L. Cox, Sandhya Dwarkadas, Li-Jie Jin, Karthick Rajamani, Willy Zwaenepoel |
Proc. IEEE | 2 |
| 1999 | Combining compile-time and run-time support for efficient software distributed shared memoryabstractWe describe an integrated compile time and run time system for efficient shared memory parallel computing on distributed memory machines. The combined system presents the user with a shared memory programming model. The run time system implements a consistent shared memory abstraction using memory access detection and automatic data caching. The compiler improves the efficiency of the shared memory implementation by directing the run time system to exploit the message passing capabilities of the underlying hardware. To do so, the compiler analyzes shared memory accesses and transforms the code to insert calls to the run time system that provide it with the access information computed by the compiler. The run time system is augmented with the appropriate entry points to use this information to implement bulk data transfer and to reduce the overhead of run time consistency maintenance. In those cases where the compiler analysis succeeds for the entire program, we demonstrate that the combined system achieves performance comparable to that produced by compilers that directly target message passing. If the compiler analysis is successful only for parts of the program, for instance, because of irregular accesses to some of the arrays, the resulting optimizations can be applied to those parts for which the analysis succeeds. If the compiler analysis fails entirely, we rely on the run time maintenance of shared memory and thereby avoid the complexity and the limitations of compilers that directly target message passing. The result is a single system that combines efficient support for both regular and irregular memory access patterns. Sandhya Dwarkadas, Honghui Lu, Alan L. Cox, Ramakrishnan Rajamony, Willy Zwaenepoel |
Proc. IEEE | 3 |
| 1997 | Software DSM Protocols that Adapt between Single Writer and Multiple WriterabstractWe present two software DSM protocols that dynamically adapt between a single writer (SW) and a multiple writer (MW) protocol based on the application's sharing patterns. The first protocol (WFS) adapts based on write-write false sharing; the second (WFS+WG) based on a combination of write-write false sharing and write granularity. The adaptation is automatic. No user or compiler information is needed. The choice between SW and MW is made on a per-page basis. We measured the performance of our adaptive protocols on an 8-node SPARC cluster connected by a 155 Mbps ATM network. We used eight applications, covering a broad spectrum in terms of write-write false sharing and write granularity. We compare our adaptive protocols against the MW-only and the SW-only approach. Adaptation to write-write false sharing proves to be the critical performance factor, while adaptation to write granularity plays only a secondary role in our environment and for the applications considered. Each of the two adaptive protocols matches or exceeds the performance of the best of MW and SW in seven out of the eight applications. Cristiana Amza, Alan L. Cox, Sandhya Dwarkadas, Willy Zwaenepoel |
HPCA | 2 |
| 1997 | Trade-offs Between False Sharing and Aggregation in Software Distributed Shared MemoryabstractSoftware Distributed Shared Memory (DSM) systems based on virtual memory techniques traditionally use the hardware page as the consistency unit. The large size of the hardware page is considered to be a performance bottleneck because of the implied false sharing overheads. Instead, we show that in the presence of a relaxed consistency model and a multiple writer protocol, a large consistency unit is generally not detrimental to performance. We study the tradeoffs between false sharing and aggregation effects when using large consistency units. In this context, this paper makes three separate contributions:1. We document the cost of false sharing in terms of extra messages and extra data being communicated. We find that, for the applications considered, when the virtual memory page is used as the consistency unit, the number of extra messages is small, while the amount of extra data can be substantial.2. We evaluate the performance when the consistency unit is increased to a multiple of the virtual memory page size. For most applications and data sets, the performance improves, except when the false sharing effects include extra messages or a large amount of extra data.3. We present a new algorithm for dynamically aggregating pages. In our algorithm, the aggregated pages do not necessarily need to be contiguous. In all cases, the performance of our dynamic aggregation algorithm is similar to that achieved with the best static page size.These results were obtained by measuring the performance of eight applications on the TreadMarks distributed shared memory system. The hardware platform used is a network of 166Mhz Pentiums connected by a switched 100Mbps Ethernet network. Cristiana Amza, Alan L. Cox, Karthick Rajamani, Willy Zwaenepoel |
PPoPP | 2 |
| 1997 | Compiler and Software Distributed Shared Memory Support for Irregular ApplicationsabstractWe investigate the use of a software distributed shared memory (DSM) layer to support irregular computations on distributed memory machines. Software DSM supports irregular computation through demand fetching of data in response to memory access faults. With the addition of a very limited form of compiler support, namely the identification of the section of the indirection array accessed by each processor, many of these on-demand page fetches can be aggregated into a single message, and prefetched prior to the access fault.We have measured the performance of this approach for two irregular applications, moldyn and nbf, using the Tread-Marks DSM system on an 8-processor IBM SP2. We find that it has similar performance to the inspector-executor method supported by the CHAOS run-time library, while requiring much simpler compile-time support. For moldyn, it is up to 23% faster than CHAOS, depending on the input problem's characteristics; and for nbf, it is no worse than 14% slower. If we include the execution time of the inspector, the software DSM-based approach is always faster than CHAOS. The advantage of this approach increases as the frequency of changes to the indirection array increases. The disadvantage of this approach is the potential for false sharing overhead when the data set is small or has poor spatial locality. Honghui Lu, Alan L. Cox, Sandhya Dwarkadas, Ramakrishnan Rajamony, Willy Zwaenepoel |
PPoPP | 2 |
| 1997 | Performance Debugging Shared Memory Parallel Programs Using Run-Time Dependence AnalysisabstractWe describe a new approach to performance debugging that focuses on automatically identifying computation transformations to reduce synchronization and communication. By grouping writes together into equivalence classes, we are able to tractably collect information from long-running programs. Our performance debugger analyzes this information and suggests computation transformations in terms of the source code. We present the transformations suggested by the debugger on a suite of four applications. For Barnes-Hut and Shallow, implementing the debugger suggestions improved the performance by a factor of 1.32 and 34 times respectively on an 8-processor IBM SP2. For Ocean, our debugger identified excess synchronization that did not have a significant impact on performance. ILINK, a genetic linkage analysis program widely used by geneticists, is already well optimized. We use it only to demonstrate the feasibility of our approach to long-running applications.We also give details on how our approach can be implemented. We use novel techniques to convert control dependences to data dependences, and to compute the source operands of stores. We report on the impact of our instrumentation on the same application suite we use for performance debugging. The instrumentation slows down the execution by a factor of between 4 and 169 times. The log files produced during execution were all less than 2.5 Mbytes in size. Ramakrishnan Rajamony, Alan L. Cox |
SIGMETRICS | 2 |
| 1997 | Java/DSM: A Platform for Heterogeneous ComputingabstractIn this paper we describe a system for programming heterogeneous computing environments based upon Java and software distributed shared memory (DSM). Compared with existing approaches for heterogeneous computing, our system transparently handles both the hardware differences and the distributed nature of the system. It is therefore much easier to program. Java is a good candidate for heterogeneous programming because of its portability. Java provides the remote method invocation (RMI) mechanism for distributed computing. However, our early experience with RMI showed that the programmer must expend significant effort on such problems as data replication and the optimization of the remote interface to improve communication efficiency. Furthermore, the need for reference marshaling is not completely eliminated by RMI's effort to facilitate the sharing of linked structures between machines. A DSM system provides a shared memory abstraction over a group of physically distributed machines. It automatically handles data communication between machines and eliminates the need for the programmer to write message-passing code. A multithreaded Java program written for a single machine will require fewer changes to run on a Java/DSM combination than with RMI. We have been implementing a JDK-1.0.2 based parallel Java Virtual Machine on the TreadMarks DSM system. Our implementation includes a distributed garbage collector and supports the Java API with very few changes. In this paper we describe our motivation and the implementation, and report our early experience with programming under both RMI and DSM. © 1997 John Wiley & Sons, Ltd. Weimin Yu, Alan L. Cox |
Concurr. Pract. Exp. | 2 |
| 1997 | Quantifying the Performance Differences between PVM and TreadMarksabstractThis paper compares two systems for parallel programming on networks of workstations: Parallel Virtual Machine (PVM), a message-passing system, and TreadMarks, a software distributed shared-memory (DSM) system. The eight applications used in this comparison are Water and Barnes–Hut from the SPLASH benchmark suite; 3-D FFT, Integer Sort (IS), and Embarrassingly Parallel (EP) from the NAS benchmarks; ILINK, a widely used genetic linkage analysis program; and Successive Over-Relaxation (SOR) and Traveling Salesman (TSP). Two different input data sets are used for five of the applications. We use two execution environments. The first is a 155 Mbps ATM network with eight Sparc-20 model 61 workstations; the second is an eight-processor IBM SP/2. The differences in speedup between TreadMarks and PVM depend mostly on the applications, and only to a much lesser extent on the platform and the data set used. In particular, the TreadMarks speedup for six of the eight applications is within 15% of that achieved with PVM. For one application, the difference in speedup is between 15% and 30%, and for another, the difference is around 50%. We identified four important factors that contribute to the lower performance of TreadMarks: (1) extra messages due to the separation of synchronization and data transfer, (2) extra messages to handle access misses caused by the use of an invalidate protocol, (3) false sharing, and (4) diff accumulation for migratory data. We have quantified the effects of the last three factors by measuring the performance gain when each is eliminated. Of the three factors, TreadMarks' use of a separate request message per page of data accessed is the most important. The effect of false sharing is comparatively low. Reducing diff accumulation benefits migratory data only when the diffs completely overlap. When these performance impediments are removed, all of the TreadMarks programs perform within 25% of PVM, and for six out of eight experiments, TreadMarks is less than 5% slower than PVM. Honghui Lu, Sandhya Dwarkadas, Alan L. Cox, Willy Zwaenepoel |
J. Parallel Distributed Comput. | 3 |
| 1996 | An Integrated Compile-Time/Run-Time Software Distributed Shared Memory SystemabstractOn a distributed memory machine, hand-coded message passing leads to the most efficient execution, but it is difficult to use. Parallelizing compilers can approach the performance of hand-coded message passing by translating data-parallel programs into message passing programs, but efficient execution is limited to those programs for which precise analysis can be carried out. Shared memory is easier to program than message passing and its domain is not constrained by the limitations of parallelizing compilers, but it lags in performance. Our goal is to close that performance gap while retaining the benefits of shared memory. In other words, our goal is (1) to make shared memory as efficient as message passing, whether hand-coded or compiler-generated, (2) to retain its ease of programming, and (3) to retain the broader class of applications it supports.To this end we have designed and implemented an integrated compile-time and run-time software DSM system. The programming model remains identical to the original pure run-time DSM system. No user intervention is required to obtain the benefits of our system. The compiler computes data access patterns for the individual processors. It then performs a source-to-source transformation, inserting in the program calls to inform the run-time system of the computed data access patterns. The run-time system uses this information to aggregate communication, to aggregate data and synchronization into a single message, to eliminate consistency overhead, and to replace global synchronization with point-to-point synchronization wherever possible.We extended the Parascope programming environment to perform the required analysis, and we augmented the TreadMarks run-time DSM library to take advantage of the analysis. We used six Fortran programs to assess the performance benefits: Jacobi, 3D-FFT, Integer Sort, Shallow, Gauss, and Modified Gramm-Schmidt, each with two different data set sizes. The experiments were run on an 8-node IBM SP/2 using user-space communication. Compiler optimization in conjunction with the augmented run-time system achieves substantial execution time improvements in comparison to the base TreadMarks, ranging from 4% to 59% on 8 processors. Relative to message passing implementations of the same applications, the compile-time run-time system is 0-29% slower than message passing, while the base run-time system is 5-212% slower. For the five programs that XHPF could parallelize (all except IS), the execution times achieved by the compiler optimized shared memory programs are within 9% of XHPF. Sandhya Dwarkadas, Alan L. Cox, Willy Zwaenepoel |
ASPLOS | 2 |
| 1996 | A Comparison of Entry Consistency and Lazy Release Consistency ImplementationsabstractThis paper compares several implementations of entry consistency (EC) and lazy release consistency (LRC), two relaxed memory models in use with software distributed shared memory (DSM) systems. We use six applications in our study: SOR, Quicksort, Water, Barnes-Hut, IS, and 3D-FFT. For these applications, EC's requirement that all shared data be associated with a synchronization object leads to a fair amount of additional programming effort. We identify, in particular, extra synchronization, lock rebinding, and object granularity as sources of extra complexity. In terms of performance, for the set of applications and for the computing environment utilized neither model is consistently better than the other. For SOR and IS, execution times are about the same, but LRC is faster for Water (33%) and Barnes-Hut (41%) and EC is faster for Quicksort (14%) and 3D-FFT (10%). Sarita V. Adve, Alan L. Cox, Sandhya Dwarkadas, Ramakrishnan Rajamony, Willy Zwaenepoel |
HPCA | 2 |
| 1996 | Conservative Garbage Collection on DSM SystemsabstractIn this paper we present the design and implementation of a conservative garbage collection algorithm for distributed shared memory (DSM) applications that use weakly-typed languages like C or C++, and evaluate its performance. In the absence of language support to identify references, our algorithm constructed a conservative approximation of the set of cross-node references based on local information only. It was also designed to tolerate memory inconsistency on DSM systems that use relaxed consistency protocols. These techniques enabled every node to perform garbage collections without communicating with others, effectively avoiding the high cost of cross-node communication in networks of workstations. We measured the performance of our garbage collector against explicit programmer management using three application programs. In two out of the three programs the performance of the GC version is within 15% of the explicit version. The results showed that the garbage collector has two effects on application programs. On one hand, it tends to reduce memory locality, increasing the communication cost; on the other hand, it may eliminate synchronization and memory accesses that would be incurred if memory were managed by the programmer reducing the communication cost. Weimin Yu, Alan L. Cox |
ICDCS | 2 |
| 1995 | Message Passing Versus Distributed Shared Memory on Networks of WorkstationsabstractThe message passing programs are executed with the Parallel Virtual Machine (PVM) library and the shared memory programs are executed using TreadMarks. The programs are Water and Barnes-Hut from the SPLASH benchmark suite; 3-D FFT, Integer Sort (IS) and Embarrassingly Parallel (EP) from the NAS benchmarks; ILINK, a widely used genetic linkage analysis program; and Successive Over-Relaxation (SOR), Traveling Salesman (TSP), and Quicksort (QSORT). Two different input data sets were used for Water (Water-288 and Water-1728), IS (IS-Small and IS-Large), and SOR (SOR-Zero and SOR-NonZero). Our execution environment is a set of eight HP735 workstations connected by a 100Mbits per second FDDI network. For Water-1728, EP, ILINK, SOR-Zero, and SOR-NonZero, the performance of TreadMarks is within 10%of PVM. For IS-Small, Water-288, Barnes-Hut, 3-D FFT, TSP, and QSORT, differences are on the order of 10%to 30%. Finally, for IS-Large, PVM performs two times better than TreadMarks. More messages and more data are sent in TreadMarks, explaining the performance differences. This extra communication is caused by 1) the separation of synchronization and data transfer, 2) extra messages to request updates for data by the invalidate protocol used in TreadMarks, 3) false sharing, and 4) diff accumulation for migratory data in TreadMarks. Honghui Lu, Sandhya Dwarkadas, Alan L. Cox, Willy Zwaenepoel |
SC | 3 |
| 1995 | An Evaluation of Software-Based Release Consistent Protocols
Peter J. Keleher, Alan L. Cox, Sandhya Dwarkadas, Willy Zwaenepoel |
J. Parallel Distributed Comput. | 2 |
| 1994 | Software Versus Hardware Shared-Memory Implementation: A Case StudyabstractCompares the performance of software-supported shared memory on a general-purpose network to hardware-supported shared memory on a dedicated interconnect. Up to eight processors, the results are based on the execution of a set of application programs on a SGI 4D/480 multiprocessor and on TreadMarks, a distributed shared memory system that runs on a Fore ATM LAN of DECstation-5000/240s. Since the DECstation and the 4D/480 use the same processor, primary cache, and compiler, the shared-memory implementation is the principal difference between the systems. Beyond eight processors, the results are based on execution-driven simulation. Specifically, the authors compare a software implementation on a general-purpose network of uniprocessor nodes, a hardware implementation using a directory-based protocol on a dedicated interconnect, and a combined implementation using software to provide shared memory between multiprocessor nodes with hardware implementing shared memory within a node.> Alan L. Cox, Sandhya Dwarkadas, Peter J. Keleher, Honghui Lu, Ramakrishnan Rajamony, Willy Zwaenepoel |
ISCA | 1 |
| 1993 | Adaptive Cache Coherency for Detecting Migratory Shared DataabstractParallel programs exhibit a small number of distinct data-sharing patterns. A common data-sharing pattern, migratory access, is characterized by exclusive read and write access by one processor at a time to a shared datum. We describe a family of adaptive cache coherency protocols that dynamically identify migratory shared data in order to reduce the cost of moving them. The protocols use a standard memory model and processor-cache interface. They do not require any compile-time or run-time software support. We describe implementations for bus-based multiprocessors and for shared-memory multiprocessors that use directory-based caches. These implementations are simple and would not significantly increase hardware cost. We use trace- and execution-driven simulation to compare the performance of the adaptive protocols to standard write-invalidate protocols. These simulations indicate that, compared to conventional protocols, the use of the adaptive protocol can almost halve the number of inter-node messages on some applications. Since cache coherency traffic represents a larger part of the total communication as cache size increases, the relative benefit of using the adaptive protocol also increases. Alan L. Cox, Robert J. Fowler |
ISCA | 1 |
| 1993 | Evaluation of Release Consistent Software Distributed Shared Memory on Emerging Network TechnologyabstractWe evaluate the effect of processor speed, network characteristics, and software overhead on the performance of release-consistent software distributed shared memory. We examine five different protocols for implementing release consistency: eager update, eager invalidate, lazy update, lazy invalidate, and a new protocol called lazy hybrid. This lazy hybrid protocol combines the benefits of both lazy update and lazy invalidate. Sandhya Dwarkadas, Peter J. Keleher, Alan L. Cox, Willy Zwaenepoel |
ISCA | 3 |
| 1992 | Lazy Release Consistency for Software Distributed Shared MemoryabstractRelaxed memory consistency models, such as release consistency, were introduced in order to reduce the impact of remote memory access latency in both software and hardware distributed shared memory (DSM). However, in a software DSM, it is also important to reduce the number of messages and the amount of data exchanged for remote memory access. Lazy release consistency is a new algorithm for implementing release consistency that lazily pulls modifications across the interconnect only when necessary. Trace-driven simulation using the SPLASH benchmarks indicates that lazy release consistency reduces both the number of messages and the amount of data transferred between processors. These reductions are especially significant for programs that exhibit false sharing and make extensive use of locks. Peter J. Keleher, Alan L. Cox, Willy Zwaenepoel |
ISCA | 2 |
| 1991 | NUMA Policies and Their Relation to Memory ArchitectureabstractMultiprocessor memory reference traces provide a wealth of information on the behavior of parallel programs.We have used this information to explore the relationship between kernel-based NUMA management policies and multiprocessor memory architecture.Our trace analysis techniques employ an off-line, optimal cost policy as a baseline against which to compare on-line policies, and as a policyinsensitive tool for evaluating architectural design alternatives.We compare the performance of our optimal policy with that of three implementable policies (two of which appear in previous work), on a variety of applications, with varying relative speeds for page moves and local, global, and remote memory references.Our results indicate that a good NUMA policy must be chosen to match its machine, and confirm that such policies can be both simple and effective.They also indicate that programs for NUMA machines must be written with care to obtain the best performance. William J. Bolosky, Michael L. Scott, Robert P. Fitzgerald, Robert J. Fowler, Alan L. Cox |
ASPLOS | 5 |
| 1989 | The Implementation of a Coherent Memory Abstraction on a NUMA Multiprocessor: Experiences with PLATINUMabstractPLATINUM is an operating system kernel with a novel memory management system for Non-Uniform Memory Access (NUMA) multiprocessor architectures. This memory management system implements a coherent memory abstraction. Coherent memory is uniformly accessible from all processors in the system. When used by applications coded with appropriate programming styles it appears to be nearly as fast as local physical memory and it reduces memory contention. Coherent memory makes programming NUMA multiprocessors easier for the user while attaining a level of performance comparable with hand-tuned programs. Alan L. Cox, Robert J. Fowler |
SOSP | 1 |