Daniel S. Katz

dblp:78/3211 · DBLP profile ↗
← Back
62ranked-venue papers
7as first author
9since 2021 · last 2026
0000-0001-5934-7525ORCID · verified

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

Systems, architecture and hardware · 42 · 5 first-author · 4 since 2021Applied, interdisciplinary, general and emerging computing · 16 · 1 first-author · 4 since 2021Software engineering, systems software and programming languages · 15 · 1 first-author · 3 since 2021Security and privacy · 4Databases, data management, data science and information retrieval · 2 · 1 first-author · 1 since 2021Theory of computation · 1
YearPublicationVenuePosition
2026 A terminology for scientific workflow systems
Frédéric Suter, Tainã Coleman, Ilkay Altintas, Rosa M. Badia, Bartosz Balis, Kyle Chard, Iacopo Colonnelli, Ewa Deelman, Paolo Di Tommaso, Thomas Fahringer, Carole A. Goble, Shantenu Jha, Daniel S. Katz, Johannes Köster, Ulf Leser, Kshitij Mehta, Hilary Oliver, Jayson Luc Peterson, Giovanni Pizzi, Loïc Pottier, Raül Sirvent, Eric Suchyta, Douglas Thain, Sean R. Wilkinson, Justin M. Wozniak, Rafael Ferreira da Silva
Future Gener. Comput. Syst.13
2025 Preface of Special Issue on Highlights from the Joint-Laboratory on Extreme Scale Computing
Franck Cappello, Ruth Partzsch, Daniel S. Katz
Future Gener. Comput. Syst.3
2024 FAIR-USE4OS: Guidelines for creating impactful open-source software
abstract
This paper extends the FAIR (Findable, Accessible, Interoperable, Reusable) guidelines to provide criteria for assessing if software conforms to best practices in open source. By adding "USE" (User-Centered, Sustainable, Equitable), software development can adhere to open source best practice by incorporating user-input early on, ensuring front-end designs are accessible to all possible stakeholders, and planning long-term sustainability alongside software design. The FAIR-USE4OS guidelines will allow funders and researchers to more effectively evaluate and plan open-source software projects. There is good evidence of funders increasingly mandating that all funded research software is open source; however, even under the FAIR guidelines, this could simply mean software released on public repositories with a Zenodo DOI. By creating FAIR-USE software, best practice can be demonstrated from the very beginning of the design process and the software has the greatest chance of success by being impactful.
Raphael Sonabend, Hugo Gruson, Leo Wolansky, Agnes Kiragga, Daniel S. Katz
PLoS Comput. Biol.5
2023 Research Software Engineering in 2030
abstract
This position paper for an invited talk on the “Future of eScience” discusses the Research Software Engineering Movement and where it might be in 2030. Because of the authors' experiences, it is aimed globally but with examples that focus on the United States and United Kingdom.
Daniel S. Katz, Simon Hettrick
e-Science1
2023 Fine-grained Policy-driven I/O Sharing for Burst Buffers
abstract
A burst buffer is a common method to bridge the performance gap between the I/O needs of modern supercomputing applications and the performance of the shared file system on large-scale supercomputers. However, existing I/O sharing methods require resource isolation, offline profiling, or repeated execution that significantly limit the utilization and applicability of these systems. Here we present ThemisIO, a policy-driven I/O sharing framework for a remote-shared burst buffer: a dedicated group of I/O nodes, each with a local storage device. ThemisIO preserves high utilization by implementing opportunity fairness so that it can reallocate unused I/O resources to other applications. ThemisIO accurately and efficiently allocates I/O cycles among applications, purely based on real-time I/O behavior without requiring user-supplied information or offline-profiled application characteristics. ThemisIO supports a variety of fair sharing policies, such as user-fair, size-fair, as well as composite policies, e.g., group-then-user-fair. All these features are enabled by its statistical token design. ThemisIO can alter the execution order of incoming I/O requests based on assigned tokens to precisely balance I/O cycles between applications via time slicing, thereby enforcing processing isolation. Experiments using I/O benchmarks show that ThemisIO sustains 13.5--13.7% higher I/O throughput and 19.5--40.4% lower performance variation than existing algorithms. For real applications, ThemisIO significantly reduces the slowdown by 59.1--99.8% caused by I/O interference.
Ed Karrels, Lei Huang 0019, Yuhong Kan, Ishank Arora, Yinzhi Wang, Daniel S. Katz, William Gropp, Zhao Zhang 0007
SC6
2022 $f$funcX: Federated Function as a Service for Science
abstract
ƒuncX is a distributed function as a service (FaaS) platform that enables flexible, scalable, and high performance remote function execution. Unlike centralized FaaS systems, ƒuncX decouples the cloud-hosted management functionality from the edge-hosted execution functionality. ƒuncX's endpoint software can be deployed, by users or administrators, on arbitrary laptops, clouds, clusters, and supercomputers, in effect turning them into function serving systems. ƒuncX's cloud-hosted service provides a single location for registering, sharing, and managing both functions and endpoints. It allows for transparent, secure, and reliable function execution across the federated ecosystem of endpoints—enabling users to route functions to endpoints based on specific needs. ƒuncX uses containers (e.g., Docker, Singularity, and Shifter) to provide common execution environments across endpoints. ƒuncX implements various container management strategies to execute functions with high performance and efficiency on diverse ƒuncX endpoints. ƒuncX also integrates with an in-memory data store and Globus for managing data that may span endpoints. We motivate the need for ƒuncX, present our prototype design and implementation, and demonstrate, via experiments on two supercomputers, that ƒuncX can scale to more than 130000 concurrent workers. We show that ƒuncX's container warming-aware routing algorithm can reduce the completion time for 3,000 functions by up to 61% compared to a randomized algorithm and the in-memory data store can speed up data transfers by up to 3x compared to a shared file system.
Zhuozhao Li, Ryan Chard, Yadu N. Babuji, Ben Galewsky, Tyler J. Skluzacek, Kirill Nagaitsev, Anna Woodard, Ben Blaiszik, Josh Bryan, Daniel S. Katz, Ian T. Foster, Kyle Chard
IEEE Trans. Parallel Distributed Syst.10
2021 Federated Function as a Service for eScience
abstract
The function as a service paradigm aims to abstract the complexities of managing computing infrastructure for users. While adoption in industry has been swift, we have yet to see widespread adoption in academia. This is in part due to barriers such as the need to access large research data, diverse hardware requirements, monolithic code bases, and existing systems available to researchers. We describe funcX, a federated functionas-a-service platform that addresses important requirements for use of FaaS in research computing. We outline how funcX has been used in early science deployments.
Yadu N. Babuji, Josh Bryan, Ryan Chard, Kyle Chard, Ian T. Foster, Ben Galewsky, Daniel S. Katz, Zhuozhao Li
e-Science7
2021 Extreme Scale Survey Simulation with Python Workflows
abstract
The Vera C. Rubin Observatory Legacy Survey of Space and Time (LSST) will soon carry out an unprecedented wide, fast, and deep survey of the sky in multiple optical bands. The data from LSST will open up a new discovery space in astronomy and cosmology, simultaneously providing clues toward addressing burning issues of the day, such as the origin of dark energy and and the nature of dark matter, while at the same time yielding data that will, in turn, pose fresh new questions. To prepare for the imminent arrival of this remarkable data set, it is crucial that the associated scientific communities be able to develop the software needed to analyze it. Computational power now available allows us to generate synthetic data sets that can be used as a realistic training ground for such an effort. This effort raises its own challenges—the need to generate very large simulations of the night sky, scaling up simulation campaigns to large numbers of compute nodes across multiple computing centers with different architectures, and optimizing the complex workload around memory requirements and widely varying wall clock times. We describe here a large-scale workflow that melds together Python code to steer the workflow, Parsl to manage the large-scale distributed execution of workflow components, and containers to carry out the image simulation campaign across multiple sites. Taking advantage of these tools, we developed an extreme-scale computational framework and used it to simulate five years of observations for 300 square degrees of sky area. We describe our experiences and lessons learned in developing this workflow capability, and highlight how the scalability and portability of our approach enabled us to efficiently execute it on up to 4000 compute nodes on two supercomputers.
A. S. Villarreal, Yadu N. Babuji, Thomas D. Uram, Daniel S. Katz, Kyle Chard, Katrin Heitmann
e-Science4
2021 Understanding the multifaceted geospatial software ecosystem: a survey approach
abstract
Rebecca C. Vandewalleab , William C. Barleyc , Anand Padmanabhanabd, Daniel S. Katzd & Shaowen Wangab* a Department of Geography and Geographic Information Science, University of Illinois at Urbana-Champaign, Urbana, IL, USAb CyberGIS Center for Advanced Digital and Spatial Studies, University of Illinois at Urbana-Champaign, Urbana, IL, USAc Department of Communication, University of Illinois at Urbana-Champaign, Urbana, IL, USAd National Center for Supercomputing Applications, University of Illinois at Urbana-Champaign, Urbana, IL, USARebecca Vandewalle is a PhD student at the University of Illinois at Urbana-Champaign. Her research interests include spatially-explicit agent-based modeling, spatial network analysis, coupled human and natural systems in emergency contexts, and cyberGIS.William C. Barley is an Assistant Professor in the Department of Communication at the University of Illinois Urbana-Champaign. His research interests include organizational communication, collaboration and coordination, data representation, and field studies of technology design, adoption, and use.Anand Padmanabhan is a Research Associate Professor in the Department of Geography and Geographic Information Science at the University of Illinois Urbana-Champaign. His research interests include distributed systems, cyberinfrastructure, and cyberGIS.Daniel S. Katz is Assistant Director for Scientific Software and Applications at the National Center for Supercomputing Applications and Research Associate Professor in Computer Science, Electrical and Computer Engineering, and the School of Information Sciences at the University of Illinois Urbana-Champaign. His research interests include the interaction of people and software.Shaowen Wang is a Professor and Head of the Department of Geography and Geographic Information Science; and an Affiliate Professor of the Department of Computer Science, Department of Urban and Regional Planning, and School of Information Sciences at the University of Illinois at Urbana-Champaign. His research interests include geographic information science and systems (GIS), advanced cyberinfrastructure and cyberGIS, complex environmental and geospatial problems, computational and data sciences, high-performance and distributed computing, and spatial analysis and modeling.CONTACT Shaowen Wang [email protected] the characteristics of the rapidly evolving geospatial software ecosystem in the United States is critical to enable convergence research and education that are dependent on geospatial data and software. This paper describes a survey approach to better understand geospatial use cases, software and tools, and limitations encountered while using and developing geospatial software. The survey was broadcast through a variety of geospatial-related academic mailing lists and listservs. We report both quantitative responses and qualitative insights. As 42% of respondents indicated that they viewed their work as limited by inadequacies in geospatial software, ample room for improvement exists. In general, respondents expressed concerns about steep learning curves and insufficient time for mastering geospatial software, and often limited access to high-performance computing resources. If adequate efforts were taken to resolve software limitations, respondents believed they would be able to better handle big data, cover broader study areas, integrate more types of data, and pursue new research. Insights gained from this survey play an important role in supporting the conceptualization of a national geospatial software institute in the United States with the aim to drastically advance the geospatial software ecosystem to enable broad and significant research and education advances.
Rebecca Vandewalle, William C. Barley, Anand Padmanabhan, Daniel S. Katz, Shaowen Wang 0001
Int. J. Geogr. Inf. Sci.4
2019 Quantifying the Impact of Memory Errors in Deep Learning
abstract
The use of deep learning (DL) on HPC resources has become common as scientists explore and exploit DL methods to solve domain problems. On the other hand, in the coming exascale computing era, a high error rate is expected to be problematic for most HPC applications. The impact of errors on DL applications, especially DL training, remains unclear given their stochastic nature. In this paper, we focus on understanding DL training applications on HPC in the presence of silent data corruption. Specifically, we design and perform a quantification study with three representative applications by manually injecting silent data corruption errors (SDCs) across the design space and compare training results with the error-free baseline. The results show only 0.61-1.76% of SDCs cause training failures, and taking the SDC rate in modern hardware into account, the actual chance of a failure is one in thousands to millions of executions. With this quantitatively measured impact, computing centers can make rational design decisions based on their application portfolio, the acceptable failure rate, and financial constraints; for example, they might determine their confidence in the correctness of training results performed on processors without error correction code (ECC) RAM. We also discover that over 75-90% of the SDCs that cause catastrophic errors can be easily detected by a training loss in the next iteration. Thus we propose this error-aware software solution to correct catastrophic errors, as it has significantly lower time and space overhead compared to algorithm-based fault-tolerance (ABFT) and ECC.
Zhao Zhang 0007, Lei Huang 0019, Ruizhu Huang, Weijia Xu, Daniel S. Katz
CLUSTER5
2019 Parsl: Pervasive Parallel Programming in Python
abstract
High-level programming languages such as Python are increasingly used to provide intuitive interfaces to libraries written in lower-level languages and for assembling applications from various components. This migration towards orchestration rather than implementation, coupled with the growing need for parallel computing (e.g., due to big data and the end of Moore's law), necessitates rethinking how parallelism is expressed in programs. Here, we present Parsl, a parallel scripting library that augments Python with simple, scalable, and flexible constructs for encoding parallelism. These constructs allow Parsl to construct a dynamic dependency graph of components that it can then execute efficiently on one or many processors. Parsl is designed for scalability, with an extensible set of executors tailored to different use cases, such as low-latency, high-throughput, or extreme-scale execution. We show, via experiments on the Blue Waters supercomputer, that Parsl executors can allow Python scripts to execute components with as little as 5 ms of overhead, scale to more than 250000 workers across more than 8000 nodes, and process upward of 1200 tasks per second. Other Parsl features simplify the construction and execution of composite programs by supporting elastic provisioning and scaling of infrastructure, fault-tolerant execution, and integrated wide-area data management. We show that these capabilities satisfy the needs of many-task, interactive, online, and machine learning applications in fields such as biology, cosmology, and materials science.
Yadu N. Babuji, Anna Woodard, Zhuozhao Li, Daniel S. Katz, Ben Clifford, Lukasz Lacinski, Ryan Chard, Justin M. Wozniak, Ian T. Foster, Michael Wilde, Kyle Chard
HPDC4
2019 The global impact of science gateways, virtual research environments and virtual laboratories
Michelle Barker, Sílvia Delgado Olabarriaga, Nancy Wilkins-Diehr, Sandra Gesing, Daniel S. Katz, Shayan Shahand, Scott Henwood, Tristan Glatard, Keith G. Jeffery, Brian Corrie, Andrew E. Treloar, Helen M. Glaves, Lesley Wyborn, Neil P. Chue Hong, Alessandro Costa
Future Gener. Comput. Syst.5
2018 Building a Sustainable Structure for Research Software Engineering Activities
abstract
The profile of research software engineering has been greatly enhanced by developments at institutions around the world to form groups and communities that can support effective, sustainable development of research software. We observe, however, that there is still a long way to go to build a clear understanding about what approaches provide the best support for research software developers in different contexts, and how such understanding can be used to suggest more formal structures, models or frameworks that can help to further support the growth of research software engineering. This short paper provides an overview of some preliminary thoughts and proposes an initial high-level framework based on discussions between the authors around the concept of a set of pillars representing key activities and processes that form the core structure of a successful research software engineering offering.
Jeremy Cohen 0002, Daniel S. Katz, Michelle Barker, Robert Haines, Neil P. Chue Hong
eScience2
2018 Mapping the Research Software Sustainability Space
abstract
A growing number of largely uncoordinated initiatives focus on research software sustainability. A comprehensive mapping of the research software sustainability space can help identify gaps in their efforts, track results, and avoid duplication of work. To this end, this paper suggests enhancing an existing schematic of activities in research software sustainability, and formalizing it in a directed graph model. Such a model can be further used to define a classification schema which, applied to research results in the field, can drive the identification of past activities and the planning of future efforts.
Stephan Druskat, Daniel S. Katz
eScience2
2017 BOSS-LDG: A Novel Computational Framework That Brings Together Blue Waters, Open Science Grid, Shifter and the LIGO Data Grid to Accelerate Gravitational Wave Discovery
abstract
We present a novel computational framework that connects Blue Waters, the NSF-supported, leadership-class supercomputer operated by NCSA, to the Laser Interferometer Gravitational-Wave Observatory (LIGO) Data Grid via Open Science Grid technology. To enable this computational infrastructure, we configured, for the first time, a LIGO Data Grid Tier-1 Center that can submit heterogeneous LIGO workflows using Open Science Grid facilities. In order to enable a seamless connection between the LIGO Data Grid and Blue Waters via Open Science Grid, we utilize Shifter to containerize LIGO's workflow software. This work represents the first time Open Science Grid, Shifter, and Blue Waters are unified to tackle a scientific problem and, in particular, it is the first time a framework of this nature is used in the context of large scale gravitational wave data analysis. This new framework has been used in the last several weeks of LIGO's second discovery campaign to run the most computationally demanding gravitational wave search workflows on Blue Waters, and accelerate discovery in the emergent field of gravitational wave astrophysics. We discuss the implications of this novel framework for a wider ecosystem of Higher Performance Computing users.
Eliu A. Huerta, Roland Haas, Edgar Fajardo Hernandez, Daniel S. Katz, Peter Couvares, Josh Willis 0002, Timothy Bouvet, Jeremy Enos, William T. Kramer, Hon Wai Leong, David Wheeler
eScience4
2017 Understanding Software in Research: Initial Results from Examining Nature and a Call for Collaboration
abstract
This lightning talk paper discusses an initial data set that has been gathered to understand the use of software in research, and is intended to spark wider interest in gathering more data. The initial data analyzes three months of articles in the journal Nature for software mentions. The wider activity that we seek is a community effort to analyze a wider set of articles, including both a longer timespan of Nature articles as well as articles in other journals. Such a collection of data could be used to understand how the role of software has changed over time and how it varies across fields.
Udit Nangia, Daniel S. Katz
eScience2
2017 Evaluating Distributed Execution of Workloads
abstract
Resource selection and task placement for distributed execution poses conceptual and implementation difficulties. Although resource selection and task placement are at the core of many tools and workflow systems, the methods are ad hoc rather than being based on models. Consequently, partial and non-interoperable implementations proliferate. We address both the conceptual and implementation difficulties by experimentally characterizing diverse modalities of resource selection and task placement. We compare the architectures and capabilities of two systems: the AIMES middleware and Swift workflow scripting language and runtime. We integrate these systems to enable the distributed execution of Swift workflows on Pilot-Jobs managed by the AIMES middleware. Our experiments characterize and compare alternative execution strategies by measuring the time to completion of heterogeneous uncoupled workloads executed at diverse scale and on multiple resources. We measure the adverse effects of pilot fragmentation and early binding of tasks to resources and the benefits of backfill scheduling across pilots on multiple resources. We then use this insight to execute a multi-stage workflow across five production-grade resources. We discuss the importance and implications for other tools and workflow systems
Matteo Turilli, Yadu N. Babuji, André Merzky, Ming Tai Ha, Michael Wilde, Daniel S. Katz, Shantenu Jha
eScience6
2017 A social content delivery network for e-Science
abstract
Summary We are in the midst of a scientific data explosion in which the rate of data growth is rapidly increasing. While large‐scale research projects have developed sophisticated data distribution networks to share their data with researchers globally, there is no such support for the many millions of research projects generating data of interest to much smaller audiences (as exemplified by the long tail scientist). In data‐oriented research, every aspect of the research process is influenced by data access. However, sharing and accessing data efficiently as well as lowering access barriers are difficult. In the absence of dedicated large‐scale storage, many have noted that there is an enormous storage capacity available via connected peers, none more so than the storage resources of many research groups. With widespread usage of the content delivery network model for disseminating web content, we believe a similar model can be applied to distributing, sharing, and accessing long tail research data in an e‐Science context. We describe the vision and architecture of a social content delivery network – a model that leverages the social networks of researchers to automatically share and replicate data on peers' resources based upon shared interests and trust. Using this model, we describe a simulator and investigate how aspects such as user activity, geographic distribution, trust, and replica selection algorithms affect data access and storage performance. From these results, we show that socially informed replication strategies are comparable with more general strategies in terms of availability and outperform them in terms of spatial efficiency. Copyright © 2016 John Wiley & Sons, Ltd.
Kyle Chard, Simon Caton, Kai Kugler 0002, Omer F. Rana, Daniel S. Katz
Concurr. Comput. Pract. Exp.5
2017 Introducing distributed dynamic data-intensive (D3) science: Understanding applications and infrastructure
abstract
Summary A common feature across many science and engineering applications is the amount and diversity of data and computation that must be integrated to yield insights. Datasets are growing larger and becoming distributed; their location, availability, and properties are often time‐dependent. Collectively, these characteristics give rise to dynamic distributed data‐intensive applications. While “static” data applications have received significant attention, the characteristics, requirements, and software systems for the analysis of large volumes of dynamic, distributed data, and data‐intensive applications have received relatively less attention. This paper surveys several representative dynamic distributed data‐intensive application scenarios, provides a common conceptual framework to understand them, and examines the infrastructure used in support of applications.
Shantenu Jha, Daniel S. Katz, André Luckow, Neil P. Chue Hong, Omer F. Rana, Yogesh L. Simmhan
Concurr. Comput. Pract. Exp.2
2017 Report on the first workshop on negative and null results in eScience
abstract
New techniques and technologies, such as the use of large-scale computing, influence research approaches, methods, and scales and are rapidly changing the scientific landscape. Research projects in eScience ∗ thus start with many assumptions and many unknowns and are often complex. While the scientific process is sometimes viewed, at least in hindsight, as a linear progression from one good idea to the next, it is in fact fraught with false starts, wrong assumptions, and dead ends. The increasing reliance on computation adds to the scope of problems that occur. Researchers invest a significant amount of time and effort in their research. Funding agencies similarly make large investments to support such research, on the assumption that most of the research will be successful. When the research assumptions and hypotheses turn out to be false, causing results that are "negative" or "null", the natural bias is to judge that the research project "failed." The history of science, however, shows that negative results may be an opportunity to revolutionize a field of study. For example, Fleming noticed that his flu cultures were contaminated by mold, but that there was infection around that mold, leading to his discovery of Penicillin. Similarly, a project today may fail because of the misuse or failure of computational support. Such "failures" actually indicate that there is an opportunity for the cyberinfrastructure research community to improve computing resources and tools. The interaction of these modes of failure is multi-faceted. Negative results have been difficult to find in published papers in all scientific domains. We identify three reasons for this. First, negative results may not be identified as such but simply considered mistakes. Such cases may never be investigated further. Secondly, paper referees may demand a higher standard from such results, because they are more difficult to understand or challenge the conventional narrative. Third, researchers may self-select against publishing such results in light of the previous point. This paper contributes to the discussion about null or negative results in eScience. It also attempts to organize concepts about negative or null results in eScience in the form of a taxonomy. Falsifiability is the concept that a given statement can be refuted by a real-world measurement or observation 4. The empirical sciences are dominated by the construction of such statements and efforts to confirm or refute them. In eScience, such statements are rarely formally presented in a refutable manner. eScience projects typically merge goals from the physical science with computer science and engineering aspects. A failure in eScience may often be attributed to a computer engineering failure (software defects or unresolved performance shortcomings) or a collaboration misfit (the groups never came together). However, many important statements are never answered definitely, such as whether a given computational approach is effective for the physical science investigation. The formalization and confirmation/refutation of such statements have the potential to prevent efforts lost to engineering aspects. Post-mortem analysis of failed experiments provides "clues suggesting deeper lying forces," as Galison notes in How Experiments End, "Any historical reconstruction that ignores what seems in retrospect to be erroneous will be an inadequate account" 5. This means that the study of errors is not only relevant to students of history, as these "forces" can guide future investigations, suggest fundamental problems in experimental approaches, or even challenge prevailing theories. For example, relational database systems have been a well-accepted solution for information structuring, storage, and retrieval. This model is now being challenged by other database concepts, largely motivated by the need to cope with increasing data volumes. While the boundary between failure and success is not sharp in transitioning from relational to noSQL databases, the transition demonstrates a need to adapt and improve. This is often the normal path in research in computer science and cyberinfrastructure in particular, which could learn a lot from the various "failures" in eScience projects. Similarly, "short-term" examples of such "failures" are abound, such as the limit of being able to be a part of at most 16 Unix groups in NFS. It is likely that someone architecting an eScience collaboration system will face this limitation rather quickly. Are negative results as valuable as positive results in general? As Ayer points out, "What justifies scientific procedure ... is the success of the predictions to which it gives rise" 6. Following this line of thought, negative results are subordinate to the positive results that validate useful predictions. A negative result invalidates a previously held prediction, challenging or demolishing a theory or model. It is, however, incomplete. A negative result is an opportunity to pick up the pieces and fix the theory. Negative results are thus an important reminder of the limitations of science at any given point in time. "We forget about unpredictability when it is our turn to predict," Taleb says in The Black Swan 7, a book that attempts to analyze tumultuous events, including several scientific cases. Taleb makes the case that studying such cognitive upheavals is worthwhile in its own right. Professionals who act with the history of failed ideas and efforts in mind will be more resilient against similar changes in the future. Taleb describes the social and mental impact of experiencing (repeated) failure, indicating that without support, researchers can easily become demoralized and shy away from challenging, long-term problems. However, he notes that "Your finding nothing is very valuable... —hey, you know where not to look" 7. Venues such as the ERROR workshop are intended to encourage discussion of specific negative results. By co-locating with the eScience conference in Munich, the workshop attracted significant attention from this scientific community, with about 20 participants. The workshop accepted four papers out of six submitted after a peer review process. Each paper was reviewed by two to three members of the program committee, who evaluated the works based on originality, scientific rigor, significance, and presentation. The accepted papers were presented orally at the workshop, followed by a panel discussion on the topic "theory versus practice in eScience: gaps and gaping holes." The first presentation, by Gomes et al. 9, considered problems regarding interoperability between scientific workflows. Specifically, they discussed the problem of reusing workflows previously developed and implemented using one particular scientific workflow management framework with another one. To solve this problem, the authors developed an "intermediate" workflow language, with the idea that this intermediate language would preserve the workflow's semantic information across frameworks. However, they observed a loss of information about workflow semantics during the translation from a first workflow language to the "neutral" language and from the "neutral" language to the second workflow language. This happens because there are no ideal or standardized semantics for workflow languages, which is the key negative result in this research. A solution proposed to this problem is through the adoption of workflow patterns to describe richer workflow semantics. The second presentation, by Groen and Portgies Zwart 10, provided a high level overview of the authors' experience in constructing a distributed supercomputing system: CosmoGrid. The authors discussed how ambitious ideas can often be stymied by site-local resource allocation decisions. One of the negative results is the conclusion that harnessing multiple large machines is not feasible and therefore one should focus on harnessing a larger number of smaller machines. Additionally, the authors pointed out that a task as simple as getting software installed is significantly difficult at major computational sites, exposing the often overlooked reality of working with large scale computational infrastructures. The third presentation, by Cebrian et al. 11, presented an experience in designing two separate cache stores—for private and shared data— for multicore system architectures. The premise of the work was to improve efficiency by excluding private and shared read-only cache contents from coherency management. From the experiments and analysis of results obtained with this approach, the authors concluded that systems are less efficient with this kind of design, which is a negative result. This is because the overhead of classification mechanisms and increased concentration of access to shared data cause a bandwidth bottleneck to a particular portion of cache, resulting in higher latencies. The fourth presentation, by Jackson et al. 12, discussed an experience of performing an experiment in the context of a larger body of work. The experiment was on latency measurement between nodes within a single cluster, as well as across different clusters. Some of the main takeaways from the experiments as described by the authors were the technical and administrative obstacles faced when the experiment involves dependencies on several independently managed computational systems across administrative boundaries. Despite these obstacles, the authors were able to produce a significant body of latency data. One finding from this dataset was that an exhaustive study of latencies among systems was not necessarily, by itself, a good predictor of actual application performance. The topic of discussion for the panel session was "theory versus practice in eScience: gaps and gaping holes." This theme emerged from a recurring observation in the submitted papers, in which negative or null results are attributed to a mismatch between expectations, which are based on theory, and what is found in reality. The panelists were Daniel S. Katz, Simon Portegies Zwart, Kyle Chard, Juan M. Cebrián, and Gary Jackson. Each panelist spoke for 2–3 min, and then there was an open discussion between panelists and the audience. The rest of this subsection presents the highlights of the panelists talks and the discussion that followed. Jackson spoke about the importance of not losing research focus because of infrastructure complexities and problems—the proverbial "missing the forest for trees." Chard said that there are no good definitions of eScience, although we provide one taken from the eScience conference series website in Section 2. Furthermore, he raised questions as to how the scientific process, which had been mostly unchanged for hundreds of years and has long review cycles, has recently been changing with online data publication and open access journals. Katz spoke of the phenomenon of failures among research projects and endeavors by quoting from the opening of Tolstoy's Anna Karenina, "Happy families are all alike; every unhappy family is unhappy in its own way." He noted that successful research endeavors must have all their critical factors right to be successful and failing even one of them could jeopardize the complete project. This is popularly known as Anna Karenina Principle 13. Katz also emphasized the importance and value of scientific results in general and negative results in particular with the question/statement: "How do we decide if there is value in a result?" Portgies Zwart spoke about the importance of understanding the difference between core computer science and other sciences, as well as the scientists associated with each one. He argued that computer science is currently undergoing a crisis because it is hard to find interesting problems, because of competition with the industry. In particular, computer science is challenged by reproducibility. One solution, he suggested, is an establishment of a software museum to prevent loss of software. The open discussion that followed focused on diverse topics such as software preservation, publication and credit, training, and the definition of negative or null results. It began with participants expressing concerns about issues related to software in particular. Scenarios were discussed that introduce the "gaps between eScience theory and practice" connecting technologies, ideas, and people. One gap is that digital products in general and software in particular, including methods and knowledge (algorithms), can be lost over time, sometimes known as bit rot 14. Software hosting services such as GitHub can address this problem to a certain extent by preserving the files, but they still require much human effort to preserve the function delivered by the software as meaningful. A curation service for algorithms could be another solution, requiring additional effort. Can the software and algorithm hosting services be linked as concepts and implementation? Preservation and publication of negative results are a challenge. While there are no technical barriers, from a publishing culture point of view, there are few or no incentives for publishing negative results. In the presence of such incentives, people would develop the culture about explaining not only what they did but also why they did not do so in some other way. And, if the negative results are actually published, it is likely that the same approach will not be taken by other researchers and groups. Another identified gap is the lack of a comprehensive understanding of negative results because of the lack of a conceptual framework, for example, a taxonomy. The discussion also raised the gap introduced by a lack of a credit model for discovering, identifying, and reporting negative results. For example, can negative results and/or methods to obtain them be patented? For instance, who receives credit if a succession of graduate students working on a problem arrive at a negative result followed by positive result? Or what happens to the positive results obtained before further investigation leads to their negation and nullification? Another gap arises from the lack of proper eScience training of domain scientists—can domain scientists be trained to become eScientists? This also applies to principal investigators, many of whom were trained in an era when science was carried out differently than it is carried out today. Training imparting the knowledge of modern computational methods and capabilities could play a key role in filling such a gap. The difference between incomplete (such as obtained from samples of insufficient size) and negative results can also be unclear. This can often result in negative results that are subject to interpretation. For instance, in an MD simulation 14, it cannot be shown if sampling was sufficient. In the same vein, should incremental competitive results be considered negative? In general, there is no standard on how many simulated timesteps are needed to obtain the correct answer, although sometimes, one can validate against lab experiments. Another topic raised during the discussion compared research in academia versus science in the commercial sector. One prominent sentiment expressed in the discussion was that, in some areas, research done in the commercial domain is "ahead" of research in academia, particularly, where industry has larger-scale problems and data than academia. One possible reason for this could be that there are more negative and null results in academia compared to industry. However, academic research can be transferred to the commercial domain and vice versa. For instance, the patent system is in place to enable commercial contribution to public research. One question that arises here is should there be a distinction between commercial and academic research? How will this distinction manifest itself? It was noted that there is an asymmetry between positive and negative results: With negative results, it is more likely that some error was made. It may be harder to truly prove a negative result. For instance an "existence proof" is sufficient for a positive result while a "non-possible proof" is needed for negative results. The workshop led to concrete outcomes before, during, and as a follow-up of its realization. After the announcement of the workshop, the Mozilla Science Foundation hosted a guest post about the workshop by the organizers 15. The post discussed the importance of the theme of "negative" results and the goals for the workshop. One of the outcomes of the panel discussion was the call for a taxonomy of negative and null results in eScience. We respond to this call in this paper by proposing a taxonomy in Section 5. In Figure 1, we present a taxonomy of eScience results. The three kinds of results in eScience are positive, null, and negative. In the taxonomy, we focus on negative and null results. Negative results may be caused by one or more of the following reasons: technological, technical, human, and domain. For instance, an erroneous result obtained because of insufficient precision resulting from a limitation of a system library is an example of negative result caused by technical and technological limitations. Similarly, a simulation algorithm resulting from a flawed understanding of a natural phenomenon could be considered a negative result caused by domain and human factors. A mismatch between the problem/data size and the technology/methodology used is an example of a technological reason. Examples of technical causes include software bugs and cyberinfrastructure faults. Human causes include both incidental issues such as mistakes in measurements and systemic issues such as a false hypothesis or an insufficient sample size. Null results are obtained because of the lack of discriminating conditions to confirm or refute a hypothesis. Such situation may be caused by similar reasons as for negative results. An example of technological/technical reason is a statistical test has poor performance on the data because the implementation uses limited precision. Null results can also have a human cause, when insufficient samples are used in the experiments or when some bias in the data goes unnoticed. We are aware of two workshops with similar themes in related fields (Information and Communication Technologies). The first is NoISE (Workshop on Negative or Inconclusive Results in Semantic Web) 19. The second is NOPE (Workshop on Negative Outcomes, Post-mortems, and Experiences) 20. Both the workshops were organized for the first time in 2015. Similarly to the current special issue, there have been two special issues in the prominent journals focused on negative or null results: the Journal on Negative results in Empirical Software Engineering 21 and PLOS ONE Collection 22. The PLOS ONE Collection focuses on inconclusive results as a distinct type of negative results in addition to null results. Similarly to the workshops, both special issues were launched for the first time in 2015. In this section, we discuss some of the key implications that are drawn from the previous sections. These are the issues that are directly impacted by the occurrence of negative and null results in the research as conducted by the scientific community. We classify these implications into four categories: technological, technical, cultural, and domain specific. Publication, credit, and citation of the work that has yielded negative results are an important consideration from the research community point of view. Citations and credit are important measures of success for a research publication. Given the current trends of publishing positive results, it is a crucial decision for a researcher to invest efforts in publishing a negative result. Technical issues such as hardware faults and software bugs often go undetected until late in the research work. In these cases, the negative results are not necessarily of the same nature as the science domain unless the domain is computer science itself. It becomes difficult for a domain scientist to draw value from the publication and dissemination of such results. As a consequence, they often are ignored or fixed after the results were obtained and disseminated. For example, a bug in the third party library call up the toolchain of an application that limited the results of the actual science in scale or precision can be considered a negative result. In some cases, the problems are mismatched to the available computational infrastructure. Sometimes, the problems are too small for a given environment, leading to an inefficient use. In others cases, the problems are too large for the infrastructure, leading to generation of incomplete or no results. Publishing details about such cases can benefit the community by allowing it to better understand how to more optimally match problems and solutions. Finally, a cumulative effect from more than one cause is possible. One of the biggest concern about such issues is that they often go unnoticed by the larger community and hence the appropriate correction measures are not adapted. We are grateful to the program committee members, the panelists, and the authors who supported the realization of this first workshop. We also thank the reviewers of this special issue. The work by Katz was supported in part by the National Science Foundation while working at the Foundation. Any opinion, finding, and conclusions or recommendations expressed in this material are those of the author(s) and do not necessarily reflect the views of the National Science Foundation.
Ketan Maheshwari, Daniel S. Katz, Sílvia Delgado Olabarriaga, Justin M. Wozniak, Douglas Thain
Concurr. Comput. Pract. Exp.2
2017 Leading-edge research in cluster, cloud, and grid computing: Best papers from the IEEE/ACM CCGrid 2015 conference
Daniel S. Katz
Future Gener. Comput. Syst.1
2016 Integrating Abstractions to Enhance the Execution of Distributed Applications
abstract
One of the factors that limits the scale, performance, and sophistication of distributed applications is the difficulty of concurrently executing them on multiple distributed computing resources. In part, this is due to a poor understanding of the general properties and performance of the coupling between applications and dynamic resources. This paper addresses this issue by integrating abstractions representing distributed applications, resources, and execution processes into a pilot-based middleware. The middleware provides a platform that can specify distributed applications, execute them on multiple resource and for different configurations, and is instrumented to support investigative analysis. We analyzed the execution of distributed applications using experiments that measure the benefits of using multiple resources, the late-binding of scheduling decisions, and the use of backfill scheduling.
Matteo Turilli, Zhao Zhang 0007, André Merzky, Michael Wilde, Jon B. Weissman, Daniel S. Katz, Shantenu Jha
IPDPS7
2016 Application skeletons: Construction and use in eScience
Daniel S. Katz, André Merzky, Zhao Zhang 0007, Shantenu Jha
Future Gener. Comput. Syst.1
2016 eScience today and tomorrow
Claudia Bauzer Medeiros, Daniel S. Katz
Future Gener. Comput. Syst.2
2016 eScience today and tomorrow - Part 2
Claudia Bauzer Medeiros, Daniel S. Katz
Future Gener. Comput. Syst.2
2016 Ten Simple Rules for Taking Advantage of Git and GitHub
abstract
Bioinformatics is a broad discipline in which one common denominator is the need to produce and/or use software that can be applied to biological data in different contexts. To enable and ensure the replicability and traceability of scientific claims, it is essential that the scientific publication, the corresponding datasets, and the data analysis are made publicly available [1,2]. All software used for the analysis should be either carefully documented (e.g., for commercial software) or, better yet, openly shared and directly accessible to others [3,4]. The rise of openly available software and source code alongside concomitant collaborative development is facilitated by the existence of several code repository services such as SourceForge, Bitbucket, GitLab, and GitHub, among others. These resources are also essential for collaborative software projects because they enable the organization and sharing of programming tasks between different remote contributors. Here, we introduce the main features of GitHub, a popular web-based platform that offers a free and integrated environment for hosting the source code, documentation, and project-related web content for open-source projects. GitHub also offers paid plans for private repositories (see Box 1) for individuals and businesses as well as free plans including private repositories for research and educational use.
Yasset Pérez-Riverol, Laurent Gatto, Timo Sachsenberg, Julian Uszkoreit, Felipe da Veiga Leprevost, Christian Fufezan, Tobias Ternent, Stephen J. Eglen, Daniel S. Katz, Tom J. Pollard, Olexandr Konovalov, Robert M. Flight, Kai Blin, Juan Antonio Vizcaíno
PLoS Comput. Biol.10
2015 Toward Interlanguage Parallel Scripting for Distributed-Memory Scientific Computing
abstract
Scripting languages such as Python and R have been widely adopted as tools for the productive development of scientific software because of the power and expressiveness of the languages and available libraries. However, deploying scripted applications on large-scale parallel computer systems such as the IBM Blue Gene/Q or Cray XE6 is a challenge because of issues including operating system limitations, interoperability challenges, parallel filesystem overheads due to the small file system accesses common in scripted approaches, and other issues. We present here a new approach to these problems in which the Swift scripting system is used to integrate high-level scripts written in Python, R, and Tcl, with native code developed in C, C++, and Fortran, by linking Swift to the library interfaces to the script interpreters. In this approach, Swift handles data management, movement, and marshaling among distributed-memory processes without direct user manipulation of low-level communication libraries such as MPI. We present a technique to efficiently launch scripted applications on large-scale supercomputers using a hierarchical programming model.
Justin M. Wozniak, Timothy G. Armstrong, Ketan Maheshwari, Daniel S. Katz, Michael Wilde, Ian T. Foster
CLUSTER4
2015 Porting Ordinary Applications to Blue Gene/Q Supercomputers
abstract
Efficiently porting ordinary applications to Blue Gene/Q supercomputers is a significant challenge. Codes are often originally developed without considering advanced architectures and related tool chains. Science needs frequently lead users to want to run large numbers of relatively small jobs (often called many-task computing, an ensemble, or a workflow), which can conflict with supercomputer configurations. In this paper, we discuss techniques developed to execute ordinary applications over leadership class supercomputers. We use the high-performance Swift parallel scripting framework and build two workflow execution techniques -- sub-jobs and main-wrap. The sub-jobs technique, built on top of the IBM Blue Gene/Q resource manager Cobalt's sub-block jobs, lets users submit multiple, independent, repeated smaller jobs within a single larger resource block. The main-wrap technique is a scheme that enables C/C++ programs to be defined as functions that are wrapped by a high-performance Swift wrapper and that are invoked as a Swift script. We discuss the needs, benefits, technicalities, and current limitations of these techniques. We further discuss the real-world science enabled by these techniques and the results obtained.
Ketan Maheshwari, Justin M. Wozniak, Timothy G. Armstrong, Daniel S. Katz, T. Andrew Binkowski, Xiaoliang Zhong, Olle Heinonen, Dmitry Karpeyev, Michael Wilde
e-Science4
2015 The Case for Workflow-Aware Storage: An Opportunity Study
Lauro Beltrão Costa, Hao Yang 0039, Emalayan Vairavanathan, Abmar Barros, Ketan Maheshwari, Gilles Fedak, Daniel S. Katz, Michael Wilde, Matei Ripeanu, Samer Al-Kiswany
J. Grid Comput.7
2014 On Replica Placement in a Social CDN for e-Science
abstract
Research data is experiencing a seemingly endless increase in both volume and production rate. At the same time, efficiently transferring, storing, and analyzing large scale research data have become major research foci. In this paper, we expand on our approach to sharing data for e-Science: a Social Content Delivery Network (S-CDN). A S-CDN leverages the social networks of researchers to automatically share data and place replicas on peers' resources based upon the premises of trust and interest in shared data. We denote a consumer of shared data as a data follower, similar to the notion of Twitter followers, except we add the element of bilateral authorization to capture a notion of trust. We describe a prototypical implementation for a S-CDN that captures an efficient asynchronous transfer mechanism for data management and replication. In addition, we study via simulation the interplay of user behavior with different replication strategies that capture social as well as more general premises for data sharing. Our results illustrate the opportunities and pitfalls of various replication and data access management strategies. Specifically, we show that socially-informed replication strategies are competitive with more general strategies in terms of availability, and outperform them in terms of spatial efficiency.
Kai Kugler 0002, Simon Caton, Kyle Chard, Daniel S. Katz
eScience4
2014 Using Application Skeletons to Improve eScience Infrastructure
abstract
Computer scientists who work on tools and systems to support eScience (a variety of parallel and distributed) applications usually use actual applications to prove that their systems will benefit science and engineering (e.g., improve application performance). Accessing and building the applications and necessary data sets can be difficult because of policy or technical issues, and it can be difficult to modify the characteristics of the applications to understand corner cases in the system design. In this paper, we present the Application Skeleton, a simple yet powerful tool to build synthetic applications that represent real applications, with runtime and I/O close to those of the real applications. This allows computer scientists to focus on the system they are building, they can work with the simpler skeleton applications and be sure that their work will also be applicable to the real applications. In addition, skeleton applications support simple reproducible system experiments since they are represented by a compact set of parameters. Our Application Skeleton tool (available as open source at https://github.com/applicationskeleton/Skeleton) currently can create easy-to-access, easy-to-build, and easy-to-run bag-of-task, (iterative) map-reduce, and (iterative) multistage workflow applications. The tasks can be serial or parallel or a mix of both. We select three representative applications (Montage, BLAST, CyberShake Postprocessing), then describe and generate skeleton applications for each. We show that the skeleton applications have identical (or close) performance to that of the real applications. We then show examples of using skeleton applications to verify system optimizations such as data caching, I/O tuning, and task scheduling, as well as the system resilience mechanism, in some cases modifying the skeleton applications to emphasize some characteristic, and thus show that using skeleton applications simplifies the process of designing, implementing, and testing these optimizations.
Zhao Zhang 0007, Daniel S. Katz
eScience2
2014 Design and evaluation of the gemtc framework for GPU-enabled many-task computing
abstract
We present the design and first performance and usability evaluation of GeMTC, a novel execution model and runtime system that enables accelerators to be programmed with many concurrent and independent tasks of potentially short or variable duration. With GeMTC, a broad class of such "many-task" applications can leverage the increasing number of accelerated and hybrid high-end computing systems. GeMTC overcomes the obstacles to using GPUs in a many-task manner by scheduling and launching independent tasks on hardware designed for SIMD-style vector processing. We demonstrate the use of a high-level MTC programming model (the Swift parallel dataflow language) to run tasks on many accelerators and thus provide a high-productivity programming model for the growing number of supercomputers that are accelerator-enabled. While still in an experimental stage, GeMTC can already support tasks of fine (subsecond) granularity and execute concurrent heterogeneous tasks on 86,000 independent GPU warps spanning 2.7M GPU threads on the Blue Waters supercomputer.
Scott J. Krieder, Justin M. Wozniak, Timothy G. Armstrong, Michael Wilde, Daniel S. Katz, Benjamin Grimmer, Ian T. Foster, Ioan Raicu
HPDC5
2014 Exploring Automatic, Online Failure Recovery for Scientific Applications at Extreme Scales
abstract
Application resilience is a key challenge that must be addressed in order to realize the exascale vision. Process/node failures, an important class of failures, are typically handled today by terminating the job and restarting it from the last stored checkpoint. This approach is not expected to scale to exascale. In this paper we present Fenix, a framework for enabling recovery from process/node/blade/cabinet failures for MPI-based parallel applications in an online (i.e., Without disrupting the job) and transparent manner. Fenix provides mechanisms for transparently capturing failures, re-spawning new processes, fixing failed communicators, restoring application state, and returning execution control back to the application. To enable automatic data recovery, Fenix relies on application-driven, diskless, implicitly coordinated check pointing. Using the S3D combustion simulation running on the Titan Cray-XK7 production system at ORNL, we experimentally demonstrate Felix's ability to tolerate high failure rates (e.g., More than one per minute) with low overhead while sustaining performance.
Marc Gamell, Daniel S. Katz, Hemanth Kolla, Jacqueline Chen, Scott Klasky, Manish Parashar
SC2
2014 Special issue on eScience infrastructure and applications
Daniel S. Katz, Zhao Zhang 0007
Future Gener. Comput. Syst.1
2013 Swift/T: Large-Scale Application Composition via Distributed-Memory Dataflow Processing
abstract
Many scientific applications are conceptually built up from independent component tasks as a parameter study, optimization, or other search. Large batches of these tasks may be executed on high-end computing systems, however, the coordination of the independent processes, their data, and their data dependencies is a significant scalability challenge. Many problems must be addressed, including load balancing, data distribution, notifications, concurrent programming, and linking to existing codes. In this work, we present Swift/T, a programming language and runtime that enables the rapid development of highly concurrent, task-parallel applications. Swift/Tis composed of several enabling technologies to address scalability challenges, offers a high-level optimizing compiler for user programming and debugging, and provides tools for binding user code in C/C++/Fortran into a logical script. In this work, we describe the Swift/T solution and present scaling results from the IBM Blue Gene/Pand Blue Gene/Q.
Justin M. Wozniak, Timothy G. Armstrong, Michael Wilde, Daniel S. Katz, Ewing L. Lusk, Ian T. Foster
CCGRID4
2013 Constructing a Social Content Delivery Network for eScience
abstract
Increases in the size of research data and the move towards citizen science, in which everyday users contribute data and analyses, have resulted in a research data deluge. Researchers must now carefully determine how to store, transfer and analyze "Big Data" in collaborative environments. This task is even more complicated when considering budget and locality constraints on data storage and access. In this paper we investigate the potential to construct a Social Content Delivery Network (S-CDN) based upon the social networks that exist between researchers. The S-CDN model builds upon the incentives of collaborative researchers within a given scientific community to address their data challenges collaboratively and in proven trusted settings. In this paper we present a prototype implementation of a S-CDN and investigate the performance of the data transfer mechanisms (using Glob us Online) and the potential cost advantages of this approach.
Kai Kugler 0002, Kyle Chard, Simon Caton, Omer F. Rana, Daniel S. Katz
e-Science5
2013 MTC envelope: defining the capability of large scale computers in the context of parallel scripting applications
Zhao Zhang 0007, Daniel S. Katz, Michael Wilde, Justin M. Wozniak, Ian T. Foster
HPDC2
2013 Swift/T: scalable data flow programming for many-task applications
abstract
Swift/T, a novel programming language implementation for highly scalable data flow programs, is presented.
Justin M. Wozniak, Timothy G. Armstrong, Michael Wilde, Daniel S. Katz, Ewing L. Lusk, Ian T. Foster
PPoPP4
2013 Parallelizing the execution of sequential scripts
abstract
Scripting is often used in science to create applications via the composition of existing programs. Parallel scripting systems allow the creation of such applications, but each system introduces the need to adopt a somewhat specialized programming model. We present an alternative scripting approach, AMFS Shell, that lets programmers express parallel scripting applications via minor extensions to existing sequential scripting languages, such as Bash, and then execute them in-memory on large-scale computers. We define a small set of commands between the scripts and a parallel scripting runtime system, so that programmers can compose their scripts in a familiar scripting language. The underlying AMFS implements both collective (fast file movement) and functional (transformation based on content) file management. Tasks are handled by AMFS's built-in execution engine. AMFS Shell is expressive enough for a wide range of applications, and the framework can run such applications efficiently on large-scale computers.
Zhao Zhang 0007, Daniel S. Katz, Timothy G. Armstrong, Justin M. Wozniak, Ian T. Foster
SC2
2013 Distributed computing practice for large-scale science and engineering applications
abstract
SUMMARY It is generally accepted that the ability to develop large‐scale distributed applications has lagged seriously behind other developments in cyberinfrastructure. In this paper, we provide insight into how such applications have been developed and an understanding of why developing applications for distributed infrastructure is hard. Our approach is unique in the sense that it is centered around half a dozen existing scientific applications; we posit that these scientific applications are representative of the characteristics, requirements, as well as the challenges of the bulk of current distributed applications on production cyberinfrastructure (such as the US TeraGrid). We provide a novel and comprehensive analysis of such distributed scientific applications. Specifically, we survey existing models and methods for large‐scale distributed applications and identify commonalities, recurring structures, patterns and abstractions. We find that there are many ad hoc solutions employed to develop and execute distributed applications, which result in a lack of generality and the inability of distributed applications to be extensible and independent of infrastructure details. In our analysis, we introduce the notion of application vectors: a novel way of understanding the structure of distributed applications. Important contributions of this paper include identifying patterns that are derived from a wide range of real distributed applications, as well as an integrated approach to analyzing applications, programming systems and patterns, resulting in the ability to provide a critical assessment of the current practice of developing, deploying and executing distributed applications. Gaps and omissions in the state of the art are identified, and directions for future research are outlined. Copyright © 2012 John Wiley & Sons, Ltd.
Shantenu Jha, Murray Cole, Daniel S. Katz, Manish Parashar, Omer F. Rana, Jon B. Weissman
Concurr. Comput. Pract. Exp.3
2013 Recent advances in e-Science
Daniel S. Katz, David Abramson 0001
Future Gener. Comput. Syst.1
2013 Turbine: A Distributed-memory Dataflow Engine for High Performance Many-task Applications
abstract
Efficiently utilizing the rapidly increasing concurrency of multi-petaflop computing systems is a significant programming challenge. One approach is to structure applications with an upper layer of many loosely coupled coarse-grained tasks, each comp
Justin M. Wozniak, Timothy G. Armstrong, Ketan Maheshwari, Ewing L. Lusk, Daniel S. Katz, Michael Wilde, Ian T. Foster
Fundam. Informaticae5
2013 JETS: Language and System Support for Many-Parallel-Task Workflows
Justin M. Wozniak, Michael Wilde, Daniel S. Katz
J. Grid Comput.3
2012 A Workflow-Aware Storage System: An Opportunity Study
abstract
This paper evaluates the potential gains a workflow-aware storage system can bring. Two observations make us believe such storage system is crucial to efficiently support workflow-based applications: First, workflows generate irregular and application-dependent data access patterns. These patterns render existing storage systems unable to harness all optimization opportunities as this often requires conflicting optimization options or even conflicting design decision at the level of the storage system. Second, when scheduling, workflow runtime engines make suboptimal decisions as they lack detailed data location information. This paper discusses the feasibility, and evaluates the potential performance benefits brought by, building a workflow-aware storage system that supports per-file access optimizations and exposes data location. To this end, this paper presents approaches to determine the application-specific data access patterns, and evaluates experimentally the performance gains of a workflow-aware storage approach. Our evaluation using synthetic benchmarks shows that a workflow-aware storage system can bring significant performance gains: up to 7× performance gain compared to the distributed storage system - MosaStore and up to 16× compared to a central, well provisioned, NFS server.
Emalayan Vairavanathan, Samer Al-Kiswany, Lauro Beltrão Costa, Zhao Zhang 0007, Daniel S. Katz, Michael Wilde, Matei Ripeanu
CCGRID5
2012 Pilot abstractions for compute, data, and network
abstract
Scientific experiments in a variety of domains are producing increasing amounts of data that need to be processed efficiently. Distributed Computing Infrastructures are increasingly important in fulfilling these large-scale computational requirements.
Mark Santcroos, Sílvia Delgado Olabarriaga, Daniel S. Katz, Shantenu Jha
eScience3
2012 Topic 1: Support Tools and Environments
Omer F. Rana, Marios D. Dikaiakos, Daniel S. Katz, Christine Morin
Euro-Par3
2012 Design and analysis of data management in scalable parallel scripting
abstract
We seek to enable efficient large-scale parallel execution of applications in which a shared filesystem abstraction is used to couple many tasks. Such parallel scripting (many-task computing, MTC) applications suffer poor performance and utilization on large parallel computers because of the volume of filesystem I/O and a lack of appropriate optimizations in the shared filesystem. Thus, we design and implement a scalable MTC data management system that uses aggregated compute node local storage for more efficient data movement strategies. We co-design the data management system with the data-aware scheduler to enable dataflow pattern identification and automatic optimization. The framework reduces the time to solution of parallel stages of an astronomy data analysis application, Montage, by 83.2% on 512 cores; decreases the time to solution of a seismology application, CyberShake, by 7.9% on 2,048 cores; and delivers BLAST performance better than mpiBLAST at various scales up to 32,768 cores, while preserving the flexibility of the original BLAST application.
Zhao Zhang 0007, Daniel S. Katz, Justin M. Wozniak, Allan Espinosa, Ian T. Foster
SC2
2011 Swift: A language for distributed parallel scripting
Michael Wilde, Mihael Hategan, Justin M. Wozniak, Ben Clifford, Daniel S. Katz, Ian T. Foster
Parallel Comput.5
2010 Distributed Systems and Algorithms
Omer F. Rana, Giandomenico Spezzano, Michael Gerndt, Daniel S. Katz
Euro-Par (1)4
2010 Global-scale distributed I/O with ParaMEDIC
abstract
Abstract Achieving high performance for distributed I/O on a wide‐area network continues to be an elusive holy grail. Despite enhancements in network hardware as well as software stacks, achieving high‐performance remains a challenge. In this paper, our worldwide team took a completely new and non‐traditional approach to distributed I/O, calledParaMEDIC: Parallel Metadata Environment for Distributed I/O and Computing, by utilizing application‐specifictransformationof data to orders of magnitude smaller metadata before performing the actual I/O. Specifically, this paper details our experiences in deploying a large‐scale system to facilitate the discovery of missing genes and constructing a genome similarity tree by encapsulating the mpiBLAST sequence‐search algorithm into ParaMEDIC. The overall project involved nine computational sites spread across the U.S. and generated more than a petabyte of data that was ‘teleported’ to a large‐scale facility in Tokyo for storage. Copyright © 2010 John Wiley & Sons, Ltd.
Pavan Balaji, Wu-chun Feng, Heshan Lin, Jeremy S. Archuleta, Satoshi Matsuoka, Andrew S. Warren, João Carlos Setubal, Ewing L. Lusk, Rajeev Thakur, Ian T. Foster, Daniel S. Katz, Shantenu Jha, K. Shinpaugh, Susan Coghlan, Daniel A. Reed
Concurr. Comput. Pract. Exp.11
2010 Special Issue: Grid Computing, High Performance and Distributed Application
abstract
In the recent decades we have witnessed a major revolution in the computer field. The major challenges posed by applications in fields of bioinformatics, earth sciences or weather forecasting, among others, have caused the proliferation of complex solutions, such as grid, cloud and highperformance computing. The common objective of all these disciplines is the sharing of hardware and software resources to provide an infrastructure in which to run efficiently these applications. Particularly, grid computing has been one of the most important computing topics in the last years. Within this context, the GADA workshop arose in 2004 as a forum for researchers in grid computing and its application to data analysis. From then until 2008, GADA became a reference conference for researchers in grid, covering also a broader set of disciplines, although grid computing continued to play a key role in the set of main topics of the conference. This special issue includes the eight best papers presented at the International Conference on Grid Computing, High Performance and Distributed Application (GADA 2008) from a total of 31 submissions with an acceptance rate of 26%. GADA 2008 was held in conjunction with the On The Move Conferences during November 2008 in Monterrey, Mexico. Each submitted paper was reviewed by three reviewers and one of the program chairs, and a total of 70 reviewers were involved in the review process of GADA 2008. A further review stage was performed to select the papers for this special issue. Topics covered by these papers include grid modelling, performance and scalability. Three different papers address the important topic of grid modelling. Branco et al. [1] describe the experience in the development of the data management middleware DQ2, used in the ATLAS experiment for the Large Hadron Collider (HLC). From this experience, they have identified an important degree of uncertainty over the behaviour of large grid infrastructures. From the analysis of this uncertainty, they propose novel modelling and simulation techniques for Data Grids. van der Aalst et al. [2] show a formal description of the grid in terms of a Colored Petri Net (CPN). This formalism can be used as a conceptual model of a grid environment, and also allows you to perform various types of analysis, including performance analysis. The model has been validated by means of experiments in a testbed grid architecture. Montes et al. [3] present a methodology based on data mining techniques for the building of a global behaviour model of the grid. This methodology deals with the complexity of grid environments, providing a simple model that can be used for administrative purposes. The validation of the model has been performed by means of real and simulated case studies.
María S. Pérez 0001, Pilar Herrero, Dennis Gannon, Daniel S. Katz
Concurr. Comput. Pract. Exp.4
2010 Special Section: Grid computing, high-performance and distributed applications
Pilar Herrero, Daniel S. Katz, María S. Pérez 0001, Domenico Talia
Future Gener. Comput. Syst.2
2009 An innovative application execution toolkit for multicluster grids
abstract
Multicluster grids provide one promising solution to satisfying growing computation demands of compute-intensive applications by collaborating various networked clusters. However, it is challenging to seamlessly integrate all participating clusters in different domains into a virtual computation platform. In order to take full advantages of multicluster grids capability, computer scientists need to deal with how to collaborate practically and efficiently participating autonomic systems to execute Grid-enabled applications. We make efforts on grid resource management and implement a toolkit called Pelecanus to improve the overall performance of application execution in multicluster grids environment. The Pelecanus takes advantages of the DA-TC (Dynamic Assignment with Task Containers) execution model to improve resource interoperability and enhance application execution and monitoring. Experiments show that it can significantly reduce turnaround time and increase resource utilization for certain applications with large number of sequential jobs.
Zhifeng Yun, Zhou Lei 0001, Gabrielle Allen, Daniel S. Katz, Tevfik Kosar, Shantenu Jha, J. Ramanujam
CLUSTER4
2009 A Fresh Perspective on Developing and Executing DAG-Based Distributed Applications: A Case-Study of SAGA-Based Montage
abstract
Most workflow based applications currently have to adapt to available tools. While this keeps the cost of development low, it can lead to performance and flexibility tradeoffs that the application developer and deployer must make. In this paper, we use the Montage astronomical image mosaicking application as prototypical DAG-based workflow application to layout the development and deployment decisions for distributed applications. We discuss and explain the lack of simple (easy-to-use), scalable, and extensible distributed applications. We then introduce SAGA as a technology that permits the construction of abstractions that aid the development and execution of the applications, and thus addresses some of common shortcomings of traditional distributed applications development. We use Montage together with SAGA to examine how legacy applications can be made to run on distributed infrastructures, to see if our reasons are valid, and to compare potential new methods for creating distributed applications with existing technologies that are currently used. We demonstrate the ability to (i) scale-out and (ii) use different production infrastructure, while maintaining performance comparable to established systems. Our hope is that by demonstrating the simplicity of development along with other advantages (performance, scalability, extensibility, and infrastructure independence), this example will encourage others to think more broadly about how distributed applications are created and how new programming models such as Dryad can be supported in an infrastructure independent way, thus eventually leading to more applications that can seamlessly scale-out.
André Merzky, Katerina Stamou, Shantenu Jha, Daniel S. Katz
eScience4
2004 Application-Level Fault Tolerance in the Orbital Thermal Imaging Spectrometer
abstract
Systems that operate in extremely volatile environments, such as orbiting satellites, must be designed with a strong emphasis on fault tolerance. Rather than rely solely on the system hardware, it may be beneficial to entrust some of the fault handling to software at the application level, which can utilize semantic information and software communication channels to achieve fault tolerance with considerably less power and performance overhead. We show the implementation and evaluation of such a software-level approach, application-level fault tolerance and detection (ALFTD) into the orbital thermal imaging spectrometer (OTIS).
E. Ciocca, Israel Koren, Zahava Koren, C. Mani Krishna 0001, Daniel S. Katz
PRDC5
2004 Accessing and Visualizing Scientific Spatiotemporal Data
Daniel S. Katz, Attila Bergou, G. Bruce Berriman, Gary L. Block, Jim Collier, David W. Curkendall, John Good, Laura Husman, Joseph C. Jacob, Anastasia C. Laity, Peggy Li, Craig Miller, Tom Prince, Herb Siegel, Roy Williams
SSDBM1
2003 Tests and Tolerances for High-Performance Software-Implemented Fault Detection
abstract
We describe and test a software approach to fault detection in common numerical algorithms. Such result checking or algorithm-based fault tolerance (ABFT) methods may be used, for example, to overcome single-event upsets in computational hardware or to detect errors in complex, high-efficiency implementations of the algorithms. Following earlier work, we use checksum methods to validate results returned by a numerical subroutine operating subject to unpredictable errors in data. We consider common matrix and Fourier algorithms which return results satisfying a necessary condition having a linear form; the checksum tests compliance with this condition. We discuss the theory and practice of setting numerical tolerances to separate errors caused by a fault from those inherent in finite-precision floating-point calculations. We concentrate on comprehensively defining and evaluating tests having various accuracy/computational burden tradeoffs, and we emphasize average-case algorithm behavior rather than using worst-case upper, bounds on error.
Michael J. Turmon, Robert A. Granat, Daniel S. Katz, John Z. Lou
IEEE Trans. Computers3
2001 Fault-Tolerant High-Performance Matrix Multiplication: Theory and Practice
abstract
We extend the theory and practice regarding algorithmic fault-tolerant matrix-matrix multiplication, C=AB, in a number of ways. First, we propose low-overhead methods for detecting errors introduced not only in C but also in A and/or B. Second, we show that, theoretically, these methods will detect all errors as long as only one entry, is corrupted. Third we propose a low-overhead roll-back approach to correct errors once detected. Finally, we give a high-performance implementation of matrix-matrix multiplication that incorporates these error detection and correction methods. Empirical results demonstrate that these methods work well in practice while imposing an acceptable level of overhead relative to high-performance implementations without fault-tolerance.
John A. Gunnels, Robert A. van de Geijn, Daniel S. Katz, Enrique S. Quintana-Ortí
DSN3
2000 Development of a Spaceborne Embedded Cluster
abstract
Over the last decade and continuing into the foreseeable future, a trend has developed in the spacecraft industry of both number of missions and the amount of data taken by each mission increasing faster than bandwidth capabilities to send these data to Earth. The result of this trend is a bottleneck between data gathering (on-board) and data analysis (on the ground). This bottleneck can be overcome by performing data analysis on-board and only transferring the results of this analysis to the ground, rather than the raw data. One attempt to do this is being made by the NASA HPCC Remote Exploration and Experimentation (REE) Project, which is developing spaceborne embedded clusters. Spaceborne embedded clusters share many characteristics of traditional, ground-based clusters such as POSIX-compliant operating systems and message-passing applications, but also have significant differences, including packaging and the need for fault-tolerance and real-time scheduling in software. This paper discusses these similarities and differences, and how they impact application development.
Daniel S. Katz, Paul L. Springer
CLUSTER1
2000 Demonstration of the Remote Exploration and Experimentation (REE) Fault-Tolerant Parallel-Processing Supercomputer for Spacecraft Onboard Scientific Data Processing
abstract
Concerns a demonstration of the REE Project's work to date. The demonstration is intended to simulate an REE system that might exist on a Mars rover, consisting of multiple COTS processors, a COTS network, a COTS node-level operating system, REE middleware, and an REE application. The specific application performs texture processing of images. It was chosen as a building block of automated geological processing that will eventually be used for both navigation and data processing. Because the COTS hardware is not radiation hardened, single-event-upset-induced soft errors will occur. These errors are simulated in the demonstration by use of a software-implemented fault-injector, and are injected at a rate much higher than is realistic for the sake of viewer interest. Both the application and the middleware contain mechanisms for both detection of and recovery from these faults, and these mechanisms are tested by this very high fault-rate. The consequence of the REE system being able to tolerate this fault rate while continuing to process data is that the system will easily be able to handle the true fault rate.
Fannie Chen, Loring Craymer, Jeff Deifik, Alvin J. Fogel, Daniel S. Katz, Alfred G. Silliman Jr., Raphael R. Some, Sean A. Upchurch, Keith Whisnant
DSN5
2000 Software-Implemented Fault Detection for High-Performance Space Applications
abstract
We describe and test a software approach to overcoming radiation-induced errors in spaceborne applications running on commercial off-the-shelf components. The approach uses checksum methods to validate results returned by a numerical subroutine operating subject to unpredictable errors in data. We can treat subroutines that return results satisfying a necessary condition having a linear form; the checksum tests compliance with this condition. We discuss the theory and practice of setting numerical tolerances to separate errors caused by a fault from those inherent infinite-precision numerical calculations. We test both the general effectiveness of the linear fault tolerant schemes we propose, and the correct behavior of our parallel implementation of them.
Michael J. Turmon, Robert A. Granat, Daniel S. Katz
DSN3
1997 Optimization of a Parallel Ocean General Circulation Model
abstract
Global climate modeling is one of the grand challenges of computational science, and ocean modeling plays an important role in both understanding the current climatic conditions and predicting the future climate change. Three-dimensional time-dependent ocean general circulation models (OGCMs) require a large amount of memory and processing time to run realistic simulations. Recent advances in computing hardware have dramatically affected the prospect of studying the global climate. The significant computational resources of massively parallel supercomputers promise to make such studies feasible. In addition to using advanced hardware, designing and implementing a well-optimized parallel ocean code will significantly improve the computational performance and reduce the total research time to complete these studies.In our present work, we chose the most widely used OGCM code as our base code. This OGCM is based on the Parallel Ocean Program (POP) developed in FORTRAN 90 on the Los Alamos CM-2 Connection Machine by the Los Alamos ocean modeling research group. During the first half of 1994, the code was ported to the Cray T3D by Cray Research using SHMEM-based message passing. Since the code on the T3D was still time-consuming when large problems were encountered, improving the code performance was considered essential.We have developed several general strategies to optimize the ocean general circulation model on the Cray T3D. These strategies include memory optimization, effective use of arithmetic pipelines, and usage of optimized libraries. The optimized code runs 2 to 2.5 times faster than the original code, which gives significant performance improvements for modeling large scaled ocean flows. Many test runs for both of the original and the optimized code have been carried out on the Cray T3D using various numbers of processors (1-256). Comparisons are made for a variety of real-world problems. A nearly linear scaling performance line is obtained for the optimized code, while the speed up data of the optimized code also shows excellent improvement over the original code.In addition to discussing the optimization of the code, we also address the issue of portability. Given the short life cycle of the massively parallel computer, usually on the order of three to five years, we emphasize the portability of the ocean model and the associated optimization routines across several computing platforms. Currently, the ocean modeling code has been ported successfully to the Hewlett Packard (HP)/Convex SPP-2000, and is readily portable to Cray T3E.This paper reports our efforts to optimize the parallel implementations of the oceanic model. So far, the work has focused on improving the load balancing and single node performance of the code on the Cray T3D. As a result, the atmosphere and ocean model components running side-by-side can achieve a performance level of slightly more than 10 GFLOPS on 512 processors of that machine. We have also developed a user-friendly coupling interface with atmospheric and biogeochemical models, in order to make the global climate modeling more complete and more realistic.
Daniel S. Katz, Yi Chao
SC2