Liana L. Fong

dblp:83/4382 · also Liana Fong · DBLP profile ↗
← Back
33ranked-venue papers
3as first author
0since 2021 · last 2018
—ORCID · none

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

Systems, architecture and hardware · 18 · 1 first-authorSoftware engineering, systems software and programming languages · 8 · 1 first-authorArtificial intelligence and machine learning · 3Computer networks · 2Databases, data management, data science and information retrieval · 1Theory of computation · 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
4 papers
GPUs and heterogeneous computing · 38% Parallel and multicore computing · 38% High-performance computing · 18%
Artificial intelligence
1 paper
Optimization for machine learning · 56% Representation and self-supervised learning · 44%

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

TopicWeightPapersLastEvidence papers
GPUs and heterogeneous computing
GPU computing
0.522017
CuMF_SGD: Parallelized Stochastic Gradient Descent for Matrix Factorization on GPUs · HPDC 2017
Faster and Cheaper: Parallelizing Large-Scale Matrix Factorization on GPUs · HPDC 2016
Machine learning › Representation and self-supervised learning
matrix factorization
0.312017
CuMF_SGD: Parallelized Stochastic Gradient Descent for Matrix Factorization on GPUs · HPDC 2017
Machine learning › Optimization for machine learning › distributed optimization
parallel stochastic gradient descent
0.312017
CuMF_SGD: Parallelized Stochastic Gradient Descent for Matrix Factorization on GPUs · HPDC 2017
Parallel and multicore computing › parallelization strategies
model and data parallelism
0.312017
CuMF_SGD: Parallelized Stochastic Gradient Descent for Matrix Factorization on GPUs · HPDC 2017
High-performance computing › numerical linear algebra
matrix factorization
0.212016
Faster and Cheaper: Parallelizing Large-Scale Matrix Factorization on GPUs · HPDC 2016
Parallel and multicore computing
parallel algorithms
0.212016
Faster and Cheaper: Parallelizing Large-Scale Matrix Factorization on GPUs · HPDC 2016
Machine learning › Optimization for machine learning
stochastic gradient descent
0.112017
CuMF_SGD: Parallelized Stochastic Gradient Descent for Matrix Factorization on GPUs · HPDC 2017
Cloud and datacenter computing
cluster resource management and scheduling
0.011998
An Infrastructure for Efficient Parallel Job Execution in Terascale Computing Environments · SC 1998
Embedded and real-time systems › real-time scheduling › multiprocessor scheduling
gang scheduling
0.011998
An Infrastructure for Efficient Parallel Job Execution in Terascale Computing Environments · SC 1998
Cloud and datacenter computing
job scheduling
0.011998
An Infrastructure for Efficient Parallel Job Execution in Terascale Computing Environments · SC 1998
Distributed systems
time-sharing systems
0.011998
An Infrastructure for Efficient Parallel Job Execution in Terascale Computing Environments · SC 1998
Operating systems › resource management › process management
CPU scheduling
0.011995
Time-Function Scheduling: A General Approach To Controllable Resource Management · SOSP 1995
Cloud and datacenter computing
resource management
0.011995
Time-Function Scheduling: A General Approach To Controllable Resource Management · SOSP 1995

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

