Kenneth Chiu

dblp:65/4238 · DBLP profile ↗
← Back
55ranked-venue papers
4as first author
9since 2021 · last 2025
0000-0001-5643-1043ORCID · corroborated

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

Systems, architecture and hardware · 30 · 3 first-author · 3 since 2021Applied, interdisciplinary, general and emerging computing · 14 · 1 first-author · 3 since 2021Software engineering, systems software and programming languages · 12 · 1 first-author · 2 since 2021Artificial intelligence and machine learning · 7 · 3 since 2021Databases, data management, data science and information retrieval · 5 · 2 since 2021Graphics, computer vision, multimedia, augmented reality and games · 1Human-computer interaction and ubiquitous computing · 1
YearPublicationVenuePosition
2025 BOOM: Benchmarking Out-Of-distribution Molecular Property Predictions of Machine Learning Models
abstract
Data-driven molecular discovery leverages artificial intelligence/machine learning (AI/ML) and generative modeling to filter and design novel molecules. Discovering novel molecules requires accurate out-of-distribution (OOD) predictions, but ML models struggle to generalize OOD. Currently, no systematic benchmarks exist for molecular OOD prediction tasks. We present BOOM, $\textbf{b}$enchmarks for $\textbf{o}$ut-$\textbf{o}f$-$\textbf{d}$istribution $\textbf{m}$olecular property predictions: a chemically-informed benchmark for OOD performance on common molecular property prediction tasks. We evaluate over 150 model-task combinations to benchmark deep learning models on OOD performance. Overall, we find that no existing model achieves strong generalization across all tasks: even the top-performing model exhibited an average OOD error 3$\times$ higher than in-distribution. Current chemical foundation models do not show strong OOD extrapolation, while models with high inductive bias can perform well on OOD tasks with simple, specific properties. We perform extensive ablation experiments, highlighting how data generation, pre-training, hyperparameter optimization, model architecture, and molecular representation impact OOD performance. Developing models with strong OOD generalization is a new frontier challenge in chemical ML. This open-source benchmark is available at https://github.com/FLASK-LLNL/BOOM
Evan R. Antoniuk, Shehtab Zaman, Tal Ben-Nun, Peggy Li, James Diffenderfer, Büsra Sahin, Obadiah Smolenski, Everett Grethel, Tim Hsu, Anna M. Hiszpanski, Kenneth Chiu, Bhavya Kailkhura, Brian Van Essen
NeurIPS11
2022 QSketch: GPU-Aware Probabilistic Sketch Data Structures
abstract
A fundamental problem in data analysis is determining the frequency distribution of items within a data set. A count-min sketch is a widely used data structure that can estimate such a distribution. General purpose computation with graphics processing units (GPUs) has become increasingly prevalent over the last decade or so. A GPU provides much higher instruction throughput and memory bandwidth than a typical CPU. However, achieving maximum throughput on GPUs for count-min sketch is challenging, especially for large, uniformly distributed data sets, because each operation needs several random memory accesses that are slow and inefficient on GPUs. In this paper, we propose a suite of novel count-min sketch structures called QSketch. QSketch structures need only one coalesced memory access on average. We also present methods and algorithms to improve the accuracy without decreasing the searching performance. We implemented the QSketch with CUDA 11 and evaluated its performance and accuracy on various GPU platforms.
Huanyi Qin, Zongpai Zhang, Madhusudhan Govindaraju, Kenneth Chiu
CCGRID5
2022 Parallelizing Graph Neural Networks via Matrix Compaction for Edge-Conditioned Networks
abstract
Graph neural networks (GNNs) are a powerful approach for machine learning on graph datasets. Such datasets often consist of millions of modestly-sized graphs, making them well-suited for data-parallel training. However, existing methods show poor scaling due to load imbalances and kernel overheads. We propose an optimized 2D scatter-gather based represen-tation of GNNs that is amenable to distributed, data-parallel training without changing the underlying mathematics of the GNN. By padding graph data to a fixed size on each process, we can simplify data ingestion, make use of efficient compute kernels, equally distribute computation load, and reduce overheads. We benchmark edge-conditioned GNNs with the PCQM4M-LSC and OGB-PPA datasets. Our implementation shows better runtime performance than the state-of-the-art, with a$12\times$strong-scaling speedup on 16 GPUs and an$89.4\times\ \text{weak}$-scaling speedup on 100 GPUs.
Shehtab Zaman, Tim Moon, Tom Benson, Sam Ade Jacobs, Kenneth Chiu, Brian Van Essen
CCGRID5
2022 ParticleGrid: Enabling Deep Learning using 3D Representation of Materials
abstract
From AlexNet to Inception, autoencoders to diffusion models, the development of novel and powerful deep learning models and learning algorithms has proceeded at breakneck speeds. In part, we believe that rapid iteration of model architecture and learning techniques by a large community of researchers over a common representation of the underlying entities has resulted in transferable deep learning knowledge. As a result, model scale, accuracy, fidelity, and compute performance have dramatically increased in computer vision and natural language processing. On the other hand, the lack of a common representation for chemical structure has hampered similar progress. To enable transferable deep learning, we identify the need for a robust 3-dimensional representation of materials such as molecules and crystals. The goal is to enable both materials property prediction and materials generation with 3D structures. While computationally costly, such representations can model a large set of chemical structures. We propose ParticleGrid, a SIMD-optimized library for 3D structures, that is designed for deep learning applications and to seamlessly integrate with deep learning frameworks. Our highly optimized grid generation allows for generating grids on the fly on the CPU, reducing storage and GPU compute and memory requirements. We show the efficacy of 3D grids generated via ParticleGrid and accurately predict molecular energy properties using a 3D convolutional neural network. Our model is able to get 0.006 mean square error and nearly match the values calculated using computationally costly density functional theory at a fraction of the time.
Shehtab Zaman, Ethan Ferguson, Cécile Pereira, Denis Akhiyarov, Mauricio Araya-Polo, Kenneth Chiu
e-Science6
2022 The essence of online data processing
abstract
Data processing systems are a fundamental component of the modern computing stack. These systems are routinely deployed online: they continuously receive the requests of data processing operations, and continuously return the results to end users or client applications. Online data processing systems have unique features beyond conventional data processing, and the optimizations designed for them are complex, especially when data themselves are structured and dynamic. This paper describes DON Calculus, the first rigorous foundation for online data processing. It captures the essential behavior of both the backend data processing engine and the frontend application, with the focus on two design dimensions essential yet unique to online data processing systems: incremental operation processing (IOP) and temporal locality optimization (TLO). A novel design insight is that the operations continuously applied to the data can be defined as an operation stream flowing through the data structure, and this abstraction unifies diverse designs of IOP and TLO in one calculus. DON Calculus is endowed with a mechanized metatheory centering around a key observable equivalence property: despite the significant non-deterministic executions introduced by IOP and TLO, the observable result of DON Calculus data processing is identical to that of conventional data processing without IOP and TLO. Broadly, DON Calculus is a novel instance in the active pursuit of providing rigorous guarantees to the software system stack. The specification and mechanization of DON Calculus provide a sound base for the designers of future data processing systems to build upon, helping them embrace rigorous semantic engineering without the need of developing from scratch.
Philip Dexter, Yu David Liu, Kenneth Chiu
Proc. ACM Program. Lang.3
2021 Deep Learning Techniques for Unmixing of Hyperspectral Stimulated Raman Scattering Images
abstract
Stimulated Raman Scattering (SRS) microscopy is a stain-free, laser-scanning imaging technology that utilizes two coherent laser beams (i.e., the pump and Stokes) to stimulate vibration of chemical bonds in molecules. Different bonds have different resonant vibrational frequencies, and thus SRS microscopy can achieve rapid chemical imaging at high-resolution, enabling live cell imaging and near-instant, stain-free pathological imaging. To increase the ability to resolve different chemical species, multiple Raman wavenumbers can be used with the hyperspectral SRS imaging data. In particular, this approach holds promise for quantifying DNA content, which is important to characterize cancer cell polyploidy. The SRS spectra is the mixture of the spectra of various pure substances present in each pixel so unmixing must be performed to find the relative abundances of these substances. We ran our SRS hyperspectral data of cancer cells through SciPy’s Least Square Error Linear Optimization algorithm (LSQ) [1] but found that it was not able to return the correct DNA content. Our proposed solution to this problem is to use an autoencoder neural network to unmix the spectra. We based the network on the findings in Palsson et al. (2018) [2]. Our initial results show that the network is effective at finding an accurate linear combination, but the noise in the collection of the SRS hyperspectral data significantly increases the number of low error solutions which makes it difficult for the network to find the true linear combination. Future work will be focused on using noise reduction techniques to help the network find the true abundance values.
Nikolas Burzynski, Yuhao Yuan, Adiel Felsen, David Reitano, Zhibo Wang 0006, Khalid A. Sethi, Fake Lu, Kenneth Chiu
IEEE BigData8
2021 Cell Nuclei and Lipid Droplets Quantification in Stimulated Raman Images
abstract
Stimulated Raman scattering (SRS) microscopy is a stain-free, laser-scanning imaging technology that allows for rapid chemical imaging at high-resolution. When developing cancer diagnosis based on cellular and tissue pathology, it is critical to understand the number, size, and density of the cell nuclei, as well as other metabolic features, such as the lipid droplets. In our research, we compare the U-Net and Mask R-CNN convolutional neural network architectures to segment cell nuclei from SRS images of cultured cancer cells. We also use a modified version of U-Net to identify the centroids of nuclei and lipid droplets. Combining these centroids with a segmentation, we can generate a Voronoi diagram to estimate the size of each nucleus and lipid droplet. Future work will focus on applying these methods to identify and segment various cellular structure in both cells and human cancer tissues with SRS imaging.
Adiel Felsen, Yuhao Yuan, Nikolas Burzynski, David Reitano, Zhibo Wang 0006, Khalid A. Sethi, Fake Lu, Kenneth Chiu
IEEE BigData8
2021 GVT-Guided Demand-Driven Scheduling in Parallel Discrete Event Simulation
abstract
The performance and scalability of Parallel Discrete Event Simulation (PDES) can be significantly impacted by temporarily inactive threads that occupy CPU resources but do no useful processing. A recent design called Demand-Driven PDES (DD-PDES) identifies such threads and de-schedules them from CPU cores to eliminate the unnecessary overhead. In this paper, we propose significant further improvements to DD-PDES. First, we introduce a new GVT (Global Virtual Time)-guided algorithm named GG-PDES to perform de-scheduling operations in a lock-free fashion and without relying on a centralized controller thread as was used previously. Second, we introduce the Dynamic CPU Affinity algorithm built on top of GG-PDES that adaptively pins simulation threads to CPU cores to achieve a balanced execution. We demonstrate that these optimizations can yield performance improvements in the range of 13% to 50% over the original DD-PDES system.
Ali Eker, David Timmerman, Barry Williams, Kenneth Chiu, Dmitry V. Ponomarev
ICPP4
2021 High-Performance PDES on Manycore Clusters
abstract
Performance and scalability of Parallel Discrete Event Simulation (PDES) is often limited by fine-grain communication, especially in execution environments with high communication cost. Low latencies of on-chip communication in emerging manycore processors promise to substantially alleviate conventional PDES bottlenecks. However, scaling to manycore clusters requires balancing faster on chip communication with slower traditional network communication between cluster nodes. In this work, we investigate performance of PDES on a cluster of Intel's Knights Landing (KNL) processors, identify performance bottlenecks, and propose techniques to address them. Specifically, we propose three performance optimizations: (1) a new design of the communication buffer centered around the use of atomic compare-and-swap operations to reduce synchronization overhead between a dedicated communication thread and computation threads; (2) careful selection of the number of computation threads per communication thread to limit the pressure on each communication thread; and (3) balancing the timing of communication and computation threads to ensure their synchronized forward progress. Combined, these optimizations result in a 2X - 16X speedup over baseline implementations in ROSS simulator.
Barry Williams, Ali Eker, Kenneth Chiu, Dmitry V. Ponomarev
SIGSIM-PADS3
2020 Detecting and Reacting to Anomalies in Relaxed Uses of Raft
abstract
The Raft consensus algorithm is used in many popular distributed key-value stores to offer strong consistency. Due to the cost of implementing strong consistency, its performance characteristics may not meet the requirements of some users. To satisfy these users, many distributed key-value stores allow users to bypass Raft when serving read requests. Unfortunately, yet predictably, this introduces anomalies. While this is a tradeoff many users may be willing to make, the effects of the tradeoff are not properly accounted for: i.e., it is impossible to know how much consistency is being traded away for the increased speed. We propose the use of reflective consistency-a design space of consistency models used to expose anomaly statistics to the system and its users-to regain transparency in the tradeoff space. This work presents the complete lifecycle of implementing an instance of reflective consistency. We first describe how a popular feature in strongly consistent distributed key-value stores causes anomalies. We then design a reflective consistency implementation which can quantify the existence of anomalies. Finally, we evaluate the implementation, showing that, with nearly zero overhead, users are able to regain control over the anomaly behavior of their distributed storage systems.
Philip Dexter, Bedri Sendir, Kenneth Chiu
CCGRID3
2020 SphericRTC: A System for Content-Adaptive Real-Time 360-Degree Video Communication
abstract
We present the SphericRTC system for real-time 360-degree video communication. 360-degree video allows the viewer to observe the environment in any direction from the camera location. This more-immersive streaming experience allows users to more-efficiently exchange information and can be beneficial in the real-time setting. Our system applies a novel approach to select representations of 360-degree frames to allow efficient, content-adaptive delivery. The system performs joint content and bitrate adaptation in real-time by offloading expensive transformation operations to the GPU via CUDA. The system demonstrates that the multiple sub-components -- viewport feedback, representation selection, and joint content and bitrate adaptation -- can be effectively integrated within a single framework. Compared to a baseline implementation, views in SphericRTC have consistently higher visual quality. The median Viewport-PSNR of such views is 2.25 dB higher than views in the baseline system.
Shuoqian Wang, Mengbai Xiao, Kenneth Chiu, Yao Liu 0001
ACM Multimedia4
2020 Demand-Driven PDES: Exploiting Locality in Simulation Models
abstract
Traditional parallel discrete event simulation (PDES) systems treat each simulation thread in the same manner, regardless of whether a thread has events to process in its input queue or not. At the same time, many real-life simulation models exhibit significant execution locality, where only part of the model (and thus a subset of threads) are actively sending or receiving messages in a given time period. These inactive threads still continuously check their queues and participate in simulation-wide time synchronization mechanisms, such as computing Global Virtual Time (GVT). This wastes resources, ties up CPU cores with threads that offer no contribution to event processing and limits the performance and scalability of the simulation.
Ali Eker, Barry Williams, Kenneth Chiu, Dmitry V. Ponomarev
SIGSIM-PADS3
2019 Controlled Asynchronous GVT: Accelerating Parallel Discrete Event Simulation on Many-Core Clusters
abstract
In this paper, we investigate the performance of Parallel Discrete Event Simulation (PDES) on a cluster of many-core Intel KNL processors. Specifically, we analyze the impact of different Global Virtual Time (GVT) algorithms in this environment and contribute three significant results. First, we show that it is essential to isolate the thread performing MPI communications from the task of processing simulation events, otherwise the simulation is significantly imbalanced and performs poorly. This applies to both synchronous and asynchronous GVT algorithms. Second, we demonstrate that synchronous GVT algorithm based on barrier synchronization is a better choice for communication-dominated models, while asynchronous GVT based on Mattern's algorithm performs better for computation-dominated scenarios. Third, we propose Controlled Asynchronous GVT (CA-GVT) algorithm that selectively adds synchronization to Mattern-style GVT based on simulation conditions. We demonstrate that CA-GVT outperforms both barrier and Mattern's GVT and achieves about 8% performance improvement on mixed computation-communication models. This is a reasonable improvement for a simple modification to a GVT algorithm.
Ali Eker, Barry Williams, Kenneth Chiu, Dmitry V. Ponomarev
ICPP3
2019 An Error-Reflective Consistency Model for Distributed Data Stores
abstract
Consistency models for distributed data stores offer insights and paths to reasoning about what a user of such a system can expect. However, often consistency models are defined or implemented in coarse-grained manners, making it difficult to achieve precisely the consistency required. Further, many domains are already written to handle anomalies in distributed systems, yet they have little opportunity for expressing or taking advantage of their leniency. We propose reflective consistency-an active solution which adapts an underlying data store to changing loads and resource availability to meet a given consistency level. We implement reflective consistency in Cassandra, an existing distributed data store supporting per-read and perwrite consistency. Our implementation allows users to express their anomaly leniency directly and the system will react to the presence of anomalies, changing Cassandra's consistency only when needed. Users of Reflective Cassandra can expect minimal overhead (anywhere from 1% to 14% depending on configuration) and a 50% decrease in the amount of costly strong reads.
Philip Dexter, Kenneth Chiu, Bedri Sendir
IPDPS2
2018 Performance Implications of Global Virtual Time Algorithms on a Knights Landing Processor
abstract
Recent studies investigated the performance of Parallel Discrete Event Simulation (PDES) on Intel Xeon Phi manycore processors, but generally reported underwhelming performance results, especially at high scales when all cores and thread contexts are fully loaded. While the lack of scalability in an earlier study on a Knights Corner (KC) processor is an artifact of physical limitations of the KC system, performance challenges on a Knights Landing (KNL) system partially stem from a slower global virtual time (GVT) computation algorithm used in that study. In this paper, we re-examine PDES performance on KNL under more efficient GVT algorithms to alleviate the GVT bottleneck. Specifically, we compare a synchronous GVT algorithm based on barrier synchronization, and two asynchronous GVT implementations: a modified Mattern's algorithm for shared memory systems and a recently-proposed wait-free algorithm. Using the ROSS simulator, we demonstrate that minimizing the GVT bottleneck results in significant improvement in scalability, allowing the simulation to scale with performance all the way to 250 threads (per chip). Interestingly, we observe that while for the balanced models the wait-free algorithm is a clear winner, barrier-based GVT provides significantly better results for imbalanced models executed at high scale. We also perform detailed simulation profiling to understand the underlying reasons for these performance trends.
Ali Eker, Barry Williams, Nitesh Mishra, Dushyant Thakur, Kenneth Chiu, Dmitry V. Ponomarev, Nael B. Abu-Ghazaleh
DS-RT5
2016 Lazy graph processing in Haskell
abstract
This paper presents a Haskell library for graph processing: DeltaGraph. One unique feature of this system is that intentions to perform graph updates can be memoized in-graph in a decentralized fashion, and the propagation of these intentions within the graph can be decoupled from the realization of the updates. As a result, DeltaGraph can respond to updates in constant time and work elegantly with parallelism support. We build a Twitter-like application on top of DeltaGraph to demonstrate its effectiveness and explore parallelism and opportunistic computing optimizations.
Philip Dexter, Yu David Liu, Kenneth Chiu
Haskell3
2014 Scaling up Prioritized Grammar Enumeration for scientific discovery in the cloud
abstract
Symbolic Regression (SR) is the data driven search for mathematical relations as performed by a computer. In essence, SR is a search over all possible equations to find those which best model the data on hand. Prioritized Grammar Enumeration (PGE) is a recently proposed algorithm which has been shown to have great efficacy and efficiency on the Symbolic Regression problem, using just a single compute core. PGE reformulates the SR problem as a search over a grammar, makes reductions in the magnitude of the search space, and introduces mechanisms for exploring that space efficiently. Notably, PGE provides reliability and reproducibility of results, a key aspect to any system used by scientists at large. In this paper, we enhance the PGE algorithm in several ways. First, we extend PGE to discover differential equations. Second, we incorporate multiple prioritization heaps into PGE, reducing point evaluations while maintaining efficacy. Finally, we decouple the PGE subroutines into a set of services, contain each with Docker, and deploy them onto the cloud. Our algorithm experiments cover a range of dynamical systems from a multitude of domains. and our cloud experiments explore a variety of architectural setups. Our results show PGE to have great promise and efficacy in automating the discovery of equations at the scales needed by tomorrow's scientific data problems.
Tony Worm, Kenneth Chiu
IEEE BigData2
2013 A stream partitioning approach to processing large scale distributed graph datasets
abstract
RDF datasets are an important source of big data. Many of them, however, are too large to fit on a single machine. One approach to address this is to partition the RDF graph across multiple machines, with each component residing on a single machine. A poor partition can incur significant communication costs, however, if as a result many queries involve multiple machines. A number of existing partitioning schemes seek to reduce these costs by finding partitions that avoid cutting edges in the RDF graph. While these can successfully find good partitions the partitioning process itself is often not very scalable, and not capable of handling incrementally-generated RDF data. In this paper, we develop a more scalable, effective and low complexity approach, online graph dataset partitioning, to produce high quality dataset partitions with fewer links between partitions. We show experimentally that it works well in reducing the communication cost of query processing, while at the same time improving scalability of the partitioning itself.
Kenneth Chiu
IEEE BigData2
2013 Automatic Performance Prediction for Load-Balancing Coupled Models
abstract
Computationally-demanding, parallel coupled models are crucial to understanding many important multi-physics/multiscale phenomena. Load-balancing such simulation son large clusters is often done through off-line, static means that often require significant manual input. Dynamic, runtime load-balancing has been shown in our previous work to be effective, but we still used a manually generated performance predictor to guide the load-balancing decisions. In this paper, we show how timing and interaction information obtained by instrumenting the middleware can be used to automatically generate a performance predictor that relates the overall execution time to the execution time of each individual sub model. The performance predictor is evaluated through the new coupled model benchmark employing five constituent sub models that simulates the CCSM coupled climate model.
Daihee Kim, Jay Walter Larson, Kenneth Chiu
CCGRID3
2013 Prioritized grammar enumeration: symbolic regression by dynamic programming
abstract
We introduce Prioritized Grammar Enumeration (PGE), a deterministic Symbolic Regression (SR) algorithm using dynamic programming techniques. PGE maintains the tree-based representation and Pareto non-dominated sorting from Genetic Programming (GP), but replaces genetic operators and random number use with grammar production rules and systematic choices. PGE uses non-linear regression and abstract parameters to fit the coefficients of an equation, effectively separating the exploration for form, from the optimization of a form. Memoization enables PGE to evaluate each point of the search space only once, and a Pareto Priority Queue provides direction to the search. Sorting and simplification algorithms are used to transform candidate expressions into a canonical form, reducing the size of the search space. Our results show that PGE performs well on 22 benchmarks from the SR literature, returning exact formulas in many cases. As a deterministic algorithm, PGE offers reliability and reproducibility of results, a key aspect to any system used by scientists at large. We believe PGE is a capable SR implementation, following an alternative perspective we hope leads the community to new ideas.
Tony Worm, Kenneth Chiu
GECCO2
2012 Malleable Model Coupling with Prediction
abstract
Achieving ultra scalability in coupled multiphysics and multiscale models requires dynamic load balancing both within and between their constituent subsystems. Interconstituent dynamic load balance requires runtime resizing -- or malleability -- of subsystem processing element (PE) cohorts. We enhance the Malleable Model Coupling Toolkit's Load Balance Manager (LBM) to incorporate prediction of a coupled system's constituent computation times and coupled model global iteration time. The prediction system employs piecewise linear and cubic spline interpolation of timing measurements to guide constituent cohort resizing. Performance studies of the new LBM using a simplified coupled model test bed similar to a coupled climate model show dramatic improvement ( 77%) in the LBM's convergence rate.
Daihee Kim, Jay Walter Larson, Kenneth Chiu
CCGRID3
2012 Optimizing Distributed RDF Triplestores via a Locally Indexed Graph Partitioning
abstract
Semantic web techniques based on RDF are a promising approach to help make sense of the growing deluge in scientific data and conceptually integrate separately administered datasets. However, storing large datasets entirely on a single machine is not scalable, which has led to the concept of distributed triple stores. Existing triple stores, however, fail to take advantage of the non-uniform nature of semantic web data, leading to inefficient data allocation. In this paper, we extend our previous work on uniform graph partitioning to include a local index which is used to filter intermediate results and optimize sub-querying, and a new system design which can better exploit graph partitioned datasets. This also allows us to measure the total communication cost in our simulation. As shown in our experiments, our approach can effectively reduce the communication cost of query-processing messages compared with other approaches, balance the size of partitions compared with other approaches, and enhanced parallelism through independent sub-querying.
Kenneth Chiu
ICPP2
2012 Dynamic Load Balancing for Malleable Model Coupling
abstract
Dynamic load balancing both within and between constituent subsystems is required to achieve ultrascalability in coupled multiphysics and multiscale models. Interconstituent dynamic load balancing requires runtime resizing-or malleability-of subsystem processing element (PE) cohorts. In our previous work, we developed and introduced the Malleable Model Coupling Toolkit with a load balance manager implementing and providing a runtime load-balancing algorithm using PE reallocation across subsystems. In this paper, we extend that work by adding the ability to adapt to coupled models that have changing loads during execution. We evaluate the algorithm through a synthetic coupled-model benchmark that uses the LogP performance model as applied to parallel LU decomposition.
Daihee Kim, Jay Walter Larson, Kenneth Chiu
ISPA3
2012 A Graph Partitioning Approach to Distributed RDF Stores
abstract
With growing of Semantic Web data, especially RDF data, managing large RDF dataset on a single machine does not scale well. Previous work has explored how to distribute RDF triples to multiple machines. However due to inefficient dataset partitioning used by these solutions, the performance of distributed store system is significantly affected. In this paper, we proposed a promising approach that utilized the graph nature of RDF datasets to minimize relations between partitions after dataset partitioning, and optimized system design based on it. As shown in our experiments, our approach can effectively reduce communication cost of query-processing messages, balance size of partitions compared with other approaches, and enhance parallelism through independent sub-querying.
Kenneth Chiu
ISPA2
2010 A Parallel XPath Engine Based on Concurrent NFA Execution
abstract
The importance of XPath in XML filtering systems has led to a significant body of research on improving the processing performance of XPath queries. Most of the work, however, has been in the context of a single processing core. Given the prevalence of multicore processors, we believe that a parallel approach can provide significant benefits for a number of application scenarios. In this paper we thus investigate the use of multiple threads to concurrently process XPath queries on a shared incoming XML document. Using an approach that builds on YFilter, we divide the NFA into several smaller ones for concurrent processing. We implement and test two strategies for load balancing: a static approach and a dynamic approach. We test our approach on an eight-core machine, and show that it provides reasonable speedup up to eight cores.
Yinfei Pan, Kenneth Chiu
ICPADS3
2010 Towards Realistic Networks for Simulating Large-Scale Distributed Systems
abstract
Research in large-scale distributed systems, such as P2P systems, often relies critically on simulations to validate research results. Though systems such as PlanetLab can be used to test on real networks in some cases, there are still significant practical challenges to evaluating large-scale distributed systems on actual hardware. Actual measured datasets such as the ones measured with King's method also have an important role, but often do not provide enough scale or are not representative of the network for which the research is intended. Typically, in such cases, tools such as GT-ITM are used to generate network topologies for the evaluation simulation. These tools work reasonably well at generating physical topologies that are representative of real systems. The behavior of distributed systems, however, depends not just on the physical topologies, but also the routing policies and other factors that affect latency and bandwidth. These aspects may have a considerable impact on any evaluation performed on the generated network, and can lead to significant differences between simulated performance and actual performance. In particular, triangle inequality violations and path inflation can adversely impact large-scale distributed systems. In this paper, we present techniques and approaches for adding such real-world effects to generated networks, in a parametrized approach. We show the parameters can be varied to generate a variety of networks with different characteristics, and compare them to measured datasets.
Ketan Bahulkar, Daihee Kim, Dmytro Zhydkov, Kenneth Chiu
ISPA5
2010 A Graph Clustering Approach to Computing Network Coordinates
abstract
In the technique known as network coordinates, the network latency between nodes is modeled as the distance between points in a metric space. Actual network latencies, however, exhibit numerous triangle inequality violations, which result in significant error between the actual latency and the distance as determined by the network coordinates. In this work, we show how graph clustering techniques can be used to find regions of the network space that show low triangle inequality violation within the region. By using techniques to increase the relative edge density in these regions, we improve the accuracy of network coordinates in these regions. We reduce the relative error within a cluster by 15% on average for the Meridian dataset, and by 7% over all; when compared to a single spring relaxation over the whole network.
Beilan Wang, Kenneth Chiu
PDP3
2009 Brain Image Registration Analysis Workflow for fMRI Studies on Global Grids
abstract
Scientific applications like neuroscience data analysis are usually compute and data-intensive. With the use of globally distributed resources and suitable middlewares, we can achieve much shorter execution time, distribute compute and storage load, and add greater flexibility to the execution of these scientific applications than we could ever achieve in a single compute resource.In this paper, we present the processing of Image Registration (IR) for Functional Magnetic Resonance Imaging(fMRI) studies on global Grids. We characterize the application, list its requirements and transform it to a workflow. We then execute the application on Grid’5000 platform and present extensive performance results. We show that the IR application can have 1) significantly improved makespan, 2) distribution of compute and storage load among resources used, and 3) flexibility when executing multiple times on global Grids.
Suraj Pandey, William Voorsluys, Mustafizur Rahman 0003, Rajkumar Buyya, James E. Dobson, Kenneth Chiu
AINA6
2009 Speculative p-DFAs for parallel XML parsing
abstract
XML has seen wide acceptance in a number of application domains, and contributed to the success of wide-scale grid and scientific computing environments. Performance, however, is still an issue, and limits adoption under some situations where it might otherwise be able to provide significant interoperability, flexibility, and extensibility. As CPUs increasingly have multiple cores, parallel XML parsing can help to address this concern. This paper explores the use of speculation to improve the performance of parallel XML parsing. Building on previous work, we use an initial preparsing stage to build a sketch of the document which we called the skeleton. This skeleton contains enough information so that we can then proceed to do the full parse in parallel using unmodified libxml2. The preparsing itself is parallelized using product machines which we call p-DFAs. During execution, unlikely possibilities are discarded in favor of more likely ones. Statistics are gathered to decide which possibilities are not likely. The results show good performance and scalability on both a 30 CPU Sun E6500 machine running Solaris and a Linux machine with two Intel Xeon L5320 CPUs for a total of 8 physical cores.
Yinfei Pan, Kenneth Chiu
HiPC3
2009 A grid workflow environment for brain imaging analysis on distributed systems
abstract
Abstract Scientific applications like neuroscience data analysis are usually compute and data‐intensive. With the use of the additional capacity offered by distributed resources and suitable middlewares, we can achieve much shorter execution time, distribute compute and storage load, and add greater flexibility to the execution of these scientific applications than we could ever achieve in a single compute resource. In this paper, we present the processing of image registration (IR) for functional magnetic resonance imaging studies on Global Grids. We characterize the application, list its requirements and then transform it to a workflow. We use Gridbus Broker and Gridbus Workflow Engine technologies for executing the neuroscience application on the Grid. We developed a complete web‐based portal integrating GUI‐based workflow editor, execution management, monitoring and visualization of tasks and resources. We describe each component of the system in detail. We then execute the application on Grid'5000 platform and present extensive performance results. We show that the IR application can have (1) significantly improved makespan, (2) distribution of compute and storage load among resources used, and (3) flexibility when executing multiple times on Grid resources. Copyright © 2009 John Wiley & Sons, Ltd.
Suraj Pandey, William Voorsluys, Mustafizur Rahman 0003, Rajkumar Buyya, James E. Dobson, Kenneth Chiu
Concurr. Comput. Pract. Exp.6
2009 Special Section: Third IEEE International Conference on e-Science and Grid Computing
Kenneth Chiu, Geoffrey C. Fox
Future Gener. Comput. Syst.1
2008 Parsing XML Using Parallel Traversal of Streaming Trees
Yinfei Pan, Kenneth Chiu
HiPC3
2008 Hybrid Parallelism for XML SAX Parsing
abstract
XML has been widely adopted across a wide spectrum of applications. Its parsing efficiency, however, remains a concern, and can be a bottleneck. At the same time, with the trend towards multicore CPUs, parallelization to improve performance has become increasingly relevant. In previous work, we have investigated parallelizing DOM-style parsing and gained significant speedup. For streaming XML applications, however, SAX-style parsing is often required. In this paper, we present a technique and implementation of a parallel XML SAX parser. To handle inherent data dependencies in XML while still allowing reasonable scalability, we use a 4-stage software pipeline with a combination of strictly sequential stages and stages that can be further data-parallelized within the stage. We thus utilize a hybrid between pipelined parallelism and data parallelism. To demonstrate effectiveness, we test this approach on a Linux machine with two Intel Xeon L5320 CPUs for a total of 8 physical cores, and obtain good speedup up to about 8 CPUs.
Yinfei Pan, Kenneth Chiu
ICWS3
2008 Grid-based research, development, and deployment in New York State
abstract
In this paper, we present cyberinfrastructure and grid computing efforts in New York State. In particular, we focus on fundamental efforts in Binghamton and Buffalo, including the design, development, and deployment of the New York State Grid, as well as a grass-roots New York State Initiative.
Russ Miller, Jonathan J. Bednasz, Kenneth Chiu, Steven M. Gallo, Madhusudhan Govindaraju
IPDPS3
2008 Simultaneous transducers for data-parallel XML parsing
abstract
Though XML has gained significant acceptance in a number of application domains, XML parsing can still be a vexing performance bottleneck. With the growing prevalence of multicore CPUs, parallel XML parsing could be one option for addressing this bottleneck. Achieving data parallelism by dividing the XML document into chunks and then independently processing all chunks in parallel is difficult, however, because the state of an XML parser at the first character of a chunk depends potentially on the characters in all preceding chunks. In previous work, we have used a sequential preparser implementing a preparsing pass to determine the document structure, followed by a parallel full parse. The preparsing is sequential, however, and thus limits speedup. In this work, we parallelize the preparsing pass itself by using a simultaneous finite transducer (SFT), which implicitly maintains multiple preparser results. Each result corresponds to starting the preparser in a different state at the beginning of the chunk. This addresses the challenge of determining the correct initial state at beginning of a chunk by simply considering all possible initial states simultaneously. Since the SFT is finite, the simultaneity can be implemented efficiently simply by enumerating the states, which limits the overhead. To demonstrate effectiveness, we use an SFT to build a parallel XML parsing implementation on an unmodified version of libxml2, and obtained good scalability on both a 30 CPU Sun E6500 machine running Solaris and a Linux machine with two Intel Xeon L5320 CPUs for a total of 8 physical cores.
Yinfei Pan, Kenneth Chiu
IPDPS3
2007 A Static Load-Balancing Scheme for Parallel XML Parsing on Multicore CPUs
abstract
A number of techniques to improve the parsing performance of XML have been developed. Generally, however, these techniques have limited impact on the construction of a DOM tree, which can be a significant bottleneck. Meanwhile, the trend in hardware technology is toward an increasing number of cores per CPU. As we have shown in previous work, these cores can be used to parse XML in parallel, resulting in significant speedups. In this paper, we introduce a new static partitioning and load-balancing mechanism. By using a static, global approach, we reduce synchronization and load-balancing overhead, thus improving performance over dynamic schemes for a large class of XML documents. Our approach leverages libxm12 without modification, which reduces development effort and shows that our approach is applicable to real-world, production parsers. Our scheme works well with Sun's Niagara class of CMT architectures, and shows that multiple hardware threads can be effectively used for XML parsing.
Yinfei Pan, Kenneth Chiu
CCGRID4
2007 Portal Services for Collaborative Remote Instrument Control, Monitoring and Data Access
abstract
A two component portal system is being developed for collaborative remote instrument and data control and monitoring. The system builds on and enhances the common instrument middleware architecture (CIMA) model for Web services based monitoring of remote scientific instruments and sensors. The architecture supports remote access to multiple instruments from a single portal. Plugin modules are used to provide flexibility and re-use, and the notion of plugin control is being developed. The use of Web 2.0 Pushlet and AJAX technologies has been introduced for push based portlet refresh and updating. An X3D based 3D virtual representation of the instrument provides data collection simulation and (pseudo) real time instrument representation. An important component of the system is a Webs services driven portlet for collaborative image viewing.
Douglas du Boulay, Clinton Chee, Kenneth Chiu, Richard Leow, Donald F. McMullen, Romain Quilici, Peter Turner
eScience3
2007 Index Structures for Efficient Querying of Distributed Triplestores
abstract
Data is dynamically structured by nature and can be highly diverse and multifaceted. Often, such diverse and complex information needs to be linked. Conventional data-stores, such as relational databases, do not conveniently accommodate dynamically varying structures, as frequently modifying database schemas is not feasible. RDF triplestores offer a flexible solution for handling such data, where any property about an entity can be described by a triple having a subject, a predicate, and an object. Also, data is inherently distributed due to origination points, ownership and many other reasons. Furthermore, storing data in triplestores gives rise to the need to distribute data due to the large number of triples that would result by migrating existing data from a database, for example. In this paper, we present our work on designing index structures in order to facilitate efficient querying of a distributed triplestore (DTS). The distributed querying algorithm in DTS makes use of a sub-graph isomorphism approach to eliminate traversing edges between triplestores that does not have the potential to produce any results. We show that our triplestore has equivalent performance as 3Store when used in a non-distributed mode. Our performance tests in the distributed mode show that the indexes improve efficiency of querying.
Tharaka Devadithya, Kenneth Chiu
eScience2
2007 A High Performance Schema-Specific XML Parser
abstract
Performance of XML parsers with validation are usually suffer. This is because such parsers need first parsing and undertanding XML schemas, and thus are limited by the very complexity of XML schemas. Schema-specific approach, however, may adjust such problem. In this paper, we introduce a high performance SAX like validating XML parser using a schema-specific approach. In this approach, a schema compiler first transforms the schema into an intermediate representation, called generalized automata, which abstracts the computations required to parse XML documents as well as validate them against a schema. The generalized automaton is then translated to a schema specific parser, which is capable of parsing and validating XML documents with namespaces through a schema specific modified SAX API. Our performance evaluation shows good results when compared with other validating parsers.
Zhenghong Gao, Yinfei Pan, Kenneth Chiu
eScience4
2007 Parallel XML Parsing Using Meta-DFAs
abstract
By leveraging the growing prevalence of multicore CPUs, parallel XML parsing(PXP) can significantly improve the performance of XML, enhancing its suitability for scientific data which is often dominated by floating-point numbers. One approach is to divide the XML document into equal-sized chunks, and parse each chunk in parallel. XML parsing is inherently sequential, however, because the state of an XML parser when reading a given character depends potentially on all preceding characters. In previous work, we addressed this by using a fast preparsing scan to build an outline of the document which we called the skeleton. The skeleton is then used to guide the parallel full parse. The preparse is a sequential phase that limits scalability, however, and so in this paper, we show how the preparse itself can be parallelized using a mechanism we call a meta-DFA. For each state q of the original preparser the meta-DFA incorporates a complete copy of the preparser state machine as a sub-DFA which starts in state q. The meta-DFA thus runs multiple instances of the preparser simultaneously when parsing a chunk, with each possible preparser state at the beginning of a chunk represented by an instance. By pursuing all possibilities simultaneously, the meta-DFA allows each chunk to be preparsed independently in parallel. The parallel full parse following the preparse is performed using libxml2, and outputs DOM trees that are fully compatible with existing applications that use libxml2. Our implementation scales well on a 30 CPU Sun E6500 machine.
Yinfei Pan, Kenneth Chiu
eScience3
2007 Design and implementation issues for distributed CCA framework interoperability
abstract
Abstract Component frameworks, including those that support the Common Component Architecture (CCA), represent a promising approach to addressing the challenge of building and deploying high‐performance scientific applications in Grid environments, one that is being realized, for example, in our LegionCCA and XCAT‐C++ frameworks. The next step beyond building independent individual frameworks is making them interoperate. Component‐based applications should be able to transparently span multiple disjoint component frameworks with low overhead as compared with the same applications running within a single framework. Interoperable frameworks enable applications to take advantage of more resources, and to better match constituent parts to the underlying resources that best support them. The CCA specification does not prescribe a wire format for inter‐component calls in distributed frameworks, thereby promoting considerable flexibility and customization for the framework developer. This approach thus requires an additional specific strategy outside of the CCA to support interoperability between distributed frameworks. Mandating one common wire format, however, risks choosing the wrong format. We discuss in detail five underlying component framework interoperability requirements, and three general approaches to addressing them. We then discuss how the approaches can be applied to meet the requirements, and address the advantages, issues, and implications of doing so. This effectively defines a design space for framework interoperability approaches. We then address the communication interoperability in detail via a single multi‐protocol communication library called Proteus, and discuss how we have incorporated it into two distinct distributed framework implementations of the CCA specification: LegionCCA and XCAT‐C++. Copyright © 2006 John Wiley & Sons, Ltd.
Madhusudhan Govindaraju, Michael J. Lewis, Kenneth Chiu
Concurr. Comput. Pract. Exp.3
2006 CIMA Based Remote Instrument and Data Access: An Extension into the Australian e-Science Environment
abstract
The Common Instrument Middleware Architecture (CIMA) is being used as a core component of a portal based remote instrument access system being developed as an Australian e-Science project. The CIMA model is being enhanced to use federated Grid storage infrastructure (SRB), and the Kepler workflow system to, as much as possible, automate data management, and the facile extraction and generation of instrument and experimental metadata. The Personal Grid Library is introduced as a user friendly portlet interface to SRB data and metadata, and which supports customisable metadata schemas. An Instrument Instruction Module has been introduced as a CIMA plug-in for instrument control. A virtual instrument portlet provides a simulation of the instrument during a data collection. The system is being further augmented with a tool for collaborative data visualisation and evaluation.
Ian Atkinson, Douglas du Boulay, Clinton Chee, Kenneth Chiu, Tristan King, Donald F. McMullen, Romain Quilici, Nigel G. D. Sim, Peter Turner, Matthew Wyatt
e-Science4
2006 Toward Standards for Integration of Instruments into Grid Computing Environments
abstract
Instruments and sensors are the primary sources of data driving science and the development and refinement of theory. A critical component of eresearch yet to be clarified is the role of instruments in cyberinfrastructure. Are instruments, sensors and other real-time data sources to be mediated by file systems, or can they be fully integrated into computing and storage grids with appropriate protocol standards? Our position is that they can and must become regular grid resources and that standards for doing so represent an important research topic. Our approach, the Common Instrument Middleware Architecture, and a model application, X-ray crystallography, are described here in order to stimulate a broader discussion of requirements for grid-enabling instruments and sensors.
Donald F. McMullen, Ian Atkinson, Kenneth Chiu, Peter Turner, Kianosh Huffman, Romain Quilici, Matthew Wyatt
e-Science3
2006 Building a Generic SOAP Framework over Binary XML
abstract
The prevailing binding of SOAP to HTTP specifies that SOAP messages be encoded as an XML 1.0 document which is then sent between client and server. XML processing however can be slow and memory intensive, especially for scientific data, and consequently SOAP has been regarded as an inappropriate protocol for scientific data. Efficiency considerations thus lead to the prevailing practice of separating data from the SOAP control channel. Instead, it is stored in specialized binary formats and transmitted either via attachments or indirectly via a file sharing mechanism, such as GridFTP or HTTP. This separation invariably complicates development due to the multiple libraries and type systems to be handled; furthermore it suffers from performance issues, especially when handling small binary data. As an alternative solution, binary XML provides a highly efficient encoding scheme for binary data in the XML and SOAP messages, and with it we can gain high performance as well as unifying the development environment without unduly impacting the Web service protocol stack. In this paper we present our implementation of a generic SOAP engine that supports both textual XML and binary XML as the encoding scheme of the message. We also present our binary XML data model and encoding scheme. Our experiments show that for scientific applications binary XML together with the generic SOAP implementation not only ease development, but also provide better performance and are more widely applicable than the commonly used separated schemes
Kenneth Chiu, Dennis Gannon
HPDC2
2006 Ontology Based Publish Subscribe Framework
John Skovronski, Kenneth Chiu
iiWAS2
2006 Poster reception - Fast binary serialization for grid systems with XBS
abstract
Efficient serialization and deserialization of data is a fundamental operation of many grid systems. Some serializers are message-based, while others are stream-based. Streaming serializers can be more scalable and flexible, by promoting free form conversations not fixed to any particular static structure. We present the design and implementation of the XBS binary serializer, focusing on three important features: efficient serialization of large and small arrays, efficient pass-through of opaque data by grid intermediaries, and efficient representation of numbers with a large dynamic range. The first feature is important because such arrays dominate scientific computing. The second feature is useful for grid intermediaries such as gateways and proxies, which are becoming ever more important as grid systems become more complex. The third feature is important for efficiently supporting systems without arbitrary size limits. XBS is a freely-available C++ library and provides an object-oriented API, based on generic programming techniques.
Tharaka Devadithya, Kenneth Chiu
SC2
2005 A Binary XML for Scientific Applications
abstract
XML provides flexible, extensible data models and type systems for structured data, and has found wide-acceptance in many domains. XML processing can be slow, however, especially for scientific data, thus leading to the conventional wisdom that XML is not appropriate for such data. Instead, data is stored in specialized binary formats, and is transmitted via work-arounds such as attachments and base64 encoding. Though these work-arounds can be useful, they nonetheless relegate scientific data to second-class status within the Web services framework; and they generally require yet another API, data model, and type system. An alternative solution is to use more efficient encodings of XML, often known as "binary XML". Using XML uniformly throughout an application simplifies and unifies design and development. In this paper we present a binary XML format and implementation for scientific data called Binary XML for Scientific Applications (BXSA). We show that performance is comparable to that of commonly used scientific data formats such as netCDF. These results challenge the prevailing practice of handling control and data separately in scientific applications, with Web services for control and specialized binary formats for data
Kenneth Chiu, Tharaka Devadithya, Aleksander Slominski
e-Science1
2005 The Common Instrument Middleware Architecture: Overview of Goals and Implementation
abstract
Instruments and sensors and their accompanying actuators are essential to the conduct of scientific research. In many cases they provide observations in electronic format and can be connected to computer networks with varying degrees of remote interactivity. These devices vary in their architectures and type of data they capture and may generate data at various rates. In this paper we present an overview of the design goals and initial implementation of the common instrument middleware architecture (CIMA), a framework for making instruments and sensors network accessible in a standards-based, uniform way, and for interacting remotely with instruments and the data they produce. Some of the issues CIMA addresses include: flexibility in network transport, efficient and high throughput data transport, the availability (or lack of) computational, storage and networking resources at the instrument or sensor platform, evolution of instrument design, and reuse of data acquisition and processing codes
Tharaka Devadithya, Kenneth Chiu, Kianosh Huffman, Donald F. McMullen
e-Science2
2005 XCAT-C++: Design and Performance of a Distributed CCA Framework
Madhusudhan Govindaraju, Michael R. Head, Kenneth Chiu
HiPC3
2005 A streaming validation model for SOAP digital signature
abstract
The XML signature specification provides a rich and flexible message signature model for XML documents, and it has been adopted by SOAP applications to provide message-level security. However, the XML signature design introduces a number of complex processing steps, such as canonicalization and XPath filtering, that often lead to performance and scalability problems when encountering extremes of size and rate in the processing of XML. In this paper, we focus on the performance of validating large signed XML messages, as might be sent by a scientific application using grid Web services. We present the design and implementation of the GHPX/SSSV system for the streaming validation of SOAP digital signature. Our model consists of a streaming canonicalization and optimized SOAP signature validation. We present an empirical study of the performance characteristics of these streaming validation features. Based on our evaluations we conclude that the streaming validation model can not only provide high performance, but is also memory efficient.
Kenneth Chiu, Aleksander Slominski, Dennis Gannon
HPDC2
2005 A Benchmark Suite for SOAP-based Communication in Grid Web Services
abstract
The convergence of Web services and grid computing has promoted SOAP, a widely used Web services protocol, into a prominent protocol for a wide variety of grid applications. These applications differ widely in the characteristics of their respective SOAP messages, and also in their performance requirements. To make the right decisions, an application developer must thus understand the complex dependencies between the SOAP implementation and the application. We propose a standard benchmark suite for quantifying, comparing, and contrasting the performance of SOAP implementations under a wide range of representative use cases. The benchmarks are defined by a set of WSDL documents. To demonstrate the utility of the benchmarks and to provide a snapshot of the current SOAP implementation landscape, we report the performance of many different SOAP implementations (gSOAP, AxisJava, XSUL and bSOAP) on the benchmarks, and draw conclusions about their current performance characteristics.
Michael R. Head, Madhusudhan Govindaraju, Aleksander Slominski, Pu Liu, Nayef Abu-Ghazaleh, Robert A. van Engelen, Kenneth Chiu, Michael J. Lewis
SC7
2003 Merging the CCA Component Model with the OGSI Framework
abstract
The most important recent development in Grid systems is the adoption of the Web Services model as its basic architecture. The result is called the Open Grid Services Architecture (OGSA). This paper describes a component framework for distributed Grid applications that is consistent with that model. The framework, called XCAT, is based on the U.S. Department of Energy Common Component Architecture (CCA) but with an implementation based on the standard Web Services stack. Using this framework, an application programmer can compose an application from a set of distributed components. The result is a set of Web Services that collectively represent the executing application instance. This paper describes the basic architecture of XCAT and the design issues to be considered for a component to serve as both a CCA and Open Grid Service Infrastructure (OGSI) service.
Madhusudhan Govindaraju, Sriram Krishnan, Kenneth Chiu, Aleksander Slominski, Dennis Gannon, Randall Bramley
CCGRID3
2002 Investigating the Limits of SOAP Performance for Scientific Computing
abstract
The growing synergy between Web Services and Grid-based technologies will potentially enable profound, dynamic interactions between scientific applications dispersed in geographic, institutional, and conceptual space. Such deep interoperability requires the simplicity, robustness, and extensibility for which SOAP was conceived, thus making it a natural lingua franca. Concomitant with these advantages, however is a degree of inefficiency that may limit the applicability of SOAP to some situations. We investigate the limitations of SOAP for high-performance scientific computing. We analyze the processing of SOAP messages, and identify the issues of each stage. We present a high-performance SOAP implementation and a schema-specific parser based on the results of our investigation. After our SOAP optimizations are implemented, the most significant bottleneck is ASCII/double conversion. Instead of handling this using extensions to SOAP we recommend a multiprotocol approach that uses SOAP to negotiate faster binary protocols between messaging participants.
Kenneth Chiu, Madhusudhan Govindaraju, Randall Bramley
HPDC1
2002 The Proteus multiprotocol message library
abstract
Grid systems span manifold organizations and application domains. Because this diverse environment inevitably engenders multiple protocols, interoperability mechanisms are crucial to seamless, pervasive access. This paper presents the design, rationale, and implementation of the Proteus multiprotocol library for integrating multiple message protocols, such as SOAP and JMS, within one system. Proteus decouples application code from protocol code at run-time, allowing clients to incorporate separately developed protocols without recompiling or halting. Through generic serialization, which separates the transfer syntax from the message type, protocols can also be added independently of serialization routines. We also show performance-enhancing mechanisms for Grid services that examine metadata, but pass actual data through opaquely (such as adapters). The interface provided to protocol implementors is general enough to support protocols as disparate as our current implementations: SOAP, JMS, and binary. Proteus is written in C++; a Java port is planned.
Kenneth Chiu, Madhusudhan Govindaraju, Dennis Gannon
SC1
2000 A Component based Services Architecture for Building Distributed Applications
abstract
Describes an approach to building a distributed software component system for scientific and engineering applications that is based on representing Computational Grid services as application-level software components. These Grid services provide tools such as registry and directory services, event services and remote component creation. While a service-based architecture for grids and other distributed systems is not new, this framework provides several unique features. First, the public interfaces to each software component are described as XML documents. This allows many adaptors and user interfaces to be generated from the specification dynamically. Second, this system is designed to exploit the resources of existing Grid infrastructures like Globus and Legion, and commercial Internet frameworks like e-speak. Third, and most important, the component-based design extends throughout the system. Hence, tools such as application builders, which allow users to select components, start them on remote resources, and connect and execute them, are also interchangeable software components. Consequently, it is possible to build distributed applications using a graphical "drag-and-drop" interface, a Web-based interface, a scripting language like Python, or an existing tool such as Matlab.
Randall Bramley, Kenneth Chiu, Shridhar Diwan, Dennis Gannon, Madhusudhan Govindaraju, Nirmal Mukhi, Benjamin Temko, Madhuri Yechuri
HPDC2