VLDB 2026 Research / reviewers in the wild / expert
Xiaosong Ma
dblp:m/XiaosongMa
· DBLP profile ↗
93ranked-venue papers
9as first author
23since 2021 · last 2026
—ORCID · conflict
Domains — the database's venue-derived domains; a paper can count in several
Systems, architecture and hardware · 69 · 5 first-author · 11 since 2021Databases, data management, data science and information retrieval · 13 · 1 first-author · 4 since 2021Artificial intelligence and machine learning · 8 · 2 first-author · 7 since 2021Software engineering, systems software and programming languages · 7 · 3 since 2021Computer networks · 4 · 1 first-author · 2 since 2021Security and privacy · 1Graphics, computer vision, multimedia, augmented reality and games · 1 · 1 first-author · 1 since 2021Applied, interdisciplinary, general and emerging computing · 1
| Year | Publication | Venue | Position |
|---|---|---|---|
| 2026 | Think How to Think: Mitigating Overthinking with Autonomous Difficulty Cognition in Large Reasoning ModelsabstractRecent Large Reasoning Models (LRMs) excel at complex reasoning tasks but often suffer from overthinking, generating overly long and redundant reasoning trajectories.To explore its essence, our empirical analysis reveals that LRMs are primarily limited to recognizing task properties (i.e., difficulty levels) like humans before solving the problem, leading to a one-size-fits-all reasoning strategy.This observation motivates a fundamental question: Can we explicitly bootstrap such ability to alleviate overthinking in LRMs?To this end, we propose Think-How-to-Think (TH2T), a novel two-stage fine-tuning strategy that progressively inspires LRMs' difficulty cognition and redundancy cognition of LRMs.Specifically, we first inject Difficulty Dypnosis into output prefixes as cues for global, prospective reasoning strategy selection, stimulating the model's sharper sensitivity to task complexity and adaptive control of reasoning depth.Then, we incorporate Redundancy Hypnosis into inprogress reasoning steps, which serve as local, retrospective signals for behavior correction by identifying and eliminating superfluous reasoning detours.Experiments across 7B/14B/32B models demonstrate that TH2T significantly reduces inference costs by over 70% on easy tasks and 40% on complex ones without compromising performance.The resultant models exhibit a nascent ability for difficulty-aware reasoning, effectively mitigating behaviors like excessive reflection and looping, thereby paving the way for more cognitively efficient LRMs. Yongjiang Liu, Haoxi Li, Xiaosong Ma, Jie Zhang 0076, Song Guo 0001 |
ACL (1) | 3 |
| 2026 | LazyEviction: Lagged KV Eviction with Attention Pattern Observation for Efficient Long ReasoningabstractLarge Language Models (LLMs) exhibit enhanced capabilities by Chain-of-Thought reasoning.However, the extended reasoning sequences introduce significant GPU memory overhead due to increased key-value (KV) cache.Existing KV cache compression methods mitigate memory bottlenecks but struggle in long reasoning tasks.In this paper, we analyze attention patterns in reasoning tasks and reveal a Token Importance Recurrence phenomenon: a large proportion of tokens regain high attention after multiple decoding steps, which is failed to capture by existing works and may lead to unpredictable eviction on such periodically critical tokens.To address this, we propose LazyEviction, an observation windowbased lagged eviction framework retaining latent recurring tokens by an eviction policy informed by token recurrence patterns.Extensive experiments demonstrate that LazyEviction reduces KV cache by 50%~70% while maintaining comparable accuracy, outperforming existing KV cache compression baselines.Our implementation code can be found at https: //github.com/Halo-949/LazyEviction. Xiaosong Ma, Jie Zhang 0076, Song Guo 0001 |
ACL (1) | 3 |
| 2026 | HCT-QA: A Benchmark for Question Answering on Human-Centric TablesabstractTabular data embedded in PDF files, web pages, and other types of documents is prevalent in various domains. These tables, which we call human-centric tables (HCTs for short), are dense in information but often exhibit complex structural and semantic layouts. To query these HCTs, some existing solutions focus on transforming them into relational formats. However, they fail to handle the diverse and complex layouts of HCTs, making them not amenable to easy querying with SQL-based approaches. Another emerging option is to use Large Language Models (LLMs) and Vision Language Models (VLMs). However, there is a lack of standard evaluation benchmarks to measure and compare the performance of models to query HCTs using natural language. To address this gap, we propose the HumanCentric Tables Question-Answering extensive benchmark (HCTQA) consisting of thousands of HCTs with several thousands of natural language questions with their respective answers. More specifically, HCT-QA includes 1,880 real-world HCTs with 9,835 QA pairs in addition to 4,679 synthetic HCTs with 67.7K QA pairs. Also, we show through extensive experiments the performance of 25 and 9 different LLMS and VLMs, respectively, in an answering HCT-QA's questions. In addition, we show how finetuning an LLM on HCT-QA improves F1 scores by up to 25 percentage points compared to the off-the-shelf model. Compared to existing benchmarks, HCT-QA stands out for its broad complexity and diversity of covered HCTs and generated questions, its comprehensive metadata enabling deeper insight and analysis, and its novel synthetic data and QA generator. Mohammad Shahmeer Ahmad, Zan Ahmad Naeem, Michaël Aupetit 0001, Ahmed K. Elmagarmid, Mohamed Y. Eltabakh, Xiaosong Ma, Mourad Ouzzani, Chaoyi Ruan, Hani Al-Sayeh |
ICDE | 6 |
| 2026 | Real-Time Language Model Jamming: A Case Study for Live Music Accompaniment Generation
Andrew H. Yang, Jiaqi Ruan, Yuan-Hsin Chen, Xiaosong Ma |
RTAS | 8 |
| 2025 | Understanding Data Preprocessing for Effective End-to-End Training of DNN
Ping Gong 0009, Cheng Li 0001, Xiaosong Ma, Sam H. Noh |
APPT | 4 |
| 2025 | DHeLlam: General-Purpose, Automatic Micro-Batch Co-Execution for Distributed LLM TrainingabstractThe growth of Large Language Models (LLMs) has necessitated large-scale distributed training. Highly optimized frameworks, however, suffer significant losses in MFU (Model FLOPS Utilization) due to communication. This paper introduces DHeLlam, a novel micro-structure inspired by DNA that significantly enhances the efficiency of LLM training. Central to DHeLlam is Strand Interleaving (SI), which treats the continuous stream of training micro-batches on a GPU as two interleaved strands. DHeLlam co-schedules their forward and backward passes using operator-level overlap profiling and a dynamic programming-based search. It enables the two strands to share model states and activation memory, requiring$<3 \%$additional HBM space under common model configurations. To our best knowledge, DHeLlam is the first to co-execute two microbatches without requiring model replication. With its unique model folding design, DHeLlam seamlessly integrates with all forms of data and model parallelism, including the challenging pipeline parallelism (with a W-shaped pipeline). We evaluated DHeLlam training with the popular Llama and GPT dense models, plus the Phi Mixture of Expert (MoE) model, across 2 GPU clusters. Results show that it achieves 12-40% throughput (up to 58% MFU) and 5-24% throughput (up to 64% MFU) improvement on the 64-card A40 and A800 clusters respectively, significantly outperforming state-of-the-art methods. Chaoyi Ruan, Jiaqi Ruan, Chengjie Tang, Xiaosong Ma, Cheng Li 0001 |
ICCD | 6 |
| 2025 | Model Decomposition and Reassembly for Purified Knowledge Transfer in Personalized Federated LearningabstractPersonalized federated learning (pFL) is to collaboratively train non-identical machine learning models for different clients to adapt to their heterogeneously distributed datasets. State-of-the-art pFL approaches pay much attention on exploiting clients’ inter-similarities to facilitate the collaborative learning process, meanwhile, can barely escape from the irrelevant knowledge pooling that is inevitable during the aggregation phase, and thus hindering the optimization convergence and degrading the personalization performance. To tackle such conflicts between facilitating collaboration and promoting personalization, we propose a novel pFL framework, dubbed pFedC, to first decompose the global aggregated knowledge into several compositional branches, and then selectively reassemble the relevant branches for supporting conflicts-aware collaboration among contradictory clients. Specifically, by reconstructing each local model into a shared feature extractor and multiple decomposed task-specific classifiers, the training on each client transforms into a mutually reinforced and relatively independent multi-task learning process, which provides a new perspective for pFL. Besides, we conduct a purified knowledge aggregation mechanism via quantifying the combination weights for each client to capture clients’ common prior, as well as mitigate potential conflicts from the divergent knowledge caused by the heterogeneous data. Extensive experiments over various models and datasets demonstrate the effectiveness and superior performance of the proposed algorithm. Jie Zhang 0076, Song Guo 0001, Xiaosong Ma, Wenchao Xu 0001, Qihua Zhou, Jingcai Guo, Zicong Hong, Jun Shan |
IEEE Trans. Mob. Comput. | 3 |
| 2025 | Introduction to the Special Section on USENIX FAST 2024
Xiaosong Ma, Youjip Won |
ACM Trans. Storage | 1 |
| 2024 | Amend to Alignment: Decoupled Prompt Tuning for Mitigating Spurious Correlation in Vision-Language ModelsabstractFine-tuning the learnable prompt for a pre-trained vision-language model (VLM), such as CLIP, has demonstrated exceptional efficiency in adapting to a broad range of downstream tasks. Existing prompt tuning methods for VLMs do not distinguish spurious features introduced by biased training data from invariant features, and employ a uniform alignment process when adapting to unseen target domains. This can impair the cross-modal feature alignment when the testing data significantly deviate from the distribution of the training data, resulting in a poor out-of-distribution (OOD) generalization performance. In this paper, we reveal that the prompt tuning failure in such OOD scenarios can be attribute to the undesired alignment between the textual and the spurious feature. As a solution, we propose **CoOPood**, a fine-grained prompt tuning method that can discern the causal features and deliberately align the text modality with the invariant feature. Specifically, we design two independent contrastive phases using two lightweight projection layers during the alignment, each with different objectives: 1) pulling the text embedding closer to invariant image embedding and 2) pushing text embedding away from spurious image embedding. We have illustrated that **CoOPood** can serve as a general framework for VLMs and can be seamlessly integrated with existing prompt tuning methods. Extensive experiments on various OOD datasets demonstrate the performance superiority over state-of-the-art methods. Jie Zhang 0076, Xiaosong Ma, Song Guo 0001, Peng Li 0017, Wenchao Xu 0001, Xueyang Tang, Zicong Hong |
ICML | 2 |
| 2024 | PolyBase: Adapting to Data Affinity Changes in Geo-Replicated Database via Row-Level Paxos-Group Affiliation Re-AssignmentabstractTransaction performance in geo-replicated databases heavily relies on the request location: when not issued by the primary region, transactions are forced to involve costly wide-area communication. While existing systems distribute primary roles across regions, such assignment typically occurs at the shard level, making it difficult to align with geographically dispersed access to individual records. This paper introduces PolyBase, a pioneering architecture to address such misalignment, leveraging the widely adopted Paxos-based log replication mechanisms. It enables flexible row-level consensus group affiliation , which runs on an unchanged Paxos protocol , but dynamically re-assigns database rows between Paxos log replication groups, whose leaders become the primary region, enjoying faster writes and up-to-date versions for reads. With carefully designed data structures and protocols, PolyBase significantly reduces wide-area RTTs without compromising transaction or log replication consistency or reliability guarantees. We implemented PolyBase with optimized re-assignment policies and integrated it into two popular databases (RocksDB and MySQL). Our evaluation on AWS, using a production e-commerce workload and microbench-marks confirms that PolyBase offers significantly higher transaction throughput and lower average/tail latency compared to baselines. Chaoyi Ruan, Yingqiang Zhang, Juncheng Zhang, Cheng Li 0001, Xiaosong Ma, Hao Chen 0080, Feifei Li 0001, Xinjun Yang |
Proc. VLDB Endow. | 5 |
| 2023 | GZKP: A GPU Accelerated Zero-Knowledge Proof SystemabstractZero-knowledge proof (ZKP) is a cryptographic protocol that allows one party to prove the correctness of a statement to another party without revealing any information beyond the correctness of the statement itself. It guarantees computation integrity and confidentiality, and is therefore increasingly adopted in industry for a variety of privacy-preserving applications, such as verifiable outsource computing and digital currency. Weiliang Ma, Qian Xiong, Xuanhua Shi, Xiaosong Ma, Hai Jin 0001, Haozhao Kuang, Mingyu Gao 0001, Ye Zhang 0042, Haichen Shen, Weifang Hu |
ASPLOS (2) | 4 |
| 2023 | Persistent Memory Disaggregation for Cloud-Native Relational DatabasesabstractThe recent emergence of commodity persistent memory (PM) hardware has altered the landscape of the storage hierarchy. It brings multi-fold benefits to database systems, with its large capacity, low latency, byte addressability, and persistence. However, PM has not been incorporated into the popular disaggregated architecture of cloud-native databases. Chaoyi Ruan, Yingqiang Zhang, Chao Bi, Xiaosong Ma, Hao Chen 0080, Feifei Li 0001, Xinjun Yang, Cheng Li 0001, Ashraf Aboulnaga, Yinlong Xu 0001 |
ASPLOS (3) | 4 |
| 2023 | FrozenHot Cache: Rethinking Cache Management for Modern HardwareabstractCaching is crucial for accelerating data access, employed as a ubiquitous design in modern systems at many parts of computer systems. With increasing core count, and shrinking latency gap between cache and modern storage devices, hit-path scalability becomes increasingly critical. However, existing production in-memory caches often use list-based management with promotion on each cache hit, which requires extensive locking and poses a significant overhead for scaling beyond a few cores. Moreover, existing techniques for improving scalability either (1) only focus on the indexing structure and do not improve cache management scalability, or (2) sacrifice efficiency or miss-path scalability. Ziyue Qiu, Juncheng Yang, Juncheng Zhang, Cheng Li 0001, Xiaosong Ma, Qi Chen 0009, Mao Yang 0004, Yinlong Xu 0001 |
EuroSys | 5 |
| 2023 | Towards Unbiased Training in Federated Open-world Semi-supervised LearningabstractFederated Semi-supervised Learning (FedSSL) has emerged as a new paradigm for allowing distributed clients to collaboratively train a machine learning model over scarce labeled data and abundant unlabeled data. However, existing works for FedSSL rely on a closed-world assumption that all local training data and global testing data are from seen classes observed in the labeled dataset. It is crucial to go one step further: adapting FL models to an open-world setting, where unseen classes exist in the unlabeled data. In this paper, we propose a novel Federatedopen-world Semi-Supervised Learning (FedoSSL) framework, which can solve the key challenge in distributed and open-world settings, i.e., the biased training process for heterogeneously distributed unseen classes. Specifically, since the advent of a certain unseen class depends on a client basis, the locally unseen classes (exist in multiple clients) are likely to receive differentiated superior aggregation effects than the globally unseen classes (exist only in one client). We adopt an uncertainty-aware suppressed loss to alleviate the biased training between locally unseen and globally unseen classes. Besides, we enable a calibration module supplementary to the global aggregation to avoid potential conflicting knowledge transfer caused by inconsistent data distribution among different clients. The proposed FedoSSL can be easily adapted to state-of-the-art FL methods, which is also validated via extensive experiments on benchmarks and real-world datasets (CIFAR-10, CIFAR-100 and CINIC-10). Jie Zhang 0076, Xiaosong Ma, Song Guo 0001, Wenchao Xu 0001 |
ICML | 2 |
| 2023 | SwapPrompt: Test-Time Prompt Adaptation for Vision-Language ModelsabstractTest-time adaptation (TTA) is a special and practical setting in unsupervised domain adaptation, which allows a pre-trained model in a source domain to adapt to unlabeled test data in another target domain. To avoid the computation-intensive backbone fine-tuning process, the zero-shot generalization potentials of the emerging pre-trained vision-language models (e.g., CLIP, CoOp) are leveraged to only tune the run-time prompt for unseen test domains. However, existing solutions have yet to fully exploit the representation capabilities of pre-trained models as they only focus on the entropy-based optimization and the performance is far below the supervised prompt adaptation methods, e.g., CoOp. In this paper, we propose SwapPrompt, a novel framework that can effectively leverage the self-supervised contrastive learning to facilitate the test-time prompt adaptation. SwapPrompt employs a dual prompts paradigm, i.e., an online prompt and a target prompt that averaged from the online prompt to retain historical information. In addition, SwapPrompt applies a swapped prediction mechanism, which takes advantage of the representation capabilities of pre-trained models to enhance the online prompt via contrastive learning. Specifically, we use the online prompt together with an augmented view of the input image to predict the class assignment generated by the target prompt together with an alternative augmented view of the same image. The proposed SwapPrompt can be easily deployed on vision-language models without additional requirement, and experimental results show that it achieves state-of-the-art test-time adaptation performance on ImageNet and nine other datasets. It is also shown that SwapPrompt can even achieve comparable performance with supervised prompt adaptation methods. Xiaosong Ma, Jie Zhang 0076, Song Guo 0001, Wenchao Xu 0001 |
NeurIPS | 1 |
| 2023 | End-to-end I/O Monitoring on Leading SupercomputersabstractThis paper offers a solution to overcome the complexities of production system I/O performance monitoring. We present Beacon, an end-to-end I/O resource monitoring and diagnosis system for the 40960-node Sunway TaihuLight supercomputer, currently the fourth-ranked supercomputer in the world. Beacon simultaneously collects and correlates I/O tracing/profiling data from all the compute nodes, forwarding nodes, storage nodes, and metadata servers. With mechanisms such as aggressive online and offline trace compression and distributed caching/storage, it delivers scalable, low-overhead, and sustainable I/O diagnosis under production use. With Beacon’s deployment on TaihuLight for more than three years, we demonstrate Beacon’s effectiveness with real-world use cases for I/O performance issue identification and diagnosis. It has already successfully helped center administrators identify obscure design or configuration flaws, system anomaly occurrences, I/O performance interference, and resource under- or over-provisioning problems. Several of the exposed problems have already been fixed, with others being currently addressed. Encouraged by Beacon’s success in I/O monitoring, we extend it to monitor interconnection networks, which is another contention point on supercomputers. In addition, we demonstrate Beacon’s generality by extending it to other supercomputers. Both Beacon codes and part of collected monitoring data are released. 1 Bin Yang 0043, Wei Xue 0003, Shichao Liu 0004, Xiaosong Ma, Xiyang Wang 0003 |
ACM Trans. Storage | 5 |
| 2022 | Layer-wised Model Aggregation for Personalized Federated LearningabstractPersonalized Federated Learning (pFL) not only can capture the common priors from broad range of distributed data, but also support customized models for heterogeneous clients. Researches over the past few years have applied the weighted aggregation manner to produce personalized models, where the weights are determined by calibrating the distance of the entire model parameters or loss values, and have yet to consider the layer-level impacts to the aggregation process, leading to lagged model convergence and inadequate personalization over non-IID datasets. In this paper, we propose a novel pFL training framework dubbed Layer-wised Personalized Federated learning (pFedLA) that can discern the importance of each layer from different clients, and thus is able to optimize the personalized model aggregation for clients with heterogeneous data. Specifically, we employ a dedicated hyper-network per client on the server side, which is trained to identify the mutual contribution factors at layer granularity. Meanwhile, a parameterized mechanism is introduced to update the layer-wised aggregation weights to progressively exploit the inter-user similarity and realize accurate model personalization. Extensive experiments are conducted over different models and learning tasks, and we show that the proposed methods achieve significantly higher performance than state-of-the-art pFL methods. Xiaosong Ma, Jie Zhang 0076, Song Guo 0001, Wenchao Xu 0001 |
CVPR | 1 |
| 2021 | SpanDB: A Fast, Cost-Effective LSM-tree Based KV Store on Hybrid Storage
Hao Chen 0080, Chaoyi Ruan, Cheng Li 0001, Xiaosong Ma, Yinlong Xu 0001 |
FAST | 4 |
| 2021 | FusionRAID: Achieving Consistent Low Latency for Commodity SSD Arrays
Tianyang Jiang, Guangyan Zhang, Zican Huang, Xiaosong Ma, Junyu Wei, Zhiyue Li |
FAST | 4 |
| 2021 | You Can Hear But You Cannot Record: Privacy Protection by Jamming Audio RecordingabstractUnauthorized voice recording via smartphones can leak the talking content stealthily. This would be a serious security threat to those individuals, enterprises and the government who need to keep the conversation confidential. Furthermore, due to the size miniaturization of smartphones, it is hard to find the covert recording from malicious attendees. Existing solutions usually jam the recording with audible noise or electromagnetic emitting. However, the audible noise will seriously interfere with conversation and the effect of electromagnetic emitting will be limited by the distance. In this paper, we propose UltraArray, a pioneering silent ultrasonic anti-recording jammer, which can covertly block recording for a long distance. The principle of covert blocking is inspired by acoustic parametric array theory, which suggests that the audible frequency wave can be spread through the air silently while it is modulated to an inaudible ultrasonic frequency. The modulation used in this paper is double sideband (DSB) modulation. The microphone on the phone will record the audible frequency and filtering out the ultrasonic frequency. The jammer we developed uses an acoustics array to form a beam to spread the signal further. The evaluation shows that the device has a good jamming effect on more than 5 meters for most Android smartphones. It will also work well with more than 2.5 meters effective distance on iPhone XR, which has the active noise control (ANC) function. Those results achieve ten times the interference ability of existing solutions. Xiaosong Ma, Yubo Song, Shang Gao 0006, Bin Xiao 0001, Aiqun Hu |
ICC | 1 |
| 2021 | Parameterized Knowledge Transfer for Personalized Federated LearningabstractIn recent years, personalized federated learning (pFL) has attracted increasing attention for its potential in dealing with statistical heterogeneity among clients. However, the state-of-the-art pFL methods rely on model parameters aggregation at the server side, which require all models to have the same structure and size, and thus limits the application for more heterogeneous scenarios. To deal with such model constraints, we exploit the potentials of heterogeneous model settings and propose a novel training framework to employ personalized models for different clients. Specifically, we formulate the aggregation procedure in original pFL into a personalized group knowledge transfer training algorithm, namely, KT-pFL, which enables each client to maintain a personalized soft prediction at the server side to guide the others' local training. KT-pFL updates the personalized soft prediction of each client by a linear combination of all local soft predictions using a knowledge coefficient matrix, which can adaptively reinforce the collaboration among clients who own similar data distribution. Furthermore, to quantify the contributions of each client to others' personalized training, the knowledge coefficient matrix is parameterized so that it can be trained simultaneously with the models. The knowledge coefficient matrix and the model parameters are alternatively updated in each round following the gradient descent way. Extensive experiments on various datasets (EMNIST, Fashion_MNIST, CIFAR-10) are conducted under different settings (heterogeneous models and data distributions). It is demonstrated that the proposed framework is the first federated learning paradigm that realizes personalized model training via parameterized group knowledge transfer while achieving significant performance gain comparing with state-of-the-art algorithms. Jie Zhang 0076, Song Guo 0001, Xiaosong Ma, Haozhao Wang, Wenchao Xu 0001, Feijie Wu |
NeurIPS | 3 |
| 2021 | Random Walks on Huge Graphs at Cache EfficiencyabstractData-intensive applications dominated by random accesses to large working sets fail to utilize the computing power of modern processors. Graph random walk, an indispensable workhorse for many important graph processing and learning applications, is one prominent case of such applications. Existing graph random walk systems are currently unable to match the GPU-side node embedding training speed. Xiaosong Ma, Saravanan Thirumuruganathan, Kang Chen 0001, Yongwei Wu 0001 |
SOSP | 2 |
| 2021 | Leveraging NVMe SSDs for Building a Fast, Cost-effective, LSM-tree-based KV StoreabstractKey-value (KV) stores support many crucial applications and services. They perform fast in-memory processing but are still often limited by I/O performance. The recent emergence of high-speed commodity non-volatile memory express solid-state drives (NVMe SSDs) has propelled new KV system designs that take advantage of their ultra-low latency and high bandwidth. Meanwhile, to switch to entirely new data layouts and scale up entire databases to high-end SSDs requires considerable investment. As a compromise, we propose SpanDB, an LSM-tree-based KV store that adapts the popular RocksDB system to utilize selective deployment of high-speed SSDs . SpanDB allows users to host the bulk of their data on cheaper and larger SSDs (and even hard disc drives with certain workloads), while relocating write-ahead logs (WAL) and the top levels of the LSM-tree to a much smaller and faster NVMe SSD. To better utilize this fast disk, SpanDB provides high-speed, parallel WAL writes via SPDK, and enables asynchronous request processing to mitigate inter-thread synchronization overhead and work efficiently with polling-based I/O. To ease the live data migration between fast and slow disks, we introduce TopFS, a stripped-down file system providing familiar file interface wrappers on top of SPDK I/O. Our evaluation shows that SpanDB simultaneously improves RocksDB's throughput by up to 8.8 \times and reduces its latency by 9.5–58.3%. Compared with KVell, a system designed for high-end SSDs, SpanDB achieves 96–140% of its throughput, with a 2.3–21.6 \times lower latency, at a cheaper storage configuration. Cheng Li 0001, Hao Chen 0080, Chaoyi Ruan, Xiaosong Ma, Yinlong Xu 0001 |
ACM Trans. Storage | 4 |
| 2020 | QarSUMO: A Parallel, Congestion-optimized Traffic SimulatorabstractTraffic simulators are important tools for tasks such as urban planning and transportation management. Microscopic simulators allow per-vehicle movement simulation, but require longer simulation time. The simulation overhead is exacerbated when there is traffic congestion and most vehicles move slowly. This in particular hurts the productivity of emerging urban computing studies based on reinforcement learning, where traffic simulations are heavily and repeatedly used for designing policies to optimize traffic related tasks. Hao Chen 0080, Stefano Giovanni Rizzo, Giovanna Vantini, Phillip Taylor, Xiaosong Ma, Sanjay Chawla |
SIGSPATIAL/GIS | 6 |
| 2020 | LiveGraph: A Transactional Graph Storage System with Purely Sequential Adjacency List ScansabstractThe specific characteristics of graph workloads make it hard to design a one-size-fits-all graph storage system. Systems that support transactional updates use data structures with poor data locality, which limits the efficiency of analytical workloads or even simple edge scans. Other systems run graph analytics workloads efficiently, but cannot properly support transactions. This paper presents LiveGraph, a graph storage system that outperforms both the best graph transactional systems and the best solutions for real-time graph analytics on fresh data. LiveGraph achieves this by ensuring that adjacency list scans, a key operation in graph workloads, are purely sequential: they never require random accesses even in presence of concurrent transactions. Such pure-sequential operations are enabled by combining a novel graph-aware data structure, the Transactional Edge Log (TEL), with a concurrency control mechanism that leverages TEL's data layout. Our evaluation shows that LiveGraph significantly outperforms state-of-the-art (graph) database solutions on both transactional and real-time analytical workloads. Xiaowei Zhu 0001, Marco Serafini, Xiaosong Ma, Ashraf Aboulnaga, Guanyu Feng |
Proc. VLDB Endow. | 3 |
| 2020 | Determining Data Distribution for Large Disk Enclosures with 3-D Data TemplatesabstractConventional RAID solutions with fixed layouts partition large disk enclosures so that each RAID group uses its own disks exclusively. This achieves good performance isolation across underlying disk groups, at the cost of disk under-utilization and slow RAID reconstruction from disk failures. We propose RAID+, a new RAID construction mechanism that spreads both normal I/O and reconstruction workloads to a larger disk pool in a balanced manner. Unlike systems conducting randomized placement, RAID+ employs deterministic addressing enabled by the mathematical properties of mutually orthogonal Latin squares, based on which it constructs 3-D data templates mapping a logical data volume to uniformly distributed disk blocks across all disks. While the total read/write volume remains unchanged, with or without disk failures, many more disk drives participate in data service and disk reconstruction. Our evaluation with a 60-drive disk enclosure using both synthetic and real-world workloads shows that RAID+ significantly speeds up data recovery while delivering better normal I/O performance and higher multi-tenant system throughput. Guangyan Zhang, Zhufan Wang, Xiaosong Ma, Zican Huang |
ACM Trans. Storage | 3 |
| 2019 | Automatic, Application-Aware I/O Forwarding Resource Allocation
Bin Yang 0043, Xiaosong Ma, Xiupeng Zhu, Xiyang Wang 0003, Nosayba El-Sayed, Jidong Zhai, Wei Xue 0003 |
FAST | 4 |
| 2019 | End-to-end I/O Monitoring on a Leading Supercomputer
Bin Yang 0043, Xiaosong Ma, Xiyang Wang 0003, Xiupeng Zhu, Nosayba El-Sayed, Haidong Lan, Jidong Zhai, Wei Xue 0003 |
NSDI | 3 |
| 2019 | Spread-n-share: improving application performance and cluster throughput with resource-aware job placementabstractTraditional batch job schedulers adopt the Compact-n-Exclusive (CE) strategy, packing processes of a parallel job into as few compute nodes as possible. While CE minimizes inter-node network communication, it often brings self-contention among tasks of a resource-intensive application. Recent studies have used virtual containers to balance CPU utilization and memory capacity across physical nodes, but the imbalance in cache and memory bandwidth usage is still under-investigated. Xiongchao Tang, Haojie Wang 0004, Xiaosong Ma, Nosayba El-Sayed, Jidong Zhai, Ashraf Aboulnaga |
SC | 3 |
| 2019 | KnightKing: a fast distributed graph random walk engineabstractRandom walk on graphs has recently gained immense popularity as a tool for graph data analytics and machine learning. Currently, random walk algorithms are developed as individual implementations and suffer significant performance and scalability problems, especially with the dynamic nature of sophisticated walk strategies. Kang Chen 0001, Xiaosong Ma, Yang Bai 0011, Yong Jiang 0001 |
SOSP | 4 |
| 2018 | RAID+: Deterministic and Balanced Data Distribution for Large Disk Enclosures
Guangyan Zhang, Zican Huang, Xiaosong Ma, Zhufan Wang |
FAST | 3 |
| 2018 | KPart: A Hybrid Cache Partitioning-Sharing Technique for Commodity MulticoresabstractCache partitioning is now available in commercial hardware. In theory, software can leverage cache partitioning to use the last-level cache better and improve performance. In practice, however, current systems implement way-partitioning, which offers a limited number of partitions and often hurts performance. These limitations squander the performance potential of smart cache management. We present KPart, a hybrid cache partitioning-sharing technique that sidesteps the limitations of way-partitioning and unlocks significant performance on current systems. KPart first groups applications into clusters, then partitions the cache among these clusters. To build clusters, KPart relies on a novel technique to estimate the performance loss an application suffers when sharing a partition. KPart automatically chooses the number of clusters, balancing the isolation benefits of way-partitioning with its potential performance impact. KPart uses detailed profiling information to make these decisions. This information can be gathered either offline, or online at low overhead using a novel profiling mechanism. We evaluate KPart in a real system and in simulation. KPart improves throughput by 24% on average (up to 79%) on an Intel Broadwell-D system, whereas prior per-application partitioning policies improve throughput by just 1.7% on average and hurt 30% of workloads. Simulation results show that KPart achieves most of the performance of more advanced partitioning techniques that are not yet available in hardware. Nosayba El-Sayed, Anurag Mukkara, Po-An Tsai, Harshad Kasture, Xiaosong Ma, Daniel Sánchez 0003 |
HPCA | 5 |
| 2018 | Exploiting Locality in Graph Analytics through Hardware-Accelerated Traversal SchedulingabstractGraph processing is increasingly bottlenecked by main memory accesses. On-chip caches are of little help because the irregular structure of graphs causes seemingly random memory references. However, most real-world graphs offer significant potential locality—it is just hard to predict ahead of time. In practice, graphs have well-connected regions where relatively few vertices share edges with many common neighbors. If these vertices were processed together, graph processing would enjoy significant data reuse. Hence, a graph's traversal schedule largely determines its locality. This paper explores online traversal scheduling strategies that exploit the community structure of real-world graphs to improve locality. Software graph processing frameworks use simple, locality-oblivious scheduling because, on general-purpose cores, the benefits of locality-aware scheduling are outweighed by its overheads. Software frameworks rely on offline preprocessing to improve locality. Unfortunately, preprocessing is so expensive that its costs often negate any benefits from improved locality. Recent graph processing accelerators have inherited this design. Our insight is that this misses an opportunity: Hardware acceleration allows for more sophisticated, online locality-aware scheduling than can be realized in software, letting systems significantly improve locality without any preprocessing. To exploit this insight, we present bounded depth-first scheduling (BDFS), a simple online locality-aware scheduling strategy. BDFS restricts each core to explore one small, connected region of the graph at a time, improving locality on graphs with good community structure. We then present HATS, a hardware-accelerated traversal scheduler that adds just 0.4% area and 0.2% power over general-purpose cores. We evaluate BDFS and HATS on several algorithms using large real-world graphs. On a simulated 16-core system, BDFS reduces main memory accesses by up to 2.4x and by 30% on average. However, BDFS is too expensive in software and degrades performance by 21% on average. HATS eliminates these overheads, allowing BDFS to improve performance by 83% on average (up to 3.1x) over a locality-oblivious software implementation and by 31% on average (up to 2.1x) over specialized prefetchers. Anurag Mukkara, Nathan Beckmann, Maleen Abeydeera, Xiaosong Ma, Daniel Sánchez 0003 |
MICRO | 4 |
| 2018 | ShenTu: processing multi-trillion edge graphs on millions of cores in seconds
Heng Lin, Xiaowei Zhu 0001, Bowen Yu 0003, Xiongchao Tang, Wei Xue 0003, Lufei Zhang, Torsten Hoefler, Xiaosong Ma, Xin Liu 0081, Jingfang Xu |
SC | 9 |
| 2018 | Spindle: Informed Memory Access Monitoring
Haojie Wang 0004, Jidong Zhai, Xiongchao Tang, Bowen Yu 0003, Xiaosong Ma |
USENIX ATC | 5 |
| 2017 | POSTER: Improving Datacenter Efficiency Through Partitioning-Aware SchedulingabstractDatacenter servers often colocate multiple applications to improve utilization and efficiency. However, colocated applications interfere in shared resources, e.g., the last-level cache (LLC) and DRAM bandwidth, causing performance inefficiencies. Prior work has proposed two disjoint approaches to address interference. First, techniques that partition shared resources like the LLC can provide isolation and trade performance among colocated applications within a single node. But partitioning techniques are limited by the fixed resource demands of the applications running on the node. Second, interference-aware schedulers try to find resource-compatible applications and schedule them across nodes to improve performance. But prior schedulers are hampered by the lack of partitioning hardware in conventional multicores, and are forced to take conservative colocation decisions, leaving significant performance on the table. We show that memory-system partitioning and scheduling are complementary, and performing them in a coordinated fashion yields significant benefits. We present Shepherd, a joint scheduler and resource partitioner that seeks to maximize cluster-wide throughput. Shepherd uses detailed application profiling data to partition the shared LLC and to estimate the impact of DRAM bandwidth contention among colocated applications. Shepherd's scheduler leverages this information to colocate applications with complementary resource requirements, improving resource utilization and cluster throughput. We evaluate Shepherd in simulation and on a real cluster with hardware support for cache partitioning. When managing mixes of server and scientific applications, Shepherd improves cluster throughput over an unpartitioned system by 38% on average. Harshad Kasture, Nosayba El-Sayed, Nathan Beckmann, Xiaosong Ma, Daniel Sánchez 0003 |
PACT | 5 |
| 2017 | Understanding object-level memory access patterns across the spectrumabstractMemory accesses limit the performance and scalability of countless applications. Many design and optimization efforts will benefit from an in-depth understanding of memory access behavior, which is not offered by extant access tracing and profiling methods. Chao Wang 0056, Nosayba El-Sayed, Xiaosong Ma, Youngjae Kim 0001, Sudharshan S. Vazhkudai, Wei Xue 0003, Daniel Sánchez 0003 |
SC | 4 |
| 2016 | Gemini: A Computation-Centric Distributed Graph Processing System
Xiaowei Zhu 0001, Xiaosong Ma |
OSDI | 4 |
| 2016 | Server-side log data analytics for I/O workload characterization and coordination on large shared storage systemsabstractInter-application I/O contention and performance interference have been recognized as severe problems. In this work, we demonstrate, through measurement from Titan (world's No. 3 supercomputer), that high I/O variance co-exists with the fact that individual storage units remain under-utilized for the majority of the time. This motivates us to propose AID, a system that performs automatic application I/O characterization and I/O-aware job scheduling. AID analyzes existing I/O traffic and batch job history logs, without any prior knowledge on applications or user/developer involvement. It identifies the small set of I/O-intensive candidates among all applications running on a supercomputer and subsequently mines their I/O patterns, using more detailed per-I/O-node traffic logs. Based on such auto-extracted information, AID provides online I/O-aware scheduling recommendations to steer I/O-intensive applications away from heavy ongoing I/O activities. We evaluate AID on Titan, using both real applications (with extracted I/O patterns validated by contacting users) and our own pseudo-applications. Our results confirm that AID is able to (1) identify I/O-intensive applications and their detailed I/O characteristics, and (2) significantly reduce these applications' I/O performance degradation/variance by jointly evaluating outstanding applications' I/O pattern and real-time system l/O load. Yang Liu 0129, Raghul Gunasekaran, Xiaosong Ma, Sudharshan S. Vazhkudai |
SC | 3 |
| 2016 | S-RAC: SSD Friendly Caching for Data Center WorkloadsabstractCurrent data-center applications tend to process increasingly large volume of data sets. The caching effect of page cache is reduced by its limited capacity. Emerging flash-based solid state drives (SSD) have latency and price advantages compared to hard disk and DRAM. Thus, SSD-based caching is widely deployed in data centers. However, SSD caching faces two challenges. First, SSD has limited write endurance, which requires cache manager to reduce write amount to SSD. Second, data-center workloads exhibit a diverse I/O access patterns, which requires one to figure out SSD caching friendly access patterns. This paper first classifies 6 I/O access patterns among 32 data-center workloads using a cost-benefit analysis. We derive implications for designing SSD cache from analyzing the access patterns. We then propose an SSD cache manager S-RAC with re-adding blocks and ghost cache adaptation to retain SSD friendly blocks in SSD. The experimental evaluation shows the efficiency of S-RAC in reducing SSD write amount while improving/maintaining cache hit ratio. Yuanjiang Ni, Ji Jiang, Dejun Jiang 0001, Xiaosong Ma, Jin Xiong, Yuangang Wang |
SYSTOR | 4 |
| 2016 | MPI-ACC: Accelerator-Aware MPI for Scientific ApplicationsabstractData movement in high-performance computing systems accelerated by graphics processing units (GPUs) remains a challenging problem. Data communication in popular parallel programming models, such as the Message Passing Interface (MPI), is currently limited to the data stored in the CPU memory space. Auxiliary memory systems, such as GPU memory, are not integrated into such data movement standards, thus providing applications with no direct mechanism to perform end-to-end data movement. We introduce MPI-ACC, an integrated and extensible framework that allows end-to-end data movement in accelerator-based systems. MPI-ACC provides productivity and performance benefits by integrating support for auxiliary memory spaces into MPI. MPI-ACC supports data transfer among CUDA, OpenCL and CPU memory spaces and is extensible to other offload models as well. MPI-ACC's runtime system enables several key optimizations, including pipelining of data transfers, scalable memory management techniques, and balancing of communication based on accelerator and node architecture. MPI-ACC is designed to work concurrently with other GPU workloads with minimum contention. We describe how MPI-ACC can be used to design new communication-computation patterns in scientific applications from domains such as epidemiology simulation and seismology modeling, and we discuss the lessons learned. We present experimental results on a state-of-the-art cluster with hundreds of GPUs; and we compare the performance and productivity of MPI-ACC with MVAPICH, a popular CUDA-aware MPI solution. MPI-ACC encourages programmers to explore novel application-specific optimizations for improved overall cluster utilization. Ashwin M. Aji, Lokendra S. Panwar, Karthik Murthy, Milind Chabbi, Pavan Balaji, Keith R. Bisset, James Dinan, Wu-chun Feng, John M. Mellor-Crummey, Xiaosong Ma, Rajeev Thakur |
IEEE Trans. Parallel Distributed Syst. | 11 |
| 2016 | Building Semi-Elastic Virtual Clusters for Cost-Effective HPC Cloud Resource ProvisioningabstractRecent studies have found cloud environments increasingly appealing for executing HPC applications, including tightly coupled parallel simulations. At the same time, while public clouds offer elastic, on-demand resource provisioning and pay-as-you-go pricing, individual users setting up their on-demand virtual clusters may not be able to take full advantage of common cost-saving opportunities, such as reserved instances. In this paper, we propose a Semi-Elastic Cluster (SEC) computing model for organizations to reserve and dynamically resize a virtual cloud-based cluster. We present a set of integrated batch scheduling plus resource scaling strategies uniquely enabled by SEC, as well as an online reserved instance provisioning algorithm based on job history. Our trace-driven simulation results show that such a model has a 61.0 percent cost saving than individual users acquiring and managing cloud resources without causing longer average job wait time. Moreover, to exploit the advantages of different public clouds, we also extend SEC to a multi-cloud environment, where SEC can get a lower cost than on any single cloud. We design and implement a prototype system of the SEC model and evaluate it in terms of management overhead and average job wait time. Experimental results show that the management overhead is negligible with respect to the job wait time. Shuangcheng Niu, Jidong Zhai, Xiaosong Ma, Xiongchao Tang |
IEEE Trans. Parallel Distributed Syst. | 3 |
| 2015 | Cost-Effective Resource Configuration for Cloud Video Streaming ServicesabstractVideo streaming services are migrating to cloud environments for the economic expense with good scalability. However, cloud providers offer flexible resource configurations, e.g., on-demand, reserved and spot instances, with significant different pricing policies, of which one single configuration is suboptimal for cloud video streaming services. In this paper, we propose hybrid configuration schemes of cloud video streaming services to reduce the cost. To achieve this goal, we first introduce a lightweight prediction algorithm to predict the future video traffic. With the predicted video traffic, we then give the Hybrid-R hybrid configuration scheme by configuring both on-demand and reserved instances, and the Hybrid-RS hybrid configuration scheme by further configuring spot instances. Our evaluations using traces from real video service providers show that our configuration schemes can reduce cost by at least 20% compared to the unoptimized ones with negligible overhead. Yunyun Jiang, Xiaosong Ma |
ICPADS | 2 |
| 2015 | Combining phase identification and statistic modeling for automated parallel benchmark generationabstractParallel application benchmarks are indispensable for evaluating/optimizing HPC software and hardware. However, it is very challenging and costly to obtain high-fidelity benchmarks reflecting the scale and complexity of state-of-the-art parallel applications. Hand-extracted synthetic benchmarks are time- and labor-intensive to create. Real applications themselves, while offering most accurate performance evaluation, are expensive to compile, port, recon- figure, and often plainly inaccessible due to security or ownership concerns. This work contributes APPRIME, a novel tool for trace-based automatic parallel benchmark generation. Taking as input standard communication-I/O traces of an application’s execution, it couples accurate automatic phase identification with statistical regeneration of event parameters to create compact, portable, and to some degree reconfigurable parallel application benchmarks. Experiments with four NAS Parallel Benchmarks (NPB) and three real scientific simulation codes confirm the fidelity of APPRIME benchmarks. They retain the original applications’ performance characteristics, in particular the relative performance across platforms. Xiaosong Ma, Qing Liu 0002, Jeremy Logan, Norbert Podhorszki, Jong Choi 0001, Scott Klasky |
PPoPP | 3 |
| 2015 | Combining Phase Identification and Statistic Modeling for Automated Parallel Benchmark GenerationabstractParallel application benchmarks are indispensable for evaluating/optimizing HPC software and hardware. However, it is very challenging and costly to obtain high-fidelity benchmarks reflecting the scale and complexity of state-of-the-art parallel applications. Hand-extracted synthetic benchmarks are time- and labor-intensive to create. Real applications themselves, while offering most accurate performance evaluation, are expensive to compile, port, reconfigure, and often plainly inaccessible due to security or ownership concerns. This work contributes APPrime, a novel tool for trace-based automatic parallel benchmark generation. Taking as input standard communication-I/O traces of an application's execution, it couples accurate automatic phase identification with statistical regeneration of event parameters to create compact, portable, and to some degree reconfigurable parallel application benchmarks. Experiments with four NAS Parallel Benchmarks (NPB) and three real scientific simulation codes confirm the fidelity of APPrime benchmarks. They retain the original applications' performance characteristics, in particular their relative performance across platforms. Also, the result benchmarks, already released online, are much more compact and easy-to-port compared to the original applications. Xiaosong Ma, Qing Liu 0002, Jeremy Logan, Norbert Podhorszki, Jong Choi 0001, Scott Klasky |
SIGMETRICS | 2 |
| 2015 | Automatic Cloud I/O Configurator for I/O Intensive Parallel ApplicationsabstractAs the cloud platform becomes a promising alternative to traditional HPC (high performance computing) centers or in-house clusters, the I/O bottleneck problem is highlighted in this new environment, typically with top-of-the-line compute instances but sub-par communication and I/O facilities. It has been observed that changing the cloud I/O system configurations, such as choices of file systems, number of I/O servers and their placement strategies, etc., will lead to a considerable variation in the performance and cost efficiency of I/O intensive parallel applications. However, storage system configuration is tedious and error-prone to do manually, even for expert users, leading to solutions that are grossly over-provisioned (low cost inefficiency), substantially under-performing (poor performance) or, in the worst case, both. This paper proposes ACIC, a system which automatically searches for optimized I/O system configurations from many candidates for each individual application running on a given cloud platform. ACIC takes advantage of machine learning models to perform performance/cost predictions. To tackle the high-dimensional parameter exploration space, we enable affordable, reusable, and incremental training on cloud platforms, guided by the Plackett and Burman Matrices for experiment design. Our evaluation results with four representative parallel applications indicate that ACIC consistently identifies optimal or near-optimal configurations among a large group of candidate settings. The top ACIC-recommended configuration is capable of improving the applications' performance by a factor of up to 10.5 (3.1 on average), and cost saving of up to 89 percent (51 percent on average), compared with a commonly used baseline I/O configuration. In addition, we carried out a small-scale user study for one of the test applications, which found that ACIC consistently beat the user and even the application's developer, often by a significant margin, in selecting optimized configurations. Jidong Zhai, Xiaosong Ma |
IEEE Trans. Parallel Distributed Syst. | 4 |
| 2014 | Automatic identification of application I/O signatures from noisy server-side traces
Yang Liu 0129, Raghul Gunasekaran, Xiaosong Ma, Sudharshan S. Vazhkudai |
FAST | 3 |
| 2014 | CYPRESS: Combining Static and Dynamic Analysis for Top-Down Communication Trace CompressionabstractCommunication traces are increasingly important, both for parallel applications' performance analysis/optimization, and for designing next-generation HPC systems. Meanwhile, the problem size and the execution scale on supercomputers keep growing, producing prohibitive volume of communication traces. To reduce the size of communication traces, existing dynamic compression methods introduce large compression overhead with the job scale. We propose a hybrid static-dynamic method that leverages information acquired from static analysis to facilitate more effective and efficient dynamic trace compression. Our proposed scheme, Cypress, extracts a program communication structure tree at compile time using inter-procedural analysis. This tree naturally contains crucial iterative computing features such as the loop structure, allowing subsequent runtime compression to "fill in", in a "top-down" manner, event details into the known communication template. Results show that Cypress reduces intra-process and inter-process compression overhead up to 5× and 9× respectively over state-of-the-art dynamic methods, while only introducing very low compiling overhead. Jidong Zhai, Jianfei Hu, Xiongchao Tang, Xiaosong Ma |
SC | 4 |
| 2014 | vCacheShare: Automated Server Flash Cache Space Management in a Virtualization Environment
Xiaosong Ma, Sandeep Uttamchandani, Deng Liu |
USENIX ATC | 3 |
| 2013 | RSVM: A Region-based Software Virtual Memory for GPUabstractWhile Graphics Processing Units (GPU) have gained much success in general purpose computing in recent years, their programming is still difficult, due to, particularly, explicitly managed GPU memory and manual CPU-GPU data transfer. Despite recent calls for managing GPU resources as first-class citizens in the operating system, a mature GPU memory management mechanism is still missing, which leads to reinventing the wheels in various GPU system software. Meanwhile, due to ever enlarging problem sizes, we urgently need a system-level mechanism for unified CPU-GPU memory management. In this work, we present the design of Region-based Software Virtual Memory (RSVM), a software virtual memory running on both CPU and GPU in a distributed and cooperative way. In addition to automatic GPU memory management and GPU-CPU data transfer, RSVM offers two novel features: 1) GPU kernel-issued on-demand data fetching from the host into the GPU memory, and 2) intra-kernel transparent GPU memory swapping into the main memory. Our study reveals important insights on the challenges and opportunities of building unified virtual memory systems for heterogeneous computing. Experimental results on real GPU benchmarks demonstrate that, though it incurs a small overhead, RSVM can transparently scale GPU kernels to large problem sizes exceeding the device memory size limit; developers write the same code for different problem sizes, but still can optimize on data layout definition accordingly. Our evaluation also identifies missing GPU architecture features for better system software efficiency. Heshan Lin, Xiaosong Ma |
PACT | 3 |
| 2013 | PARLO: PArallel Run-Time Layout Optimization for Scientific Data Explorations with Heterogeneous Access PatternsabstractThe size and scope of cutting-edge scientific simulations are growing much faster than the I/O and storage capabilities of their run-time environments. The growing gap is exacerbated by exploratory, data-intensive analytics, such as querying simulation data with multivariate, spatio-temporal constraints, which induces heterogeneous access patterns that stress the performance of the underlying storage system. Previous work addresses data layout and indexing techniques to improve query performance for a single access pattern, which is not sufficient for complex analytics jobs. We present PARLO a parallel run-time layout optimization framework, to achieve multi-level data layout optimization for scientific applications at run-time before data is written to storage. The layout schemes optimize for heterogeneous access patterns with user-specified priorities. PARLO is integrated with ADIOS, a high-performance parallel I/O middleware for large-scale HPC applications, to achieve user-transparent, light-weight layout optimization for scientific datasets. It offers simple XML-based configuration for users to achieve flexible layout optimization without the need to modify or recompile application codes. Experiments show that PARLO improves performance by 2 to 26 times for queries with heterogeneous access patterns compared to state-of-the-art scientific database management systems. Compared to traditional post-processing approaches, its underlying run-time layout optimization achieves a 56% savings in processing time and a reduction in storage overhead of up to 50%. PARLO also exhibits a low run-time resource requirement, while also limiting the performance impact on running applications to a reasonable level. Zhenhuan Gong, David A. Boyuka II, Xiaocheng Zou, Qing Liu 0002, Norbert Podhorszki, Scott Klasky, Xiaosong Ma, Nagiza F. Samatova |
CCGRID | 7 |
| 2013 | Active flash: towards energy-efficient, in-situ data analytics on extreme-scale machines
Devesh Tiwari, Simona Boboila, Sudharshan S. Vazhkudai, Youngjae Kim 0001, Xiaosong Ma, Peter Desnoyers, Yan Solihin |
FAST | 5 |
| 2013 | On the efficacy of GPU-integrated MPI for scientific applications
Ashwin M. Aji, Lokendra S. Panwar, Milind Chabbi, Karthik Murthy, Pavan Balaji, Keith R. Bisset, James Dinan, Wu-chun Feng, John M. Mellor-Crummey, Xiaosong Ma, Rajeev Thakur |
HPDC | 11 |
| 2013 | Building and scaling virtual clusters with residual resources from interactive clouds
R. Benjamin Clay, Zhiming Shen, Xiaosong Ma |
HPDC | 3 |
| 2013 | ACIC: automatic cloud I/O configurator for parallel applications
Jidong Zhai, Yan Zhai, Xiaosong Ma |
HPDC | 6 |
| 2013 | Accelerating Batch Analytics with Residual Resources from Interactive CloudsabstractThe popularity of cloud-based interactive computing services (e.g., virtual desktops) brings new management challenges. Each interactive user leaves abundant but fluctuating residual resources while being intolerant to latency, precluding the use of aggressive VM consolidation. In this paper, we present the Resource Harvester for Interactive Clouds (RHIC), an autonomous management framework that harnesses dynamic residual resources aggressively without slowing the harvested interactive services. RHIC builds ad-hoc clusters for running throughput-oriented "background" workloads using a hybrid of residual and dedicated resources. These hybrid clusters offer significant gains over normal dedicated clusters: 20-40% cost and 20-29% energy savings in our test bed. For a given background job, RHIC intelligently discovers and maintains the ideal cluster size and composition, to meet user-specified goals such as cost/energy minimization or deadlines. RHIC employs black-box workload performance modeling, requiring only system-level metrics and incorporating techniques to improve modeling accuracy with bursty and heterogeneous residual resources. We demonstrate the effectiveness and adaptivity of our RHIC prototype with two parallel data analytics frameworks, Hadoop and HBase. Our results show that RHIC finds near-ideal cluster sizes and compositions across a wide range of workload/goal combinations. R. Benjamin Clay, Zhiming Shen, Xiaosong Ma |
MASCOTS | 3 |
| 2013 | ACIC: automatic cloud I/O configurator for HPC applicationsabstractThe cloud has become a promising alternative to traditional HPC centers or in-house clusters. This new environment highlights the I/O bottleneck problem, typically with top-of-the-line compute instances but sub-par communication and I/O facilities. It has been observed that changing cloud I/O system configurations leads to significant variation in the performance and cost efficiency of I/O intensive HPC applications. However, storage system configuration is tedious and error-prone to do manually, even for experts. Jidong Zhai, Yan Zhai, Xiaosong Ma |
SC | 6 |
| 2013 | Cost-effective cloud HPC resource provisioning by building semi-elastic virtual clustersabstractRecent studies have found cloud environments increasingly appealing for executing HPC applications, including tightly coupled parallel simulations. While public clouds offer elastic, on-demand resource provisioning and pay-as-you-go pricing, individual users setting up their on-demand virtual clusters may not be able to take full advantage of common cost-saving opportunities, such as reserved instances. Shuangcheng Niu, Jidong Zhai, Xiaosong Ma, Xiongchao Tang |
SC | 3 |
| 2012 | NVMalloc: Exposing an Aggregate SSD Store as a Memory Partition in Extreme-Scale MachinesabstractDRAM is a precious resource in extreme-scale machines and is increasingly becoming scarce, mainly due to the growing number of cores per node. On future multi-petaflop and exaflop machines, the memory pressure is likely to be so severe that we need to rethink our memory usage models. Fortunately, the advent of non-volatile memory (NVM) offers a unique opportunity in this space. Current NVM offerings possess several desirable properties, such as low cost and power efficiency, but suffer from high latency and lifetime issues. We need rich techniques to be able to use them alongside DRAM. In this paper, we propose a novel approach for exploiting NVM as a secondary memory partition so that applications can explicitly allocate and manipulate memory regions therein. More specifically, we propose an NVMalloc library with a suite of services that enables applications to access a distributed NVM storage system. We have devised ways within NVMalloc so that the storage system, built from compute node-local NVM devices, can be accessed in a byte-addressable fashion using the memory mapped I/O interface. Our approach has the potential to re-energize out-of-core computations on large-scale machines by having applications allocate certain variables through NVMalloc, thereby increasing the overall memory capacity available. Our evaluation on a 128-core cluster shows that NVMalloc enables applications to compute problem sizes larger than the physical memory in a cost-effective manner. It can bring more performance/efficiency gain with increased computation time between NVM memory accesses or increased data access locality. In addition, our results suggest that while NVMalloc enables transparent access to NVM-resident variables, the explicit control it provides is crucial to optimize application performance. Chao Wang 0056, Sudharshan S. Vazhkudai, Xiaosong Ma, Youngjae Kim 0001, Christian Engelmann |
IPDPS | 3 |
| 2012 | Employing Checkpoint to Improve Job Scheduling in Large-Scale Systems
Shuangcheng Niu, Jidong Zhai, Xiaosong Ma, Yan Zhai |
JSSPP | 3 |
| 2011 | Probabilistic Communication and I/O Tracing with Deterministic Replay at ScaleabstractWith today's petascale supercomputers, applications often exhibit low efficiency, such as poor communication and I/O performance, that can be diagnosed by analysis tools. However, these tools either produce extremely large trace files that complicate performance analysis, or sacrifice accuracy to collect high-level statistical information using crude averaging. This work contributes Scala-H-Trace, which features more aggressive trace compression than any previous approach, particularly for applications that do not show strict regularity in SPMD behavior. Scala-H-Trace uses histograms expressing the probabilistic distribution of arbitrary communication and I/O parameters to capture variations. Yet, where other tools fail to scale, Scala-H-Trace guarantees trace files of near constant size, even for variable communication and I/O patterns, producing trace files orders of magnitudes smaller than using prior approaches. We demonstrate the ability to collect traces of applications running on thousands of processors with the potential to scale well beyond this level. We further present the first approach to deterministically replay such probabilistic traces (a) without deadlocks and (b) in a manner closely resembling the original applications. Our results show either near constant sized traces or only sub-linear increases in trace file sizes irrespective of the number of nodes utilized. Even with the aggressively compressed histogram-based traces, our replay times are within 12% to 15% of the runtime of original codes. Such concise traces resembling the behavior of production-style codes closely and our approach of deterministic replay of probabilistic traces are without precedence. Xing Wu 0004, Karthik Vijayakumar, Frank Mueller 0001, Xiaosong Ma, Philip C. Roth |
ICPP | 4 |
| 2011 | Using Shared Memory to Accelerate MapReduce on Graphics Processing UnitsabstractModern General Purpose Graphics Processing Units (GPGPUs) provide high degrees of parallelism in computation and memory access, making them suitable for data parallel applications such as those using the elastic MapReduce model. Yet designing a MapReduce framework for GPUs faces significant challenges brought by their multi-level memory hierarchy. Due to the absence of atomic operations in the earlier generations of GPUs, existing GPU MapReduce frameworks have problems in handling input/output data with varied or unpredictable sizes. Also, existing frameworks utilize mostly a single level of memory, i.e., the relatively spacious yet slow global memory. In this work, we attempt to explore the potential benefit of enabling a GPU MapReduce framework to use multiple levels of the GPU memory hierarchy. We propose a novel GPU data staging scheme for MapReduce workloads, tailored toward the GPU memory hierarchy. Centering around the efficient utilization of the fast but very small shared memory, we designed and implemented a GPU MapReduce framework, whose key techniques include (1) shared memory staging area management, (2) thread-role partitioning, and (3) intra-block thread synchronization. We carried out evaluation with five popular MapReduce workloads and studied their performance under different GPU memory usage choices. Our results reveal that exploiting GPU shared memory is highly promising for the Map phase (with an average 2.85x speedup over using global memory only), while in the Reduce phase the benefit of using shared memory is much less pronounced, due to the high input-to-output ratio. In addition, when compared to Mars, an existing GPU MapReduce framework, our system is shown to bring a significant speedup in Map/Reduce phases. Xiaosong Ma |
IPDPS | 2 |
| 2011 | EMFS: Email-based Personal Cloud StorageabstractThough a variety of cloud storage services have been offered recently, they have not yet provided users with transparent and cost-effective personal data storage. Services like Google Docs offer easy file access and sharing, but tie storage with internal data formats and specific applications. Meanwhile, services like Drop box offer general-purpose storage. Yet they have not been widely utilized, partly due to their fee-charging nature and long-term service availability concerns. Web-based email services, on the other hand, have been offering growing email storage capacity, reliable service, and powerful search capability, making them appealing as storage resources. In this paper, we examine the efficacy of leveraging web-based email services to build a personal storage cloud. We present EMFS, which aggregates back-end storage by establishing a RAID-like system on top of virtual email disks formed by email accounts. In particular, by replicating data across accounts from different service providers, highly available storage services can be constructed based on already reliable, cloud-based email storage. This paper discusses the design and implementation of EMFS, focusing on unique challenges and opportunities associated with utilizing email services for file transfer and storage, such as email based data organization, metadata format and management, and handling provider-imposed anti-spam usage restrictions. We evaluated EMFS extensively with multiple benchmarks, and compared its performance with NFS, AFS, and a non-free cloud storage service built upon Amazon S3. Our results indicate that while EMFS cannot match the performance of highly optimized distributed file systems with dedicated servers, it performs quite closely to the commercial cloud storage solution. Jagan Srinivasan, Xiaosong Ma, Ting Yu 0001 |
NAS | 3 |
| 2011 | Transparent runtime parallelization of the R scripting language
Jiangtian Li, Xiaosong Ma, Srikanth B. Yoginath, Guruprasad Kora, Nagiza F. Samatova |
J. Parallel Distributed Comput. | 2 |
| 2011 | Coordinating Computation and I/O in Massively Parallel Sequence SearchabstractWith the explosive growth of genomic information, the searching of sequence databases has emerged as one of the most computation and data-intensive scientific applications. Our previous studies suggested that parallel genomic sequence-search possesses highly irregular computation and I/O patterns. Effectively addressing these runtime irregularities is thus the key to designing scalable sequence-search tools on massively parallel computers. While the computation scheduling for irregular scientific applications and the optimization of noncontiguous file accesses have been well-studied independently, little attention has been paid to the interplay between the two. In this paper, we systematically investigate the computation and I/O scheduling for data-intensive, irregular scientific applications within the context of genomic sequence search. Our study reveals that the lack of coordination between computation scheduling and I/O optimization could result in severe performance issues. We then propose an integrated scheduling approach that effectively improves sequence-search throughput by gracefully coordinating the dynamic load balancing of computation and high-performance noncontiguous I/O. Heshan Lin, Xiaosong Ma, Wu-chun Feng, Nagiza F. Samatova |
IEEE Trans. Parallel Distributed Syst. | 2 |
| 2010 | MOON: MapReduce On Opportunistic eNvironmentsabstractMapReduce offers an ease-of-use programming paradigm for processing large data sets, making it an attractive model for distributed volunteer computing systems. However, unlike on dedicated resources, where MapReduce has mostly been deployed, such volunteer computing systems have significantly higher rates of node unavailability. Furthermore, nodes are not fully controlled by the MapReduce framework. Consequently, we found the data and task replication scheme adopted by existing MapReduce implementations woefully inadequate for resources with high unavailability. Heshan Lin, Xiaosong Ma, Jeremy S. Archuleta, Wu-chun Feng, Mark K. Gardner, Zhe Zhang 0005 |
HPDC | 2 |
| 2010 | Functional Partitioning to Optimize End-to-End Performance on Many-core ArchitecturesabstractScaling computations on emerging massive-core supercomputers is a daunting task, which coupled with the significantly lagging system I/O capabilities exacerbates applications' end-to-end performance. The I/O bottleneck often negates potential performance benefits of assigning additional compute cores to an application. In this paper, we address this issue via a novel functional partitioning (FP) runtime environment that allocates cores to specific application tasks - checkpointing, de-duplication, and scientific data format transformation - so that the deluge of cores can be brought to bear on the entire gamut of application activities. The focus is on utilizing the extra cores to support HPC application I/O activities and also leverage solid-state disks in this context. For example, our evaluation shows that dedicating 1 core on an oct-core machine for checkpointing and its assist tasks using FP can improve overall execution time of a FLASH benchmark on 80 and 160 cores by 43.95% and 41.34%, respectively. Sudharshan S. Vazhkudai, Ali Raza Butt, Xiaosong Ma, Youngjae Kim 0001, Christian Engelmann, Galen M. Shipman |
SC | 5 |
| 2009 | Memory resource allocation for file system prefetching: from a supply chain management perspectiveabstractAs an important technique to hide disk I/O latency, prefetching has been widely studied, and dynamic adaptive prefetching techniques have been deployed in diverse storage environments. However, two issues are not well addressed by previous research: (1) how to handle the prefetching resource allocation between concurrent sequential access streams with different request rates, and (2) how to coordinate prefetching at multiple levels in the data access path. Zhe Zhang 0005, Xiaosong Ma, Yuanyuan Zhou 0001 |
EuroSys | 3 |
| 2009 | Machine learning based online performance prediction for runtime parallelization and task schedulingabstractWith the emerging many-core paradigm, parallel programming must extend beyond its traditional realm of scientific applications. Converting existing sequential applications as well as developing next-generation software requires assistance from hardware, compilers and runtime systems to exploit parallelism transparently within applications. These systems must decompose applications into tasks that can be executed in parallel and then schedule those tasks to minimize load imbalance. However, many systems lack a priori knowledge about the execution time of all tasks to perform effective load balancing with low scheduling overhead. In this paper, we approach this fundamental problem using machine learning techniques first to generate performance models for all tasks and then applying those models to perform automatic performance prediction across program executions. We also extend an existing scheduling algorithm to use generated task cost estimates for online task partitioning and scheduling. We implement the above techniques in the pR framework, which transparently parallelizes scripts in the popular R language, and evaluate their performance and overhead with both a real-world application and a large number of synthetic representative test scripts. Our experimental results show that our proposed approach significantly improves task partitioning and scheduling, with maximum improvements of 21.8%, 40.3% and 22.1% and average improvements of 15.9%, 16.9% and 4.2% for LMM (a real R application) and synthetic test cases with independent and dependent tasks, respectively. Jiangtian Li, Xiaosong Ma, Martin Schulz 0001, Bronis R. de Supinski, Sally A. McKee |
ISPASS | 2 |
| 2009 | SigLM: Signature-driven load management for cloud computing infrastructuresabstractCloud computing has emerged as a promising platform that grants users with direct yet shared access to computing resources and services without worrying about the internal complex infrastructure. Unlike traditional batch service model, cloud service model adopts a pay-as-you-go form, which demands explicit and precise resource control. In this paper, we present SigLM, a novel Signature-driven Load Management system to achieve quality-aware service delivery in shared cloud computing infrastructures. SigLM dynamically captures fine-grained signatures of different application tasks and cloud nodes using time series patterns, and performs precise resource metering and allocation based on the extracted signatures. SigLM employs dynamic time warping algorithm and multi-dimensional time series indexing to achieve efficient signature pattern matching. Our experiments using real load traces collected on the PlanetLab show that SigLM can improve resource provisioning performance by 30–80% compared to existing approaches. SigLM is scalable and efficient, which imposes less than 1% overhead to the system and can perform signature matching within tens of milliseconds. Zhenhuan Gong, Prakash Ramaswamy, Xiaohui Gu, Xiaosong Ma |
IWQoS | 4 |
| 2009 | Energy and performance impact of aggressive volunteer computing with multi-core computersabstractThe rapid advances in multi-core architecture and the predicted emergence of 100-core personal computers bring new appeal to volunteer computing. The availability of massive compute power under-utilized by personal computing tasks is a blessing to volunteer computing customers. Meanwhile the reduced performance impact of running a foreign workload, thanks to the increased hardware parallelism, makes volunteering resources more acceptable to PC owners. In addition, we suspect that with aggressive volunteer computing, which assigns foreign tasks to active computers (as opposed to idle ones in the common practice), we can obtain significant energy savings. In this paper, we assess the efficacy of such aggressive volunteer computing model by evaluating the energy saving and performance impact of co-executing resource-intensive foreign workloads with native personal computing tasks. Our results from executing 30 native-foreign workload combinations suggest that aggressive volunteer computing can achieve an average energy saving of around 52% compared to running the foreign workloads on high-end cluster nodes, and around 33% compared to using the traditional, more conservative volunteer computing model. We have also observed highly varied performance interference behavior between the workloads, and evaluated the effectiveness of foreign workload intensity throttling. Jiangtian Li, Amey Deshpande, Jagan Srinivasan, Xiaosong Ma |
MASCOTS | 4 |
| 2009 | Improving Data Availability for Better Access Performance: A Study on Caching Scientific Data on Distributed Desktop Workstations
Xiaosong Ma, Sudharshan S. Vazhkudai, Zhe Zhang 0005 |
J. Grid Comput. | 1 |
| 2008 | PFC: Transparent Optimization of Existing Prefetching Strategies for Multi-Level Storage SystemsabstractThe multi-level storage architecture has been widely adopted in servers and data centers. However, while prefetching has been shown as a crucial technique to exploit the sequentiality in accesses common for such systems and hide the increasing relative cost of disk I/O, existing multi-level storage studies have focused mostly on cache replacement strategies. In this paper, we show that prefetching algorithms designed for single-level systems may have their limitations magnified when applied to multi-level systems. Overly conservative prefetching will not be able to effectively use the lower-level cache space, while overly aggressive prefetching will be compounded across levels and generate large amounts of wasted prefetch. We take an innovative approach to this problem: rather than designing a new, multi-level prefetching algorithm, we developed prefetching-coordinator (PFC), a hierarchy-aware optimization applicable to any existing prefetching algorithms. PFC does not require any application hints, a priori knowledge on the application access pattern or the native prefetching algorithm, or modification to the I/O interface. Instead, it monitors the upper-level access patterns as well as the lower-level cache status, and dynamically adjusts the aggressiveness of the lower-level prefetching activities. We evaluated PFC with extensive simulation study using a verified multi-level storage simulator, an accurate disk simulator, and access traces with different access patterns. Our results indicate that PFC dynamically controls lower-level prefetching in reaction to multiple system and workload parameters, improving the overall system performance in all 96 test cases. Working with four well-known existing prefetching algorithms adopted in real systems, PFC obtains an improvement of up to 35% to the average request response time, with an average improvement of 14.6% over all cases. Zhe Zhang 0005, Kyuhyung Lee, Xiaosong Ma, Yuanyuan Zhou 0001 |
ICDCS | 3 |
| 2008 | On-the-Fly Recovery of Job Input Data in SupercomputersabstractStorage system failure is a serious concern as we approach Petascale computing. Even at today's sub-Petascale levels, I/O failure is the leading cause of downtimes and job failures. We contribute a novel, on-the-fly recovery framework for job input data into supercomputer parallel file systems. The framework exploits key traits ofthe HPC I/O workload to reconstruct lost input data during job execution from remote, immutable copies. Each reconstructed data stripe is made immediately accessible in the client request order due to the delayed metadata update and fine-granular locking while unrelated access to the same file remains unaffected. We have implemented the recovery component within the Lustre parallel file system, thus building a novel application-transparent online recovery solution. Our solution is integrated into Lustre's two-level locking scheme using a two-phase blocking protocol. Combining parametric and simulation studies, our experiments demonstrate a significant improvement in HPC center service ability and user job turnaround time. Chao Wang 0056, Zhe Zhang 0005, Sudharshan S. Vazhkudai, Xiaosong Ma, Frank Mueller 0001 |
ICPP | 4 |
| 2008 | Semantics-based distributed I/O for mpiBLASTabstractBLAST is a widely used software toolkit for genomic sequence search. mpiBLAST is a freely available, open-source parallelization of BLAST that uses database segmentation to allow different worker processes to search (in parallel) unique segments of the database. After searching, the workers write their output to a filesystem. While mpiBLAST has been shown to achieve high performance in clusters with fast local filesystems, its I/O processing remains a concern for scalability, especially in systems having limited I/O capabilities such as distributed filesystems spread across a wide-area network. Thus, we present ParaMEDIC---a novel environment that uses application-specific semantic information to compress I/O data and improve performance in distributed environments. Specifically, for mpiBLAST, ParaMEDIC partitions worker processes into compute and I/O workers. Compute workers, instead of directly writing the output to the filesystem, the workers process the output using semantic knowledge about the application to generate metadata and write the metadata to the filesystem. I/O workers, which physically reside closer to the actual storage, then process this metadata to re-create the actual output and write it to the filesystem. This approach allows ParaMEDIC to reduce I/O time, thus accelerating mpiBLAST by as much as 25-fold. Pavan Balaji, Wu-chun Feng, Jeremy S. Archuleta, Heshan Lin, Rajkumar Kettimuthu, Rajeev Thakur, Xiaosong Ma |
PPoPP | 7 |
| 2008 | Massively parallel genomic sequence search on the Blue Gene/P architectureabstractThis paper presents our first experiences in mapping and optimizing genomic sequence search onto the massively parallel IBM Blue Gene/P (BG/P) platform. Specifically, we performed our work on mpiBLAST, a parallel sequence-search code that has been optimized on numerous supercomputing environments. In doing so, we identify several critical performance issues. Consequently, we propose and study different approaches for mapping sequence-search and parallel I/O tasks on such massively parallel architectures.We demonstrate that our optimizations can deliver nearly linear scaling (93% efficiency) on up to 32,768 cores of BG/P. In addition, we show that such scalability enables us to complete a large-scale bioinformatics problem - sequence searching a microbial genome database against itself to support the discovery of missing genes in genomes - in only a few hours on BG/P. Previously, this problem was viewed as computationally intractable in practice. Heshan Lin, Pavan Balaji, Ruth Poole, Carlos P. Sosa, Xiaosong Ma, Wu-chun Feng |
SC | 5 |
| 2008 | Adaptive Request Scheduling for Parallel Scientific Web Services
Heshan Lin, Xiaosong Ma, Jiangtian Li, Ting Yu 0001, Nagiza F. Samatova |
SSDBM | 2 |
| 2007 | Automatic Parallelization of Scripting Languages: Toward Transparent Desktop Parallel ComputingabstractDesktop computing remains indispensable in scientific exploration, largely because it provides people with devices for human interaction and environments for interactive job execution. However, with today's rapidly growing data volume and task complexity, it is increasingly hard for individual workstations to meet the demands of interactive scientific data processing. The increasing cost of such interactive processing is hindering the productivity of end-to-end scientific computing workflows. While existing distributed computing systems allow people to aggregate desktop workstation resources for parallel computing, the burden of explicit parallel programming and parallel job execution often prohibits scientists to take advantage of such platforms. In this paper, we discuss the need for transparent desktop parallel computing in scientific data processing. As an initial step toward this goal, we present our on-going work on the automatic parallelization of the scripting language R, a popular tool for statistical computing. Our preliminary results suggest that a reasonable speedup can be achieved on real-world sequential R programs without requiring any code modification. Xiaosong Ma, Jiangtian Li, Nagiza F. Samatova |
IPDPS | 1 |
| 2007 | Optimizing center performance through coordinated data staging, scheduling and recoveryabstractProcurement and the optimized utilization of Petascale supercomputers and centers is a renewed national priority. Sustained performance and availability of such large centers is a key technical challenge significantly impacting their usability. Storage systems are known to be the primary fault source leading to data unavailability and job resubmissions. This results in reduced center performance, partially due to the lack of coordination between I/O activities and job scheduling. Zhe Zhang 0005, Chao Wang 0056, Sudharshan S. Vazhkudai, Xiaosong Ma, Gregory G. Pike, John Cobb, Frank Mueller 0001 |
SC | 4 |
| 2006 | Positioning Dynamic Storage Caches for Transient DataabstractSimulations, experiments and observatories are generating a deluge of scientific data. Even more staggering is the ever growing application demand to process and assimilate these datasets. Application users perform a range of data operations, collaborate and share data in many novel ways. The current storage landscape is struggling to keep up with these trends in scientific data processing. Application users pay the price due to over-crowded shared filesystems, or expensive storage area networks, or not enough local storage, or high-latency archival or wide-area transfers. In order to sustain and maximize I/O bandwidth relative to increasing CPU speeds, applications must take advantage of large amounts of intermediate commodity storage, However, intermediate storage presents new challenges above and beyond the traditional distributed file system paradigm: persistent scheduling, storage/CPU coallo-cation, namespace management, lifetime management, and novel application interfaces. In this paper, we describe applications that require intermediate storage management, suggest several open research problems, and illustrate two systems - Freeloader and Tactical Storage - that attack different aspects of these problems Sudharshan S. Vazhkudai, Douglas Thain, Xiaosong Ma, Vincent W. Freeh |
CLUSTER | 3 |
| 2006 | Exploring I/O Strategies for Parallel Sequence-Search Tools with S3aSimabstractParallel sequence-search tools are rising in popularity among computational biologists. With the rapid growth of sequence databases, database segmentation is the trend of the future for such search tools. While I/O currently is not a significant bottleneck for parallel sequence-search tools, future technologies including faster processors, customized computational hardware such as FPGAs, improved search algorithms, and exponentially growing databases emphasize an increasing need for efficient parallel I/O in future parallel sequence-search tools. Our paper focuses on examining different I/O strategies for these future tools in a modern parallel file system (PVFS2). Because implementing and comparing various I/O algorithms in every search tool is labor-intensive and time-consuming, we introduce S3aSim, a general simulation framework for sequence-search which allows us to quickly implement, test, and profile various I/O strategies. We examine a variety of I/O strategies (e.g., master-writing and various worker-writing strategies using individual and collective I/O methods) for storing result data in sequence-search tools such as mpiBLAST, pioBLAST, and parallel HMMer. Our experiments fully detail the interaction of computing and I/O within a full application simulation as opposed to typical I/O-only benchmarks Avery Ching, Wu-chun Feng, Heshan Lin, Xiaosong Ma, Alok N. Choudhary |
HPDC | 4 |
| 2006 | Coupling prefix caching and collective downloads for remote dataset accessabstractScientific datasets are typically archived at mass storage systems or data centers close to supercomputers/instruments. End-users of these datasets, however, usually perform parts of their workflows at their local computers. In such cases, client-side caching can offer significant gains by reducing the cost of wide-area data movement.Scientific data caches, however, traditionally cache entire data-sets, which may not be necessary. In this paper, we propose a novel combination of prefix caching and collective download. Prefix caching allows the bootstrapping of dataset downloads by caching only a prefix of the dataset, while collective download facilitates efficient parallel patching of the missing suffix from an external data source. To estimate the optimal prefix size, we further present an analytical model that considers both the initial download over-head and the downloading speed. We implemented our proposed approach in the FreeLoader distributed cache prototype. Experimental results (using multiple scientific data repositories and data transfer tools, as well as a real-world scientific dataset access trace) demonstrate that prefix caching and collective download can be implemented efficiently, our model can select an appropriate prefix size, and the cache hit rate can be improved significantly without hurting the local access rate of cached datasets. Xiaosong Ma, Vincent W. Freeh, Sudharshan S. Vazhkudai, Tyler A. Simon, Stephen L. Scott |
ICS | 1 |
| 2006 | Grid applications - Parallel genomic sequence-searching on an ad-hoc grid: experiences, lessons learned, and implicationsabstractThe Basic Local Alignment Search Tool (BLAST) allows bioinformaticists to characterize an unknown sequence by comparing it against a database of known sequences. The similarity between sequences enables biologists to detect evolutionary relationships and infer biological properties of the unknown sequence.mpiBLAST, our parallel BLAST, decreases the search time of a 300 KB query on the current NT database from over two full days to under 10 minutes on a 128-processor cluster and allows larger query files to be compared. Consequently, we propose to compare the largest query available, the entire NT database, against the largest database available, the entire NT database. The result of this comparison will provide critical information to the biology community, including insightful evolutionary, structural, and functional relationships between every sequence and family in the NT database.Preliminary projections indicated that to complete the above task in a reasonable length of time required more processors than were available to us at a single site. Hence, we assembled GreenGene, an ad-hoc grid that was constructed "on the fly" from donated computational, network, and storage resources during last year's SC|05. GreenGene consisted of 3048 processors from machines that were distributed across the United States. This paper presents a case study of mpiBLAST on GreenGene --- specifically, a pre-run characterization of the computation, the hardware and software architectural design, experimental results, and future directions. Mark K. Gardner, Wu-chun Feng, Jeremy S. Archuleta, Heshan Lin, Xiaosong Ma |
SC | 5 |
| 2006 | Constructing collaborative desktop storage caches for large scientific datasetsabstractHigh-end computing is suffering a data deluge from experiments, simulations, and apparatus that creates overwhelming application dataset sizes. This has led to the proliferation of high-end mass storage systems, storage area clusters, and data centers. These storage facilities offer a large range of choices in terms of capacity and access rate, as well as strong data availability and consistency support. However, for most end-users, the “last mile” in their analysis pipeline often requires data processing and visualization at local computers, typically local desktop workstations. End-user workstations---despite having more processing power than ever before---are ill-equipped to cope with such data demands due to insufficient secondary storage space and I/O rates. Meanwhile, a large portion of desktop storage is unused.We propose the FreeLoader framework, which aggregates unused desktop storage space and I/O bandwidth into a shared cache/scratch space, for hosting large, immutable datasets and exploiting data access locality. This article presents the FreeLoader architecture, component design, and performance results based on our proof-of-concept prototype. Its architecture comprises contributing benefactor nodes, steered by a management layer, providing services such as data integrity, high performance, load balancing, and impact control. Our experiments show that FreeLoader is an appealing low-cost solution to storing massive datasets by delivering higher data access rates than traditional storage facilities, namely, local or remote shared file systems, storage systems, and Internet data repositories. In particular, we present novel data striping techniques that allow FreeLoader to efficiently aggregate a workstation's network communication bandwidth and local I/O bandwidth. In addition, the performance impact on the native workload of donor machines is small and can be effectively controlled. Further, we show that security features such as data encryptions and integrity checks can be easily added as filters for interested clients. Finally, we demonstrate how legacy applications can use the FreeLoader API to store and retrieve datasets. Sudharshan S. Vazhkudai, Xiaosong Ma, Vincent W. Freeh, Jonathan W. Strickland, Nandan Tammineedi, Tyler A. Simon, Stephen L. Scott |
ACM Trans. Storage | 2 |
| 2006 | High-Level Buffering for Hiding Periodic Output Cost in Scientific SimulationsabstractScientific applications often need to write out large arrays and associated metadata periodically for visualization or restart purposes. In this paper, we present active buffering, a high-level transparent buffering scheme for collective I/O, in which processors actively organize their idle memory into a hierarchy of buffers for periodic output data. It utilizes idle memory on the processors, yet makes no assumption regarding runtime memory availability. Active buffering can perform background I/O while the computation is going on, is extensible to remote I/O for more efficient data migration, and can be implemented in a portable style in today's parallel I/O libraries. It can also mask performance problems of scientific data formats used by many scientists. Performance experiments with both synthetic benchmarks and real simulation codes on multiple platforms show that active buffering can greatly reduce the visible I/O cost from the application's point of view. Xiaosong Ma, Jonghyun Lee 0001, Marianne Winslett |
IEEE Trans. Parallel Distributed Syst. | 1 |
| 2005 | FreeLoader: Scavenging Desktop Storage Resources for Scientific DataabstractHigh-end computing is suffering a data deluge from experiments, simulations, and apparatus that creates overwhelming application dataset sizes. End-user workstations-despite more processing power than ever before-are ill-equipped to cope with such data demands due to insufficient secondary storage space and I/O rates. Meanwhile, a large portion of desktop storage is unused. We present the FreeLoader framework, which aggregates unused desktop storage space and I/O bandwidth into a shared cache/scratch space, for hosting large, immutable datasets and exploiting data access locality. Our experiments show that FreeLoader is an appealing low-cost solution to storing massive datasets, by delivering higher data access rates than traditional storage facilities. In particular, we present novel data striping techniques that allow FreeLoader to efficiently aggregate a workstation’s network communication bandwidth and local I/O bandwidth. In addition, the performance impact on the native workload of donor machines is small and can be effectively controlled. Sudharshan S. Vazhkudai, Xiaosong Ma, Vincent W. Freeh, Jonathan W. Strickland, Nandan Tammineedi, Stephen L. Scott |
SC | 2 |
| 2005 | Cross-Platform Performance Prediction of Parallel Applications Using Partial ExecutionabstractPerformance prediction across platforms is increasingly important as developers can choose from a wide range of execution platforms. The main challenge remains to perform accurate predictions at a low-cost across different architectures. In this paper, we derive an affordable method approaching cross-platform performance translation based on relative performance between two platforms. We argue that relative performance can be observed without running a parallel application in full. We show that it suffices to observe very short partial executions of an application since most parallel codes are iterative and behave predictably manner after a minimal startup period. This novel prediction approach is observation-based. It does not require program modeling, code analysis, or architectural simulation. Our performance results using real platforms and production codes demonstrate that prediction derived from partial executions can yield high accuracy at a low cost. We also assess the limitations of our model and identify future research directions on observationbased performance prediction. Leo T. Yang, Xiaosong Ma, Frank Mueller 0001 |
SC | 2 |
| 2004 | RFS: efficient and flexible remote file access for MPI-IOabstractScientific applications often need to access remote file systems. Because of slow networks and large data size, however, remote I/O can become an even more serious performance bottleneck than local I/O performance. In this work, we present RFS, a high-performance remote I/O facility for ROMIO, which is a well-known MPI-IO implementation. Our simple, portable, and flexible design eliminates the shortcomings of previous remote I/O efforts. In particular, RFS improves the remote I/O performance by adopting active buffering with threads (ABT), which hides I/O cost by aggressively buffering the output data using available memory and performing background I/O using threads while computation is taking place. Our experimental results show that RFS with ABT can significantly reduce the remote I/O visible cost, achieving up to 92% of the theoretical peak throughput. The computation slowdown caused by concurrent I/O activities was 0.2-6.2%, which is dwarfed by the overall performance improvement in application turnaround time. Jonghyun Lee 0001, Robert B. Ross, Rajeev Thakur, Xiaosong Ma, Marianne Winslett |
CLUSTER | 4 |
| 2004 | GODIVA: Lightweight Data Management for Scientific Visualization ApplicationsabstractScientific visualization applications are very data-intensive, with high demands for I/O and data management. Developers of many visualization tools hesitate to use traditional DBMSs, due to the lack of support for these DBMSs on parallel platforms and the risk of reducing the portability of their tools and the user data. We propose the GODIVA framework, which provides simple database-like interfaces to help visualization tool developers manage their in-memory data, and I/O optimizations such as prefetching and caching to improve input performance at run time. We implemented the GODIVA interfaces in a stand-alone, portable user library, which can be used by all types of visualization codes: interactive and batch-mode, sequential and parallel. Performance results from running a visualization tool using the GODIVA library on multiple platforms show that the GODIVA framework is easy to use, alleviates developers' data management burden, and can bring substantial I/O performance improvement. Xiaosong Ma, Marianne Winslett, John Norris, Xiangmin Jiao, Robert Fiedler |
ICDE | 1 |
| 2003 | Declustering Large Multidimensional Data Sets for Range Queries over Heterogeneous DisksabstractDeclustering is a technique to distribute data sets over multiple disks so that future retrievals can be well balanced over the disks and be performed in parallel. Although clusters often have heterogeneous disks, most declustering work has focused only on homogeneous environments. In this work, we investigate the declustering problem for a heterogeneous disk environment using virtual servers, and propose approaches for deciding the number of virtual servers and the mapping between virtual servers and physical disks. Our experimental results show that by combining our algorithm for choosing the number of virtual servers with a greedy algorithm for mapping virtual servers to disks, users can expect range query retrieval performance within 4% of the optimum achievable in practice on average, in all configurations studied. Compared to an intuitively natural approach to the problem, this represents an improvement of 8-31% in average fetch ratio, as well a 26-38% reduction in the standard deviation of performance for small queries. Jonghyun Lee 0001, Marianne Winslett, Xiaosong Ma, Shengke Yu |
SSDBM | 3 |
| 2002 | Active buffering plus compressed migration: an integrated solution to parallel simulations' data transport needsabstractScientific simulations running on parallel platforms output intermediate data periodically, typically moving the output to a remote machine for visualization. Due to the large data size and slow network, improvements in output and migration performance can significantly reduce simulation turnaround time. In this paper, we propose a novel simulation execution environment that integrates active buffering and compressed migration to hide and reduce the cost of output and migration. Our implementation removes previous shortcomings of active buffering and of compressed migration, and our experiments with real world data sets show that integrating the two techniques brings comparable application-visible I/O performance and much more flexibility compared to compressed migration without active buffering. The performance advantage will increase with slower-to-write file formats such as HDF. Jonghyun Lee 0001, Xiaosong Ma, Marianne Winslett, Shengke Yu |
ICS | 2 |
| 2001 | Tuning high-performance scientific codes: the use of performance models to control resource usage during data migration and I/OabstractLarge-scale parallel simulations are a popular tool for investigating phenomena ranging from nuclear explosions to protein folding. These codes produce copious output that must be moved to the workstation where it will be visualized. Scientists have a variety of tools to help them with this data movement, and often have several different platforms available to them for their runs. Thus questions arise such as, which data migration approach is best for a particular code and platform? Which will provide the best end-to-end response time, or lowest cost? Scientists also control how much data is output, and how often. From a scientific perspective, the more output the better; but from a cost and response time perspective, how much output is too much? To answer these questions, we built performance models for data migration approaches and verified them on parallel and sequential platforms. We use a 3D hydrodynamics code to show how scientists can use the models to predict performance and tune the I/O aspects of their codes. Jonghyun Lee 0001, Marianne Winslett, Xiaosong Ma, Shengke Yu |
ICS | 3 |
| 2000 | PRUNES: an efficient and complete strategy for automated trust negotiation over the InternetabstractArticle Free Access Share on PRUNES: an efficient and complete strategy for automated trust negotiation over the Internet Authors: Ting Yu Department of Computer Science, University of Illinois at Urbana-Champaign, 1304 W.Springfield Ave., Urbana, IL Department of Computer Science, University of Illinois at Urbana-Champaign, 1304 W.Springfield Ave., Urbana, ILView Profile , Xiaosong Ma Department of Computer Science, University of Illinois at Urbana-Champaign, 1304 W.Springfield Ave., Urbana, IL Department of Computer Science, University of Illinois at Urbana-Champaign, 1304 W.Springfield Ave., Urbana, ILView Profile , Marianne Winslett Department of Computer Science, University of Illinois at Urbana-Champaign, 1304 W.Springfield Ave., Urbana, IL Department of Computer Science, University of Illinois at Urbana-Champaign, 1304 W.Springfield Ave., Urbana, ILView Profile Authors Info & Claims CCS '00: Proceedings of the 7th ACM conference on Computer and Communications SecurityNovember 2000 Pages 210–219https://doi.org/10.1145/352600.352633Published:01 November 2000Publication History 50citation960DownloadsMetricsTotal Citations50Total Downloads960Last 12 Months37Last 6 weeks5 Get Citation AlertsNew Citation Alert added!This alert has been successfully added and will be sent to:You will be notified whenever a record that you have chosen has been cited.To manage your alert preferences, click on the button below.Manage my AlertsNew Citation Alert!Please log in to your account Save to BinderSave to BinderCreate a New BinderNameCancelCreateExport CitationPublisher SiteeReaderPDF Ting Yu 0001, Xiaosong Ma, Marianne Winslett |
CCS | 2 |