Tonglin Li

dblp:58/11205 · DBLP profile ↗
← Back
31ranked-venue papers
5as first author
12since 2021 · last 2025
0000-0002-2874-8112ORCID · corroborated

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

Applied, interdisciplinary, general and emerging computing · 19 · 1 first-author · 11 since 2021Systems, architecture and hardware · 12 · 4 first-author · 1 since 2021Artificial intelligence and machine learning · 6 · 1 first-authorDatabases, data management, data science and information retrieval · 6 · 1 first-authorSoftware engineering, systems software and programming languages · 1
YearPublicationVenuePosition
2025 Fast Forward Modeling of 3-D Gravity Data for Curved Hexahedral Grid Based on Neural Network
abstract
The unstructured grids are widely used in the processing and interpretation of geophysical data with terrain due to their excellent ability to simulate shape. Among them, the efficiency of the curved hexahedral grid in its gravity forward modeling based on the isoparametric finite-element method is poor due to the complex transformations involving numerous morphological nodes, which limits its application to large-scale data. For this reason, combining with the deep learning technology, the letter proposes a fast forward method of 3D gravity data for the curved hexahedral grid based on the back-propagation (BP) neural network. In the training phase, the method learns the complex mapping of curved hexahedral elements to their gravity sensitivities through the neural network, thereby achieving fast forward modeling during the prediction phase. Numerical examples show that the new method has good simulation accuracy and generalization ability. Under the premise that the training phase can be completed upfront with its cost excluded, its forward efficiency is tens of times higher than that of the isoparametric finite-element method. The successful application of the new method in the actual terrain model of Mount Taishan area in China further proves its practicality.
Tonglin Li, Rongzhe Zhang, Guan-Wen Gu, Zhihe Xu, Teng Luo
IEEE Geosci. Remote. Sens. Lett.2
2025 An Efficient Inversion Method for 3-D Magnetic Surveys Based on Intelligent Anomaly Identification and Adaptive Octree Mesh
abstract
Traditional 3D magnetic inversion methods often rely on global fine meshes when pursuing high precision imaging, leading to inefficient allocation of computational resources, with particularly significant waste in non target areas. To overcome this bottleneck, this paper proposes a novel adaptive octree-based 3D magnetic inversion method with an intelligent anomaly focusing mechanism, achieving an organic unity of high precision and efficiency through multi stage intelligent control. The proposed method first rapidly completes the initial inversion on a coarse mesh to obtain the overall outline of the subsurface structure. Subsequently, the fuzzy c-means (FCM) clustering algorithm is introduced to automatically identify target regions with significant magnetic anomalies from the coarse solution. Based on this, a multi stage octree mesh is constructed within the identified regions to achieve refined discretization of complex geological boundaries. The refined model then serves as a new starting point for high precision inversion. This "coarse inversion, intelligent identification, local refinement, precision inversion" workflow can be iteratively executed on demand to dynamically allocate computational resources. To address the influence of varying mesh scales, a volume-related depth weighting function is adopted, and smooth-focusing regularization is introduced to ensure inversion stability while enhancing the model’s resolution capability. Modeling experiments validate the advantages of our method in terms of both inversion quality and accuracy of the extracted anomaly. We applied this method successfully to magnetic data from a mining area in the Huzhong district of the Heilongjiang province, validating its efficiency and practicality in complex geological settings and demonstrating its broad application prospects in high precision inversion of large scale magnetic data.
Tonglin Li, Rongzhe Zhang, Hua Guo 0005, Xiaoming Pan
IEEE Trans. Geosci. Remote. Sens.2
2024 Magnetotelluric Inversion Constrained by Guided Fuzzy c-Means Clustering Using Adaptive Virtual Rock Physics Information
abstract
The magnetotelluric (MT) inversion technology is crucial for quantitatively interpreting deep mineral resources, especially when combined with rock physics information, enhancing accuracy in assessing underground structural parameters and spatial distribution. However, the traditional fuzzy c-means (FCMs) clustering-constrained inversion method requires prior rock physics information for each geological unit, limiting their application scope. We propose a guided FCMs (GFCMs) clustering-constrained inversion method based on adaptive virtual rock physics information, referred to as XG-FCM-constrained MT inversion. This method breaks free from the constraints of traditional methods by not relying on prior rock physics information. In terms of extracting virtual rock physics information, we employ a local density clustering algorithm to dynamically extract resistivity model information from MT inversion iterations, automatically determining the number of clusters and cluster centers. Regarding the inversion strategy, we construct an integrated objective function that combines data fitting, smoothing constraint, and GFCM constraint, implementing a two-stage iterative solution strategy of “smoothing first, clustering second.” Model testing demonstrates that compared to traditional smoothing-constrained MT inversion, the XG-FCM-constrained method achieves a significant improvement in the resolution of resistivity model reconstruction, clearly delineating the boundaries of underground anomalies. Even in situations where rock physics information is insufficient or absent, this method can effectively reconstruct high-quality underground resistivity models, reducing the dependence on complete prior information. The application of actual field data further highlights the advantages of the XG-FCM-constrained MT inversion method, providing robust support for accurately delineating geological unit boundaries and precisely identifying potential ore deposit target areas.
Rongzhe Zhang, Jiarong Zhang, Tonglin Li, Yang Zhang 0084, Kaixin Du, Xiaoming Pan
IEEE Trans. Geosci. Remote. Sens.3
2023 3-D Joint Inversion of DC Resistivity and Time-Domain Induced Polarization With Structural Constraints in Undulating Topography
abstract
Addressing the significant impact of undulating terrain on 3D Direct Current inversion, the unreliability of traditional linearized polarizability inversion for highly polarizable anomalies, and the non-uniqueness of separate resistivity and Time-Domain Induced Polarization inversion, this paper conducts a joint inversion study of 3D Direct Current resistivity and Time-Domain Induced Polarization with structural constraints in undulating topography. A regularly arranged deformed hexahedron mesh simulates undulating surface terrain, transformed into regular hexahedron elements for 3D undulating terrain DC resistivity modeling. Based on the exact inversion of polarization calculated from the inversion results of apparent resistivity and equivalent apparent resistivity data, a joint inversion of resistivity and polarization constrained by cross-gradient is implemented. Synthetic data examples show that the application of arbitrary hexahedron elements significantly reduces the influence of terrain on inversion, and the implementation of joint inversion markedly improves the recovery of high-polarization anomalies while enhancing both the model resolution and the inversion accuracy. The proposed algorithm is applied to the joint inversion of resistivity and polarizability in the lead-zinc mining area of Xiagalaiaoyi River in Huzhong area, the Great Khingan Mountains, northwestern Heilongjiang Province, achieving good results.
Hetian Yang, Tonglin Li, Rongzhe Zhang, Xintong Dong
IEEE Trans. Geosci. Remote. Sens.2
2023 Application of Supervised Descent Method for 3-D Gravity Data Focusing Inversion
abstract
Three-dimensional gravity inversion is an effective method for extracting underground density distribution from gravity data. However, traditional deterministic gravity inversion methods suffer from problems such as skin effect, low computational accuracy, and poor efficiency. Therefore, we propose a three-dimensional gravity data focusing inversion algorithm based on the supervised descent method. Supervised descent method (SDM) is a non-linear optimization method based on the combination of machine learning and gradient descent method. In the offline phase, we construct a training set based on a priori information and iteratively learn a set of average descent directions between the initial model and the training model. In the online phase, we introduce a focused regularization into the prediction objective function. This addition aims to obtain a sharp boundary density model that conforms to the physical distribution. Additionally, we incorporate property boundary constraints in both the offline and online phases to control the upper and lower bounds of the density values to ensure consistency with reality. Model tests show that the proposed method can effectively overcome skin effect, improve the resolution of gravity inversion. Moreover, the construction of the training set of the proposed method is less affected by prior information, and it has strong generalization ability. Furthermore, the method does not require solving large-scale linear equations, accelerating the inversion computation speed and having strong noise resistance. Field examples demonstrate that this method has good potential for improving the accuracy and efficiency of actual gravity data inversion.
Rongzhe Zhang, Xintong Dong, Tonglin Li, Cai Liu, Xinze Kang
IEEE Trans. Geosci. Remote. Sens.4
2022 Noisy2Noisy: Denoise Pre-Stack Seismic Data Without Paired Training Data With Labels
abstract
In recent years, supervised deep learning-based denoising methods have been popularized and developed rapidly in the field of seismic data processing. Supervised training, however, is limited by the quality and quantity of the paired training data (noisy clean or noisy noise). Data labeling is a time-consuming and expensive work. Compared with raw seismic data, only a small amount of seismic data have been correctly labeled, which, to some extent, limits the long-term development of supervised deep learning-based methods in the field of seismic data denoising. In this letter, we propose an improved denoising framework based on an unsupervised deep learning-based denoising method Noise2Noise, which only needs unprocessed raw seismic data to train the denoising model. Moreover, unlike Noise2Noise, the proposed method does not need to repeatedly collect seismic data to obtain a training pair with similar signal, which is more convenient and effective. Specifically, we propose a block random sampler that can generate training pairs using raw seismic data, which satisfies the training assumption of Noise2Noise that the training pair has a similar signal. In addition, our method has no requirements for the network structure and noise distribution prior and is flexible. Both synthetic seismic data and field seismic data denoising results show that our method can effectively suppress the random noise, and the denoising performance is equivalent to that of supervised deep learning-based denoising methods. In addition, our method may provide a certain reference for geophysical-related research.
Dan Shao, Yuxing Zhao, Yue Li 0003, Tonglin Li
IEEE Geosci. Remote. Sens. Lett.4
2022 3-D Joint Inversion of Gravity and Magnetic Data Using Data-Space and Truncated Gauss-Newton Methods
abstract
Gravity and magnetic inversion are important methods for comprehensive quantitative interpretation of data obtained in, e.g., mineral, oil and gas, and geothermal exploration. At present, the 3-D joint inversion technology of gravity and magnetic data is facing challenges from large-scale data exploration applications. In this letter, a new algorithm for 3-D joint inversion of gravity and magnetic data with high accuracy and low computational cost is presented. We use the geometric trellis method to perform fast forward calculations and then introduce the sparse constraint and adaptive sensitivity matrix into the model constraint terms. The inexact structural resemblance method is then used to add the cross-gradient constraint penalty term to the objective function. Finally, an algorithm (DS-TGN) combining data-space (DS) and truncated Gauss–Newton (TGN) methods is used to solve the joint inversion objective function. Numerical experiments with synthetic data show that the proposed algorithm can significantly reduce the computational cost and obtain high accuracy density and magnetization models with structural resemblance and sharp boundaries. We also apply the DS-TGN algorithm to data obtained in the area of Greater Khingan in northwestern Heilongjiang, China. The underground density and magnetization distribution results provide a high-resolution geological model for the detection of skarn-type deposits.
Rongzhe Zhang, Tonglin Li, Cai Liu, Xingguo Huang, Malte Sommer
IEEE Geosci. Remote. Sens. Lett.2
2022 2-D Magnetotelluric Multiparameter Joint Inversion Considering the Induced Polarization Effect
abstract
Magnetotelluric (MT) is an important geophysical exploration method that uses natural sources to study the electrical structure of the earth. This method is advantageous owing to its low cost, large exploration depth range, and high resolution for low-resistivity bodies. However, traditional MT modelling approaches can only invert the resistivity parameters of geological bodies, while natural geological bodies also exhibit the induced polarization (IP) effect. Applying the IP effect of geological bodies could be effective for exploring mineral resources, such as polymetallic ores, oil (gas) fields, and coal fields. To incorporate the IP information into MT inversion, we proposed a multi-parameter MT joint inversion algorithm that considers the IP effect. We used the Cole-Cole model to integrate four IP parameters into the forward algorithm: the zero-frequency resistivity ρ0, chargeability η, frequency exponentc, and time constantt. The influence of each IP parameter on the forward response was analyzed by forward simulations, and we concluded that the parameters ρ0andη should be considered in the inversion. A cross-gradient function was introduced into the objective function of the Occam inversion method to constrain the structural consistency of ρ0and η. Model testing and practical application results illustrated that the algorithm can not only obtain the subsurface resistivity structure that can be achieved by the traditional MT method but also reveal the distribution of the chargeability parameter. The additional information obtained using this algorithm is conducive to interpreting specific geological structures that cannot be distinguished by the traditional MT method.
Tonglin Li, Rongzhe Zhang
IEEE Trans. Geosci. Remote. Sens.2
2022 Joint Inversion of Multiphysical Parameters Based on a Combination of Cosine Dot-Gradient and Joint Total Variation Constraints
abstract
The joint inversion of structural constraints is a new and rapidly developing detection technology in comprehensive geophysical interpretation. In this article, a new structural constraint 2-D multiphysical parameter joint inversion algorithm for magnetotelluric (MT), gravity, and magnetic data is developed. The structural constraint term is a combination of cosine dot-gradient (CDG) and joint total variation (JTV) constraints, which not only has characteristics of traditional dot product and cross-gradient structure constraints but also avoids the uncertainty of dot product constraints predicting the gradient direction of the model parameters, overcomes the need for high-order differential approximation of the cross-gradient constraint, ignores the influence of the gradient amplitude of different model parameters on the weight of the structural constraint of different regions, and enhances the reconstruction accuracy of the underground discontinuous interface. To more easily combine multiple optimization algorithms to improve the resolution and computational efficiency of joint inversion, an adaptive inexact structural resemblance (IESR) algorithm is developed to minimize numerical solutions to the objective function. Experimental results have demonstrated that the CDG constraint has a wider use range than the traditional structural constraint, the addition of the JTV constraint can recover the underground discontinuous interface, and an inversion result of higher resolution can be obtained using the adaptive IESR algorithm.
Rongzhe Zhang, Tonglin Li, Cai Liu
IEEE Trans. Geosci. Remote. Sens.2
2022 Fast Independent Component Analysis Denoising for Magnetotelluric Data Based on a Correlation Coefficient and Fast Iterative Shrinkage Threshold Algorithm
abstract
Magnetotelluric (MT) sounding data are easily contaminated by various noise sources, especially noise with a long duration (even full-time noise), which makes it difficult to obtain accurate values when calculating the weight factors of the response function, resulting in distortion of the response results. Based on blind source separation theory, fast independent component analysis (FastICA) can separate this kind of noise. However, this method is challenged by the number of separated field sources and the unequal signal amplitude before and after processing. We have developed a novel signal noise separation method, improving the traditional FastICA, where a correlation coefficient is used to compute the number of field sources, and the fast iterative shrinkage threshold algorithm (FISTA) is used to adjust the signal amplitude problem before and after FastICA decomposition. Compared with common field source division methods and the traditional FastICA, the experimental results indicate that our method can separate and remove noise with a long duration, increase the signal-to-noise ratio of the data, and improve the MT response curves. Meanwhile, case studies of measured data illustrate that our method obtains a more robust magnetotelluric response than the conventional robust method and FastICA.
Jiangtao Han, Tonglin Li
IEEE Trans. Geosci. Remote. Sens.3
2022 Research on Magnetotelluric Long-Duration Noise Reduction Based on Adaptive Sparse Representation
abstract
When magnetotelluric (MT) sounding data are measured in mining areas and urban areas, the useful signals are buried under the surrounding interference sources in the whole period, which completely covers up the useful signals, resulting in jump points and distortion of the response curves. Sparse representation uses atoms in a dictionary to process noisy signals. Based on the arbitrariness of the length of these dictionary atoms, they can effectively suppress noise even if the noise fills the entire observation period. We propose an improved sparse representation based on an adaptive dictionary, which can construct a dictionary according to the characteristics of the data itself (extracting the noise and the useful signals separately) and then automatically filter the noise atoms (which are regarded as noise) in the dictionary to suppress long-duration noise (noise lasts for a long time). For the synthetic data and the measured data, in a comparison with the common signal noise separation methods, the results indicate that the proposed method can more completely extract the profiles of the long-duration noise and greatly increase the signal-to-noise ratio. Moreover, for various noise sources, the proposed method indicates better improved performance, obtaining smoother and more reliable response results.
Tonglin Li, Jiangtao Han, Lijia Liu
IEEE Trans. Geosci. Remote. Sens.2
2021 Battle of the Defaults: Extracting Performance Characteristics of HDF5 under Production Load
abstract
Popular parallel I/O libraries, such as HDF5, provide tuning parameters to obtain superior performance. However, the selection of effective parameters on production systems is complex due to the interdependence of I/O software and file system layers. Hence, application developers typically use the default parameters and often experience poor I/O performance. This work conducts a benchmarking-based analysis on the HDF5 behaviors with a wide variety of I/O patterns to extract performance characteristics under the production workload. To make the analysis well controlled, we exercise I/O benchmarks on POSIX-IO, MPI-IO, and HDF5 using the same I/O patterns and in the same jobs. To address high performance variability in production environments, we repeat the benchmarks across I/O patterns, storage devices, and time intervals. Based on the results, we identified consistent HDF5 behaviors that appropriate configurations and operations on dataset layout and file-metadata placement can improve performance significantly. We apply our findings and evaluate the tuned I/O library on two supercomputers: Summit and Cori. The results show that our tuned parameters can achieve more than 10× I/O performance speedup than that with default parameters on both systems, suggesting the effectiveness, stability, and generality of our solution.
Houjun Tang, Surendra Byna, Jesse Hanley, Quincey Koziol, Tonglin Li, Sarp Oral
CCGRID6
2020 Reflector: a fine-grained I/O tracker for HPC systems
abstract
We present Reflector, to support both high-level and low-level I/O monitoring through user-defined interfaces such as HDF5 and NetCDF in addition to POSIX- and MPI-IO. We evaluate Reflector on both an on-premises 500-core HPC cluster and a leadership-class supercomputer at the Lawrence Berkeley National Laboratory. Preliminary results are promising as the system prototype incurs negligible performance overhead and clearly illustrates the I/O patterns and bottlenecks of multiple applications.
Abdullah Al-Mamun 0001, Jialin Liu 0002, Tonglin Li, Quincey Koziol, Zhongyi Zhai, Junyan Qian, Haoting Shen, Dongfang Zhao 0001
PPoPP3
2020 Foreword to the special issue of the workshop on data-intensive computing in the clouds
abstract
The purpose of this special issue is to collect a selection of representative research articles that were primarily presented at the Eighth Workshop on Data-Intensive Computing in the Clouds, held in conjunction with SC'17. In particular, this annual workshop brings together domain scientists, researchers, scholars, vendors, and practitioners from the complementary fields of data science, cloud computing, and high-performance computing, in order to promote an exchange of ideas, discuss future collaborations, and develop new research directions. Data scientists increasingly rely on high-performance computers and cloud infrastructures to analyze high volumes of scientific data, automatically process data, and manage data privacy and performance. As scientific data continues to grow in volume and complexity, computational capabilities also increase at both supercomputing facilities and industry data center. Processing scientific data on emerging hardware with high performance-efficiency and privacy requires a knowledge combination from specific scientific domains and computer system techniques. This special issue presents examples of the successful collaboration from domain scientists, researchers, scholars, vendors and practitioners to address the research challenges on processing scientific data on the infrastructures of high-performance computers and cloud. The scope of this special issue is representative of the multidisciplinary nature of scientific computing in high-performance computing and cloud computing. This special issue addresses the challenges on practical experiences on processing scientific data in different domains and different platforms in high-performance computers and cloud. In particular, Zamani et al.1 show how to automatically integrate large-scale facilities with cyberinfrastructure services for automated data processing. Peng and Plale2 identify the requirements for managing computational analysis among candidate storage solutions. Koulouzis et al.3 describe their experiences that identify the time-critical requirements of environmental scientists making use of computational research support environments and provide a case study whereby their software suite is used to optimize runtime service quality for a data subscription service. Dayarathna and Suzumura4 demonstrate their approach on producing optimized stream query performance, and further compare the solution to naive deployments using two real-world stream processing applications in the domains of health care and search advertising. We encourage the readers to review the aforementioned articles to gain insight into the breadth and depth of problems and innovative solutions in the multidisciplinary field of scientific data in high-performance computing and cloud computing.
Tonglin Li
Concurr. Comput. Pract. Exp.1
2018 In-memory Blockchain: Toward Efficient and Trustworthy Data Provenance for HPC Systems
abstract
The state-of-the-art approaches for tracking data provenance on high-performance computing (HPC) systems are either supported by file systems or relational databases. These techniques shared the same critique on the provenance data’s fidelity and the associated I/O overhead. This paper envisions to track the HPC data provenance using a distributed in-memory ledger—the core technique leveraged by blockchains and proven to be highly trustworthy by many large-scale applications. We pinpoint two system challenges—storage architecture and consensus protocol—for adopting blockchains to HPC and make the following contributions: (i) We design a new in-memory blockchain architecture for HPC systems, exploiting the high-performance network infrastructure InfiniBand and greatly reducing the I/O overhead; and (ii) We develop a new consensus protocol, namely proof-of-reproducibility (PoR), crafted for the new architecture, which takes into account both proof-of-work (PoW) and proof-of-stake (PoS) mechanisms. The correctness of PoR is both theoretically proven and experimentally verified. A prototype system is implemented and evaluated with more than one million transactions, showing 32× speedup compared to the filesystem-based provenance service and four orders of magnitude speedup compared to the database-based provenance service.
Abdullah Al-Mamun 0001, Tonglin Li, Mohammad Sadoghi, Dongfang Zhao 0001
IEEE BigData2
2018 Toward Scalable Analysis of Multidimensional Scientific Data: A Case Study of Electrode Arrays
abstract
Many modern scientific applications involve large volumes of multidimensional data and extensive computation. Although distributed systems and tools are becoming increasingly scalable, they are still far away to catch up the exponential growth rate exhibited by many of those scientific big-data applications. This paper presents our early effort on overcoming the exponential complexity of one widely deployed workload over multidimensional scientific data—the n×n numerical analysis on two-dimensional arrays. More specifically, we propose a new approach to reduce the exponentially-grown data into a semantically-equivalent polynomial form in the context of two-dimensional electrode arrays, which are widely used in biomedical engineering, electrical engineering, and mechanical engineering. We have implemented a system prototype in Python, preliminary results show that the proposed approach outperforms the state-of-the-practice in various metrics: (i) the consumed space is six orders of magnitude smaller; (ii) the execution time is three orders of magnitude faster; and (iii) the scalability is improved by two orders of magnitude—from 6×6 to 100 × 100—on mainstream servers in reasonable time.
Ye Niu, Abdullah Al-Mamun 0001, Tonglin Li, Yi Zhao 0004, Dongfang Zhao 0001
IEEE BigData4
2018 Toward Performant and Energy-efficient Queries in Three-tier Wireless Sensor Networks
abstract
In a wireless sensor network (WSN) where nodes are mostly battery-powered, queries' energy consumption and response time are two of the most important metrics as they represent the network's sustainability and performance, respectively. Conventional techniques used to focus only one of the two metrics and did not attempt to optimize both in a coordinated manner. This work aims to achieve both high sustainability and high performance of WSN queries at the same time. To that end, a new mechanism is proposed to construct the topology of a three-tier WSN. The proposed mechanism eliminates routing tables and employs a novel and efficient addressing scheme inspired by the Chinese Remainder Theorem (CRT). The CRT-based topology allows for query parallelism, an unprecedented feature in WSNs. On top of the new topology encoded by CRT, a new protocol is designed to parallelly preprocess collected data on sensor nodes by effectively aggregating and deduplicating data in a neighborhood cluster. Moreover, a new algorithm is devised to allow the queries and results to be transmitted through low-power and fault-tolerant paths using recursive elections over a subset of the entire power range. With all these new techniques combined, the proposed system outperforms the state-of-the-art from various perspectives: (i) the query response is improved by up to 53%; (ii) the energy consumption is reduced by up to 70%; and (iii) the reliability is increased by up to 39%.
Jiayao Wang 0004, Abdullah Al-Mamun 0001, Tonglin Li, Linhua Jiang, Dongfang Zhao 0001
ICPP3
2017 Toward Efficient and Flexible Metadata Indexing of Big Data Systems
abstract
In Big Data era, applications are generating orders of magnitude more data in both volume and quantity. While many systems emerge to address such data explosion, the fact that these data's descriptors, i.e., metadata, are also “big” is often overlooked. The conventional approach to address the big metadata issue is to disperse metadata into multiple machines. However, it is extremely difficult to preserve both load-balance and data-locality in this approach. To this end, in this work we propose hierarchical indirection layers for indexing the underlying distributed metadata. By doing this, data locality is achieved efficiently by the indirection while load-balance is preserved. Three key challenges exist in this approach, however: first, how to achieve high resilience; second, how to ensure flexible granularity; third, how to restrain performance overhead. To address above challenges, we design Dindex, a distributed indexing service for metadata. Dindex incorporates a hierarchy of coarse-grained aggregation and horizontal key-coalition. Theoretical analysis shows that the overhead of building Dindex is compensated by only two or three queries. Dindex has been implemented by a lightweight distributed key-value store and integrated to a fully-fledged distributed filesystem. Experiments demonstrated that Dindex accelerated metadata queries by up to 60 percent with a negligible overhead.
Dongfang Zhao 0001, Kan Qiao, Zhou Zhou 0006, Tonglin Li, Zhihan Lyu, Xiaohua Xu 0002
IEEE Trans. Big Data4
2017 Understanding the Performance and Potential of Cloud Computing for Scientific Applications
abstract
Commercial clouds bring a great opportunity to the scientific computing area. Scientific applications usually require significant resources, however not all scientists have access to sufficient high-end computing systems. Cloud computing has gained the attention of scientists as a competitive resource to run HPC applications at a potentially lower cost. But as a different infrastructure, it is unclear whether clouds are capable of running scientific applications with a reasonable performance per money spent. This work provides a comprehensive evaluation of EC2 cloud in different aspects. We first analyze the potentials of the cloud by evaluating the raw performance of different services of AWS such as compute, memory, network and I/O. Based on the findings on the raw performance, we then evaluate the performance of the scientific applications running in the cloud. Finally, we compare the performance of AWS with a private cloud, in order to find the root cause of its limitations while running scientific applications. This paper aims to assess the ability of the cloud to perform well, as well as to evaluate the cost of the cloud in terms of both raw performance and scientific applications performance. Furthermore, we evaluate other services including S3, EBS and DynamoDB among many AWS services in order to assess the abilities of those to be used by scientific applications and frameworks. We also evaluate a real scientific computing application through the Swift parallel scripting system at scale. Armed with both detailed benchmarks to gauge expected performance and a detailed monetary cost analysis, we expect this paper will be a recipe cookbook for scientists to help them decide where to deploy and run their scientific applications between public clouds, private clouds, or hybrid clouds.
Iman Sadooghi, Jesus Hernandez Martin, Tonglin Li, Kevin Brandstatter, Ketan Maheshwari, Tiago Pais Pitta De Lacerda Ruivo, Gabriele Garzoglio, Steven Timm, Yong Zhao 0009, Ioan Raicu
IEEE Trans. Cloud Comput.3
2016 Albatross: An efficient cloud-enabled task scheduling and execution framework using distributed message queues
abstract
Data Analytics has become very popular on large datasets in different organizations. It is inevitable to use distributed resources such as Clouds for Data Analytics and other types of data processing at larger scales. To effectively utilize all system resources, an efficient scheduler is needed, but the traditional resource managers and job schedulers are centralized and designed for larger batch jobs which are fewer in number. Frameworks such as Hadoop and Spark, which are mainly designed for Big Data analytics, have been able to allow for more diversity in job types to some extent. However, even these systems have centralized architectures and will not be able to perform well on large scales and under heavy task loads. Modern applications generate tasks at very high rates that can cause significant slowdowns on these frameworks. Additionally, over-decomposition has shown to be very useful in increasing the system utilization. In order to achieve high efficiency, scalability, and better system utilization, it is critical for a modern scheduler to be able to handle over-decomposition and run highly granular tasks. Further, to achieve high performance, Albatross is written in C/C++, which imposes a minimal overhead to the workload process as compared to languages like Java or Python. We propose Albatross, a task level scheduling and execution framework that uses a Distributed Message Queue (DMQ) for task distribution among its workers. Unlike most scheduling systems, Albatross uses a pulling approach as opposed to the common push approach. The former would let Albatross achieve a good load balancing and scalability. Furthermore, the framework has built in support for task execution dependency on workflows. Therefore, Albatross is able to run various types of workloads, including Data Analytics and HPC applications. Finally, Albatross provides data locality support. This allows the framework to achieve higher performance through minimizing the amount of unnecessary data movement on the network. Our evaluations show that Albatross outperforms Spark and Hadoop at larger scales and in the case of running higher granularity workloads.
Iman Sadooghi, Geet Kumar, Ke Wang 0012, Dongfang Zhao 0001, Tonglin Li, Ioan Raicu
eScience5
2016 A convergence of key-value storage systems from clouds to supercomputers
abstract
Summary This paper presents a convergence of distributed key‐value storage systems in clouds and supercomputers. It specifically presents ZHT, a zero‐hop distributed key‐value store system, which has been tuned for the requirements of high‐end computing systems. ZHT aims to be a building block for future distributed systems, such as parallel and distributed file systems, distributed job management systems, and parallel programming systems. ZHT has some important properties, such as being lightweight, dynamically allowing nodes join and leave, fault tolerant through replication, persistent, scalable, and supporting unconventional operations such as append, compare and swap, callback in addition to the traditional insert/lookup/remove. We have evaluated ZHT's performance under a variety of systems, ranging from a Linux cluster with 64 nodes, an Amazon EC2 virtual cluster up to 96 nodes, to an IBM Blue Gene/P supercomputer with 8K nodes. We compared ZHT against other key‐value stores and found it offers superior performance for the features and portability it supports. This paper also presents several real systems that have adopted ZHT, namely, FusionFS (a distributed file system), IStore (a storage system with erasure coding), MATRIX (distributed scheduling), Slurm++ (distributed HPC job launch), Fabriq (distributed message queue management); all of these real systems have been simplified because of key‐value storage systems and have been shown to outperform other leading systems by orders of magnitude in some cases. It is important to highlight that some of these systems are rooted in HPC systems from supercomputers, while others are rooted in clouds and ad hoc distributed systems; through our work, we have shown how versatile key‐value storage systems can be in such a variety of environments. Copyright © 2015 John Wiley & Sons, Ltd.
Tonglin Li, Xiaobing Zhou, Ke Wang 0012, Dongfang Zhao 0001, Iman Sadooghi, Zhao Zhang 0007, Ioan Raicu
Concurr. Comput. Pract. Exp.1
2016 Load-balanced and locality-aware scheduling for data-intensive workloads at extreme scales
abstract
Summary Data‐driven programming models such as many‐task computing (MTC) have been prevalent for running data‐intensive scientific applications. MTC applies over‐decomposition to enable distributed scheduling. To achieve extreme scalability, MTC proposes a fully distributed task scheduling architecture that employs as many schedulers as the compute nodes to make scheduling decisions. Achieving distributed load balancing and best exploiting data locality are two important goals for the best performance of distributed scheduling of data‐intensive applications. Our previous research proposed a data‐aware work‐stealing technique to optimize both load balancing and data locality by using both dedicated and shared task ready queues in each scheduler. Tasks were organized in queues based on the input data size and location. Distributed key‐value store was applied to manage task metadata. We implemented the technique in MATRIX, a distributed MTC task execution framework. In this work, we devise an analytical suboptimal upper bound of the proposed technique, compare MATRIX with other scheduling systems, and explore the scalability of the technique at extreme scales. Results show that the technique is not only scalable but can achieve performance within 15% of the suboptimal solution. Copyright © 2015 John Wiley & Sons, Ltd.
Ke Wang 0012, Kan Qiao, Iman Sadooghi, Xiaobing Zhou, Tonglin Li, Michael Lang 0003, Ioan Raicu
Concurr. Comput. Pract. Exp.5
2016 Exploiting multi-cores for efficient interchange of large messages in distributed systems
abstract
Summary Conventional data serialization tools assume that objects to be coded are usually small in size so a single CPU core can encode it in a timely manner. In the era of Big Data, however, object gets increasingly complex and larger, which makes data serialization become a new performance bottleneck. This paper describes an approach to parallelize data serialization by leveraging multiple cores. Parallelizing data serialization introduces new questions such as how to split the (sub)objects, how to allocate the available cores, and how to minimize its overhead in practice. In this paper we design a framework for parallelly serializing large objects and analyze the design tradeoffs under different scenarios. To validate the proposed approach, we implemented parallel protocol buffers—the parallel version of Google's Protocol Buffers, a widely‐used data serialization utility. Experimental results confirm the effectiveness of Parallel Protocol Buffers: multiple cores employed in data serialization achieve highly scalable performance and incur negligible overhead. Copyright © 2015 John Wiley & Sons, Ltd.
Dongfang Zhao 0001, Kan Qiao, Zhou Zhou 0006, Tonglin Li, Xiaobing Zhou, Ioan Raicu
Concurr. Comput. Pract. Exp.4
2016 Toward high-performance key-value stores through GPU encoding and locality-aware encoding
Dongfang Zhao 0001, Ke Wang 0012, Kan Qiao, Tonglin Li, Iman Sadooghi, Ioan Raicu
J. Parallel Distributed Comput.4
2015 A flexible QoS fortified distributed key-value storage system for the cloud
abstract
In the era of big data and cloud, distributed key-value stores are increasingly used as building blocks of large-scale applications. Comparing to traditional relational databases, key-value stores are particularly compelling due to their low latency and excellent scalability. Many big companies, such as Facebook and Amazon, run multiple different applications and services on top of a single key-value store deployment to reduce the deployment and maintenance complexity as well as economic cost. However, every application has its performance requirement but most current key-value store systems are designed to serve every application request equally. This design works well when a single application accesses the key-value store, but it is not as good for the emerging concurrent multi-application scenario. In this paper, we present ZHT/Q, a flexible QoS (Quality of Service) fortified distributed key-value storage system for clouds and data centers. It improves the overall throughput by an order of magnitude and still satisfies different applications' latency requirements with QoS using dynamic and adaptive request batching mechanisms. The experiment results show that our new system delivers up to 28 times higher throughput than the base solution while more than 99% of requests' latency requirements are satisfied.
Tonglin Li, Ke Wang 0012, Dongfang Zhao 0001, Kan Qiao, Iman Sadooghi, Xiaobing Zhou, Ioan Raicu
IEEE BigData1
2015 MHT: A light-weight scalable zero-hop MPI enabled distributed key-value store
abstract
In this paper, we propose and implement a key-value store that supports MPI while allowing application access at any time without having to declaring in the same MPI communication world. This feature may significantly simplify the application design and allow programmers leverage the power of key-value store in an intuitive way. In our preliminary experiment results captured from a supercomputer at Los Alamos National Laboratory, our prototype shows linear scalability at up to 256 nodes.
Xiaobing Zhou, Tonglin Li, Ke Wang 0012, Dongfang Zhao 0001, Iman Sadooghi, Ioan Raicu
IEEE BigData2
2015 GRAPH/Z: A Key-Value Store Based Scalable Graph Processing System
abstract
The emerging applications in big data and social networks issue rapidly increasing demands on graph processing. Graph query operations that involve a large number of vertices and edges can be tremendously slow on traditional databases. The state-of-the-art graph processing systems and databases usually adopt master/slave architecture that potentially impairs their The contributions of this paper are as follows: scalability. This work describes the design and implementation of a new graph processing system based on Bulk Synchronous Parallel model. Our system is built on top of ZHT, a scalable distributed key-value store, which benefits the graph processing in terms of scalability, performance and persistency. The experiment results imply excellent scalability.
Tonglin Li, Chaoqi Ma, Xiaobing Zhou, Ke Wang 0012, Dongfang Zhao 0001, Iman Sadooghi, Ioan Raicu
CLUSTER1
2015 Overcoming Hadoop Scaling Limitations through Distributed Task Execution
abstract
Data driven programming models like MapReduce have gained the popularity in large-scale data processing. Although great efforts through the Hadoop implementation and framework decoupling (e.g. YARN, Mesos) have allowed Hadoop to scale to tens of thousands of commodity cluster processors, the centralized designs of the resource manager, task scheduler and metadata management of HDFS file system adversely affect Hadoop's scalability to tomorrow's extreme-scale data centers. This paper aims to address the YARN scaling issues through a distributed task execution framework, MATRIX, which was originally designed to schedule the executions of data-intensive scientific applications of many-task computing on supercomputers. We propose to leverage the distributed design wisdoms of MATRIX to schedule arbitrary data processing applications in cloud. We compare MATRIX with YARN in processing typical Hadoop workloads, such as WordCount, TeraSort, Grep and RandomWriter, and the Ligand application in Bioinformatics on the Amazon Cloud. Experimental results show that MATRIX outperforms YARN by 1.27X for the typical workloads, and by 2.04X for the real application. We also run and simulate MATRIX with fine-grained sub-second workloads. With the simulation results giving the efficiency of 86.8% at 64K cores for the 150ms workload, we show that MATRIX has the potential to enable Hadoop to scale to extreme-scale data centers for fine-grained workloads.
Ke Wang 0012, Ning Liu 0008, Iman Sadooghi, Xi Yang 0002, Xiaobing Zhou, Tonglin Li, Michael Lang 0003, Xian-He Sun, Ioan Raicu
CLUSTER6
2014 Optimizing load balancing and data-locality with data-aware scheduling
abstract
Load balancing techniques (e.g. work stealing) are important to obtain the best performance for distributed task scheduling systems that have multiple schedulers making scheduling decisions. In work stealing, tasks are randomly migrated from heavy-loaded schedulers to idle ones. However, for data-intensive applications where tasks are dependent and task execution involves processing a large amount of data, migrating tasks blindly yields poor data-locality and incurs significant data-transferring overhead. This work improves work stealing by using both dedicated and shared queues. Tasks are organized in queues based on task data size and location. We implement our technique in MATRIX, a distributed task scheduler for many-task computing. We leverage distributed key-value store to organize and scale the task metadata, task dependency, and data-locality. We evaluate the improved work stealing technique with both applications and micro-benchmarks structured as direct acyclic graphs. Results show that the proposed data-aware work stealing technique performs well.
Ke Wang 0012, Xiaobing Zhou, Tonglin Li, Dongfang Zhao 0001, Michael Lang 0003, Ioan Raicu
IEEE BigData3
2014 FusionFS: Toward supporting data-intensive scientific applications on extreme-scale high-performance computing systems
abstract
State-of-the-art, yet decades-old, architecture of high-performance computing systems has its compute and storage resources separated. It thus is limited for modern data-intensive scientific applications because every I/O needs to be transferred via the network between the compute and storage resources. In this paper we propose an architecture that hss a distributed storage layer local to the compute nodes. This layer is responsible for most of the I/O operations and saves extreme amounts of data movement between compute and storage resources. We have designed and implemented a system prototype of this architecture - which we call the FusionFS distributed file system - to support metadata-intensive and write-intensive operations, both of which are critical to the I/O performance of scientific applications. FusionFS has been deployed and evaluated on up to 16K compute nodes of an IBM Blue Gene/P supercomputer, showing more than an order of magnitude performance improvement over other popular file systems such as GPFS, PVFS, and HDFS.
Dongfang Zhao 0001, Zhao Zhang 0007, Xiaobing Zhou, Tonglin Li, Ke Wang 0012, Dries Kimpe, Philip H. Carns, Robert B. Ross, Ioan Raicu
IEEE BigData4
2013 ZHT: A Light-Weight Reliable Persistent Dynamic Scalable Zero-Hop Distributed Hash Table
abstract
This paper presents ZHT, a zero-hop distributed hash table, which has been tuned for the requirements of high-end computing systems. ZHT aims to be a building block for future distributed systems, such as parallel and distributed file systems, distributed job management systems, and parallel programming systems. The goals of ZHT are delivering high availability, good fault tolerance, high throughput, and low latencies, at extreme scales of millions of nodes. ZHT has some important properties, such as being light-weight, dynamically allowing nodes join and leave, fault tolerant through replication, persistent, scalable, and supporting unconventional operations such as append (providing lock-free concurrent key/value modifications) in addition to insert/lookup/remove. We have evaluated ZHT's performance under a variety of systems, ranging from a Linux cluster with 512-cores, to an IBM Blue Gene/P supercomputer with 160K-cores. Using micro-benchmarks, we scaled ZHT up to 32K-cores with latencies of only 1.1ms and 18M operations/sec throughput. This work provides three real systems that have integrated with ZHT, and evaluate them at modest scales. 1) ZHT was used in the FusionFS distributed file system to deliver distributed meta-data management at over 60K operations (e.g. file create) per second at 2K-core scales. 2) ZHT was used in the IStore, an information dispersal algorithm enabled distributed object storage system, to manage chunk locations, delivering more than 500 chunks/sec at 32-nodes scales. 3) ZHT was also used as a building block to MATRIX, a distributed job scheduling system, delivering 5000 jobs/sec throughputs at 2K-core scales. We compared ZHT against other distributed hash tables and key/value stores and found it offers superior performance for the features and portability it supports.
Tonglin Li, Xiaobing Zhou, Kevin Brandstatter, Dongfang Zhao 0001, Ke Wang 0012, Anupam Rajendran, Zhao Zhang 0007, Ioan Raicu
IPDPS1