model parallelism · 0.6data parallelism · 0.6matrix factorization · 0.2two-phase commit · 0.0ousterhout matrix · 0.0hierarchical distribution · 0.0
YearPublicationVenuePosition
2018 Matrix Factorization on GPUs with Memory Optimization and Approximate Computing
abstract
Matrix factorization (MF) discovers latent features from observations, which has shown great promises in the fields of collaborative filtering, data compression, feature extraction, word embedding, etc. While many problem-specific optimization techniques have been proposed, alternating least square (ALS) remains popular due to its general applicability (e.g. easy to handle positive-unlabeled inputs), fast convergence and parallelization capability. Current MF implementations are either optimized for a single machine or with a need of a large computer cluster but still are insufficent. This is because a single machine provides limited compute power for large-scale data while multiple machines suffer from the network communication bottleneck.
Wei Tan 0001, Shiyu Chang, Liana L. Fong, Cheng Li 0014, Liangliang Cao
ICPP3
2017 CuMF_SGD: Parallelized Stochastic Gradient Descent for Matrix Factorization on GPUs
abstract
Stochastic gradient descent (SGD) is widely used by many machine learning algorithms. It is efficient for big data ap- plications due to its low algorithmic complexity. SGD is inherently serial and its parallelization is not trivial. How to parallelize SGD on many-core architectures (e.g. GPUs) for high efficiency is a big challenge. In this paper, we present cuMF_SGD, a parallelized SGD solution for matrix factorization on GPUs. We first design high-performance GPU computation kernels that accelerate individual SGD updates by exploiting model parallelism. We then design efficient schemes that parallelize SGD updates by exploiting data parallelism. Finally, we scale cuMF SGD to large data sets that cannot fit into one GPU's memory. Evaluations on three public data sets show that cuMF_SGD outperforms existing solutions, including a 64- node CPU system, by a large margin using only one GPU card.
Xiaolong Xie, Wei Tan 0001, Liana L. Fong, Yun Liang 0001
HPDC3
2016 Faster and Cheaper: Parallelizing Large-Scale Matrix Factorization on GPUs
abstract
Matrix factorization (MF) is used by many popular algorithms such as collaborative filtering. GPU with massive cores and high memory bandwidth sheds light on accelerating MF much further when appropriately exploiting its architectural characteristics.
Wei Tan 0001, Liangliang Cao, Liana L. Fong
HPDC3
2015 Deferred Lightweight Indexing for Log-Structured Key-Value Stores
abstract
The recent shift towards write-intensive workload on big data (e.g., financial trading, social user-generated data streams)has pushed the proliferation of log-structured key-value stores, represented by Google's BigTable [1], Apache HBase [2] andCassandra [3]. While providing key-based data access with aPut/Get interface, these key-value stores do not support value-based access methods, which significantly limits their applicability in modern web and database applications. In this paper, we present DELI, a DEferred Lightweight Indexing scheme on the log-structured key-value stores. To index intensively updated bigdata in real time, DELI aims at making the index maintenance as lightweight as possible. The key idea is to apply an append-only design for online index maintenance and to collect index garbage at carefully chosen time. DELI optimizes the performance of index garbage collection through tightly coupling its execution with a native routine process called compaction. The DELI's system design is fault-tolerant and generic (to most key-valuestores), we implemented a prototype of DELI based on HBase without internal code modification. Our experiments show that the DELI offers significant performance advantage for the write-intensive index maintenance.
Yuzhe Tang, Arun Iyengar, Wei Tan 0001, Liana L. Fong, Ling Liu 0001, Balaji Palanisamy
CCGRID4
2014 Diff-Index: Differentiated Index in Distributed Log-Structured Data Stores
abstract
Log-Structured-Merge (LSM) Tree gains much attention re-cently because of its superior performance in write-intensive workloads. LSM Tree uses an append-only structure in memory to achieve low write latency; at memory capac-ity, in-memory data are flushed to other storage media (e.g. disk). Consequently, read access is slower comparing to write. These specific features of LSM, including no in-place update and asymmetric read/write performance raise unique challenges in index maintenance for LSM. The structural difference between LSM and B-Tree also prevents mature B-Tree based approaches from being directly applied. To address the issues of index maintenance for LSM, we pro-pose Diff-Index to support a spectrum of index maintenance schemes to suit different objectives in index consistency and performance. The schemes consist of sync-full, sync-insert, async-simple and async-session. Experiments on our HBase implementation quantitatively demonstrate that Diff-Index offers various performance/consistency balance and satisfac-tory scalability while avoiding global coordination. Sync-insert and async-simple can reduce 60%-80 % of the overall index update latency when compared to the baseline sync-full; async-simple can achieve superior index update per-formance with an acceptable inconsistency. Diff-Index ex-ploits LSM features such as versioning and the flush-compact process to achieve goals of concurrency control and failure
Wei Tan 0001, Sandeep Tata, Yuzhe Tang, Liana L. Fong
EDBT4
2014 Effectiveness Assessment of Solid-State Drive Used in Big Data Services
abstract
Big data poses challenges to the technologies required to process data of high volume, velocity, variety, and veracity. Among the challenges, the storage and computing required by big data analytics is usually huge, and as a result big data capabilities are often provisioned in cloud and delivered in the form of Web-based services. Solid-state drive (SSD) is widely used nowadays as an elementary hardware feature in cloud infrastructure for big data services. For example, Amazon Web Service (AWS) offers EC2 instances with SSD storage, and its key-value data store, DynamoDB, is backed up by SSD for superior performance. Compared to hard disk drive (HDD), SSD prevails in both access latency and bandwidth. In the foreseeable future, SSD would be readily available on commodity servers though its capacity would be neither large enough nor cost effective to accommodate big data on its own. Therefore, it is essential to investigate how to efficiently leverage SSD as one layer in a storage hierarchy in addition to HDD. In this paper, we investigate the effectiveness of using SSD in three workloads, namely standalone Hadoop MapReduce jobs, Hive jobs, and HBase queries. Firstly, we device an approach to enable Hadoop Distributed File System (HDFS) having a SSD-HDD storage hierarchy. Secondly, we investigate the IO involved in different phases of Hadoop jobs and design different schemes to place data discriminatively in the aforementioned storage hierarchy. Afterward, the effectiveness of different schemes are evaluated with respect to job run time. Finally, we summarize best practices of data placement for examined workloads in a SSD-HDD storage hierarchy.
Wei Tan 0001, Liana L. Fong
ICWS2
2013 Enabling Interoperability among Grid Meta-Schedulers
Ivan Rodero, David Villegas, Norman Bobroff, Liana L. Fong, Seyed Masoud Sadjadi
J. Grid Comput.5
2013 Extreme scale computing: Modeling the impact of system noise in multi-core clustered systems
Seetharami Seelam, Liana L. Fong, Asser N. Tantawi, John Lewars, John Divirgilio, Kevin J. Gildea
J. Parallel Distributed Comput.2
2012 Partitioned Parallel Job Scheduling for Extreme Scale Computing
David Brelsford, George Chochia, Nathan Falk, Kailash Marthi, Ravindra Sure, Norman Bobroff, Liana L. Fong, Seetharami Seelam
JSSPP7
2012 Experiences in building and scaling an enterprise application on multicore systems
abstract
SUMMARY Even though Java is the de facto programming language for enterprise applications, there exist only a limited number of Java‐based benchmarks to understand the performance on emerging multicore systems. To bridge this gap, this paper presents a report generation benchmark that is developed on top of Open Source Apache Geronimo's DayTrader benchmark. Report generation and rendering is at the heart of many enterprise business analytics and business intelligence software products, and it is used by many enterprise applications. We evaluate the performance scalability of this benchmark on a state‐of‐the‐art Power7 multicore system with 8 Power7 cores and 32 hardware threads. The benchmark throughput scales linearly up to eight hardware threads, but beyond that point, the throughput falls sharply. Significant locking in the Java class libraries for non‐shared objects results in this performance drop. Splitting the locks on these shared classes results in near linear scaling from eight to 32 threads and improved the throughput by 80%. We also show that the Linux operating system load balancing could result in a degraded application performance in hardware multithreaded systems and simultaneous‐multithreads‐aware task scheduling results in uniform core‐resource utilization as well as improved application performance. Copyright © 2011 John Wiley & Sons, Ltd.
Seetharami Seelam, Parijat Dube, Megumi Ito, Deniz Binay, Michael Dawson 0001, Pramod Nagaraja, Graeme Johnson, Liana L. Fong, Michel Hack, Xiaoqiao Meng, Li Zhang 0002
Concurr. Comput. Pract. Exp.9
2012 Cloud federation in a layered service model
David Villegas, Norman Bobroff, Ivan Rodero, Javier Delgado, Aditya Devarakonda, Liana L. Fong, Seyed Masoud Sadjadi, Manish Parashar
J. Comput. Syst. Sci.7
2011 Analysis and Modeling of Social Influence in High Performance Computing Workloads
Shuai Zheng 0002, Zon-Yin Shae, Xiangliang Zhang 0001, Hani Jamjoom, Liana L. Fong
Euro-Par (1)5
2011 Characterization of System Services and Their Performance Impact in Multi-core Nodes
abstract
The performance of parallel applications on large scale systems is shown to disproportionately degrade due to interference from system services. This interference from system services is also known as jitter. However, there is limited understanding of sources and patterns of jitter on multi-core systems. In this paper, we identify and characterize jitter sources in terms of their amplitude and execution interval distributions on multi-core IBM Power systems with UNIX-based general purpose operating systems: AIX and Linux. Our analysis shows that there are various kinds of jitter sources and their execution varies drastically between different cores and between hardware threads within each core for practical reasons. This in-depth knowledge of jitter events is leveraged to devise effective approaches to mitigate the jitter impact on application performance in large scale systems. Moreover, such knowledge would provide useful insights to a new generation of operating system designs such as multikernel or satellite kernel for multicore systems.
Seetharami Seelam, Liana L. Fong, John Lewars, John Divirgilio, Brian F. Veale, Kevin J. Gildea
IPDPS2
2010 Extreme scale computing: Modeling the impact of system noise in multicore clustered systems
abstract
System noise or Jitter is the activity of hardware, firmware, operating system, runtime system, and management software events. It is shown to disproportionately impact application performance in current generation large-scale clustered systems running general-purpose operating systems (GPOS). Jitter mitigation techniques such as co-scheduling jitter events across operating systems improve application performance but their effectiveness on future petascale systems is unknown. To understand if existing co-scheduling solutions enable scalable petascale performance, we construct two complementary jitter models based on detailed analysis of system noise from the nodes of a large-scale system running a GPOS. We validate these two models using experimental data from a system consisting of 128 GPOS instances with 4096 CPUs. Based on our models, we project a minimum slowdown of 2.1%, 5.9%, and 11.5% for applications executing on a similar one petaflop system running 1024 GPOS instances and having global synchronization operations once every 1000 msec, 100 msec, and 10 msec, respectively. Our projections indicate that additional system noise mitigation techniques are required to contain the impact of jitter on multi-petaflop systems, especially for tightly synchronized applications.
Seetharami Seelam, Liana L. Fong, Asser N. Tantawi, John Lewars, John Divirgilio, Kevin J. Gildea
IPDPS2
2010 Distributed and Adaptive Execution of Condor DAGMan Workflows
Selim Kalayci, Gargi Dasgupta, Liana L. Fong, Onyeka Ezenwoye, Seyed Masoud Sadjadi
SEKE3
2010 Grid broker selection strategies using aggregated resource information
Ivan Rodero, Francesc Guim 0001, Julita Corbalán, Liana L. Fong, Seyed Masoud Sadjadi
Future Gener. Comput. Syst.4
2009 Broker Selection Strategies in Interoperable Grid Systems
abstract
The increasing demand for resources of the high performance computing systems has led to new forms of collaboration of distributed systems such as interoperable grid systems that contain and manage their own resources. While with a single grid domain one of the most important tasks is the selection of the most appropriate set of resources to dispatch a job, in an interoperable grid environment this problem shifts to selecting the most appropriate domain containing the requiring resources for the job. In this paper, we present and evaluate broker selection strategies for interoperable grid systems. They use aggregated resource information as well as dynamic performance information of the underlying scheduling layers. From our evaluations performed with simulation tools, we conclude that aggregation techniques do not penalize performance significantly, and that delegating part of the scheduling responsibilities to the underlying scheduling layers is a good way to balance the load among the different grid systems.
Ivan Rodero, Francesc Guim 0001, Julita Corbalán, Liana L. Fong, Seyed Masoud Sadjadi
ICPP4
2009 Scalability Analysis of Job Scheduling Using Virtual Nodes
Norman Bobroff, Richard Coppinger, Liana L. Fong, Seetharami R. Seelam, Jing Xu 0012
JSSPP3
2009 Task Decomposition for Adaptive Data Staging in Workflows for Distributed Environments
Onyeka Ezenwoye, Balaji Viswanathan, Seyed Masoud Sadjadi, Liana L. Fong, Gargi Dasgupta, Selim Kalayci
SEKE4
2008 Enabling Interoperability among Meta-Schedulers
abstract
Grid computing supports shared access to computing resources from cooperating organizations or institutes in the form of virtual organizations. Resource brokering middleware, commonly known as a meta-scheduler or a resource broker, matches jobs to distributed resources. Recent advances in meta- scheduling capabilities are extended to enable resource matching across multiple virtual organizations. Several architectures have been proposed for interoperating meta-scheduling systems. This paper presents a hybrid approach, combining hierarchical and peer-to-peer architectures for flexibility and extensibility of these systems. A set of protocols are introduced to allow different meta-scheduler instances to communicate over Web Services. Interoperability between three heterogeneous and distributed organizations (namely, BSC, FIU, and IBM), each using different meta-scheduling technologies, is demonstrated under these protocols and resource models.
Norman Bobroff, Liana L. Fong, Selim Kalayci, Ivan Rodero, Seyed Masoud Sadjadi, David Villegas
CCGRID2
2008 Design and Implementation of a Fault Tolerant Job Flow Manager Using Job Flow Patterns and Recovery Policies
Selim Kalayci, Onyeka Ezenwoye, Balaji Viswanathan, Gargi Dasgupta, Seyed Masoud Sadjadi, Liana L. Fong
ICSOC6
2008 Design of a Fault-tolerant Job-flow Manager for Grid Environments Using Standard Technologies, Job-flow Patterns, and a Transparent Proxy
Gargi Dasgupta, Onyeka Ezenwoye, Liana L. Fong, Selim Kalayci, Seyed Masoud Sadjadi, Balaji Viswanathan
SEKE3
2007 BPEL4Job: A Fault-Handling Design for Job Flow Management
Wei Tan 0001, Liana L. Fong, Norman Bobroff
ICSOC2
2007 A Version-aware Approach for Web Service Directory
abstract
In real-world scenarios, the evolution of Web services to meet functional and non-functional changes ultimately leads to multiple versions of the same original service. Thus, design and implementation of version management techniques, such as version description, directory, etc, play a critical role in realizing the full promise of SOA. To address the version management issues in Web services, we propose a version-aware service model based on some architectural extensions to WSDL and UDDI. WSDL would be enhanced to describe the attributes of the service versions. UDDI would be augmented to use versions in a service directory with an event-based notification/subscription mechanism. We also design a proxy, residing in the service consumer side which can dynamically update the client application instance at runtime. We have implemented a prototype to demonstrate these models and used a weather forecast web service as an example to illustrate the usefulness of the proposed architecture.
Ru Fang, Linh Lam, Liana L. Fong, David Frank, Christopher Vignola
ICWS3
2007 A Version-aware Approach for Web Service Client Application
abstract
An increasing number of enterprises demonstrate that successful adoption of Service Oriented Architecture (SOA) using Web services technologies enables them to build enterprise applications quickly and effectively. To align with changing business requirement, services need to adapt quickly, and eventually multiple versions of the same original service would coexist. To manage all these versions and ensure continuous availability to the service consumers, innovative techniques of version management for Web services become critical to realizing the full promise of SOA. To address the version management issues in Web services, we propose to include version-awareness to various aspects of web services as an extension to the current SOA. In particular, to minimize the impact of service changes on the service consumer side, we design a version-aware web service client model (via an enhancement to the current JAX-RPC client model) which provides both consumer-aware and consumer-transparent invocation styles at build-time and dynamic service proxy generation at runtime. Leveraging the present implementation of the JAX-RPC service model, the Versioned Client APIs based on the new client model is designed to make the development process easy and intuitive. A prototype of this client model, implemented in Eclipse with an exemplary weather-forecast application, is introduced to demonstrate the usefulness of the proposed approach.
Ru Fang, Liana L. Fong, Linh Lam, David Frank, Christopher Vignola
Integrated Network Management3
2006 Resource management with stateful support for analytic applications
abstract
Analytic applications from various industrial sectors have specific attributes and requirements including relatively long processing time, parallelization, multiple interactive invocations, Web services, and expected quality of service objectives. Current parallel resource management systems for batch-oriented jobs lack the effective support for multiple interactive invocations with consideration in quality of service objectives, while transaction processing systems do not support dynamic creation of parallel application instances. To better serve the analytic applications, a set of additional resource management services, defined as stateful support, introduces the concept of service instance and service instance management. This set of stateful support services can be implemented as extension to existing parallel resource management to serve these analytic applications that rapidly increase in the demand of computing power
Liana L. Fong, Catherine H. Crawford, Hidayatullah Shaikh
IPDPS1
2004 Adaptive Memory Paging for Efficient Gang Scheduling of Parallel Applications
abstract
Summary form only given. The gang scheduling paradigm allows timesharing of computing nodes by multiple parallel applications and supports the coordinated context switches of these applications. It can improve system responsiveness and resource utilization. However, the memory paging overhead incurred during context switches can be expensive and may diminish the positive effects of gang scheduling. We investigate the reduction of paging overhead in gang scheduling environments by applying a set of simple, yet effective, adaptive paging techniques: selective page-out, aggressive page-out, adaptive page-in and background writing. Our experiments with NAS NPB2 benchmark programs show that these new adaptive paging mechanisms can reduce the job switching time significantly (up to 90%).
Kyung Dong Ryu, Nimish Pachapurkar, Liana L. Fong
IPDPS3
2002 Neptune: A Dynamic Resource Allocation and Planning System for a Cluster Computing Utility
abstract
We present Neptune - the resource director of Océano, a policy driven fabric management system that dynamically reconfigures resources in a computing utility cluster. Neptune implements an on-line control mechanism subject to policy-based performance and resource configuration objectives. Neptune reassigns servers and bandwidth among a set of service domains, based on pre-defined policy, in response to workload changes. It builds and executes a reconfiguration plan through a planning framework, breaking reconfiguration objectives into individual tasks delegated to set of lower level resource managers. We describe an example decision policy algorithm that we implemented and demonstrated in an 80 server multi-domain computing utility.
Donald P. Pazel, Tamar Eilam, Liana L. Fong, Michael H. Kalantar, Karen Appleby, Germán S. Goldszmidt
CCGRID3
2002 Dynamic resource management in an eUtility
abstract
Oceano is management software for a eUtility infrastructure capable of providing cost-effective, autonomic resource allocation for multiple customer or application domains, in response to existent performance and availability conditions. The control layer of Oceano provides mechanisms to manage resources. This layer consists of a resource director and a set of resource managers. The resource director formulates resource configurations and coordinates their execution through resource managers. These resource managers maintain restore state and carry out detailed configuration tasks. This paper describes the Oceano resource management model and its server and application deployment resource managers. A prototype of Oceano has been developed and deployed on an 80 server platform, and has been tested with multiple domains and applications.
Liana L. Fong, Michael H. Kalantar, Donald P. Pazel, Germán S. Goldszmidt, Sameh A. Fakhouri, Srirama M. Krishnakumar
NOMS1
2001 Océano - SLA Based Management of a Computing Utility
abstract
Oceano is a prototype of a highly available, scaleable, and manageable infrastructure for an e-business computing utility. It enables multiple customers to be hosted on a collection of sequentially shared resources. The hosting environment is divided into secure domains, each supporting one customer. These domains are dynamic: the resources assigned to them may be augmented when the load increases and reduced when load dips. This dynamic resource allocation enables flexible service level agreements (SLAs) with customers in an environment where peak loads are an order of magnitude greater than the normal steady state.
Karen Appleby, Sameh A. Fakhouri, Liana L. Fong, Germán S. Goldszmidt, Michael H. Kalantar, Srirama M. Krishnakumar, Donald P. Pazel, John A. Pershing, Benny Rochwerger
Integrated Network Management3
1998 An Infrastructure for Efficient Parallel Job Execution in Terascale Computing Environments
abstract
Recent Terascale computing environments, such as those in the Department of Energy Accelerated Strategic Computing Initiative, present a new challenge to job scheduling and execution systems. The traditional way to concurrently execute multiple jobs in such large machines is through space-sharing: each job is given dedicated use of a pool of processors. Previous work in this area has demonstrated the benefits of sharing the parallel machine's resources not only spatially but also temporally. Time-sharing creates virtual processors for the execution of jobs. The scheduling is typically performed cyclically and each time-slice of the cycle can be considered an independent virtual machine. When all tasks of a parallel job are scheduled to run on the same time-slice (same virtual machine), gang-scheduling is accomplished. Research has shown that gang-scheduling can greatly improve system utilization and job response time in large parallel systems. We are developing GangLL, a research prototype system for performing gang-scheduling on the ASCI Blue-Pacific machine, an IBM RS/6000 SP to be installed at Lawrence Livermore National Laboratory. This machine consists of several hundred nodes, interconnected by a high-speed communication switch. GangLL is organized as a centralized scheduler that performs global decision-making, and a local daemon in each node that controls job execution according to those decisions. The centralized scheduler builds an Ousterhout matrix that precisely defines the temporal and spatial allocation of tasks in the system. Once the matrix is built, it is distributed to each of the local daemons using a scalable hierarchical distributions scheme. A two-phase commit is used in the distribution scheme to guarantee that all local daemons have consistent information. The local daemons enforce the schedule dedicated by the Ousterhout matrix in their corresponding nodes. This requires suspending and resuming execution of tasks and multiplexing access to the communication switch. Large supercomputing centers tend to have their own job scheduling systems, to handle site specific conditions. Therefore, we are designing GangLL so that it can interact with an external site scheduler. The goal is to let the site scheduler control spatial allocation of jobs, if so desired, and to decide when jobs run. GangLL then performs the detailed temporal allocation and controls the actual execution of jobs. The site scheduler can control the fraction of a shared processor that a job receives through an execution factor parameter. To quantify the benefits of our gang-scheduling system to job execution in a large parallel system, we simulate the system with a realistic workload. We measure performance parameters under various degrees of time-sharing, characterized by the multiprogramming level. Our results show that higher multiprogramming levels lead to higher system utilization and lower job response times. We also report some results from the initial deployment of GangLL on a small multiprocessor system.
José E. Moreira, Waiman Chan, Liana L. Fong, Hubertus Franke, Morris A. Jette
SC3
1997 Extensible Resource Management for Cluster Computing
abstract
Advanced general purpose parallel systems should be able to support diverse applications with different resource requirements without compromising effectiveness and efficiency. We present a resource management model for cluster computing that allows multiple scheduling policies to co-exist dynamically. In particular, we have built Octopus, an extensible and distributed hierarchical scheduler that implements new space sharing, gang scheduling and load sharing strategies. A series of experiments performed on an IBM SP2 suggest that Octopus can effectively match application requirements to available resources, and improve the performance of a variety of parallel applications within a cluster.
Nayeem Islam, Andreas L. Prodromidis, Mark S. Squillante, Ajei S. Gopal, Liana L. Fong
ICDCS5
1995 Time-Function Scheduling: A General Approach To Controllable Resource Management
abstract
No abstract available.
Liana L. Fong, Mark S. Squillante
SOSP1