Kenjiro Taura

dblp:42/1786 · DBLP profile ↗
← Back
75ranked-venue papers
7as first author
18since 2021 · last 2026
0000-0001-5224-382XORCID · verified

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

Systems, architecture and hardware · 52 · 5 first-author · 7 since 2021Databases, data management, data science and information retrieval · 7 · 2 since 2021Artificial intelligence and machine learning · 6 · 3 since 2021Software engineering, systems software and programming languages · 5 · 2 first-author · 2 since 2021Applied, interdisciplinary, general and emerging computing · 5 · 1 first-author · 1 since 2021Human-computer interaction and ubiquitous computing · 3 · 3 since 2021Graphics, computer vision, multimedia, augmented reality and games · 2 · 1 since 2021
YearPublicationVenuePosition
2026 Importance-Aware Data Selection for Efficient LLM Instruction Tuning
abstract
Instruction tuning plays a critical role in enhancing the performance and efficiency of Large Language Models (LLMs). Its success depends not only on the quality of the instruction data but also on the inherent capabilities of the LLM itself. Some studies suggest that even a small amount of high-quality data can achieve instruction fine-tuning results that are on par with, or even exceed, those from using a full-scale dataset. However, rather than focusing solely on calculating data quality scores to evaluate instruction data, there is a growing need to select high-quality data that maximally enhances the performance of instruction tuning for a given LLM. In this paper, we propose the Model Instruction Weakness Value (MIWV) as a novel metric to quantify the importance of instruction data in enhancing model's capabilities. The MIWV metric is derived from the discrepancies in the model’s responses when using In-Context Learning (ICL), helping identify the most beneficial data for enhancing instruction tuning performance. Our experimental results demonstrate that selecting only the top 1% of data based on MIWV can outperform training on the full dataset. Furthermore, this approach extends beyond existing research that focuses on data quality scoring for data selection, offering strong empirical evidence supporting the effectiveness of our proposed method.
Tingyu Jiang, Yiyao Song, Hualei Zhu, Xiaohang Xu 0002, Kenjiro Taura, Hao Henry Wang
AAAI8
2026 LLM-Based Explainable Detection of LLM-Generated Code in Python Programming Courses
Jeonghun Baek, Tetsuro Yamazaki, Akimasa Morihata, Junichiro Mori, Yoko Yamakata, Kenjiro Taura, Shigeru Chiba
SIGCSE (1)6
2026 MaskingAgent: Preventing LLM Tutor from Providing Full Solutions in Python Programming Courses
Jeonghun Baek, Tetsuro Yamazaki, Akimasa Morihata, Junichiro Mori, Yoko Yamakata, Kenjiro Taura, Shigeru Chiba
SIGCSE (2)6
2025 How Different from the Past? Spatio-Temporal Time Series Forecasting with Self-Supervised Deviation Learning
abstract
Spatio-temporal forecasting is essential for real-world applications such as traffic management and urban computing. Although recent methods have shown improved accuracy, they often fail to account for dynamic deviations between current inputs and historical patterns. These deviations contain critical signals that can significantly affect model performance. To fill this gap, we propose $\textbf{ST-SSDL}$, a $\underline{S}$patio-$\underline{T}$emporal time series forecasting framework that incorporates a $\underline{S}$elf-$\underline{S}$upervised $\underline{D}$eviation $\underline{L}$earning scheme to capture and utilize such deviations. ST-SSDL anchors each input to its historical average and discretizes the latent space using learnable prototypes that represent typical spatio-temporal patterns. Two auxiliary objectives are proposed to refine this structure: a contrastive loss that enhances inter-prototype discriminability and a deviation loss that regularizes the distance consistency between input representations and corresponding prototypes to quantify deviation. Optimized jointly with the forecasting objective, these components guide the model to organize its hidden space and improve generalization across diverse input conditions. Experiments on six benchmark datasets show that ST-SSDL consistently outperforms state-of-the-art baselines across multiple metrics. Visualizations further demonstrate its ability to adaptively respond to varying levels of deviation in complex spatio-temporal scenarios. Our code and datasets are available at https://github.com/Jimmy-7664/ST-SSDL.
Zheng Dong 0006, Jiawei Yong, Shintaro Fukushima, Kenjiro Taura, Renhe Jiang
NeurIPS5
2025 Fairer and More Scalable Reader-Writer Locks by Optimizing Queue Management
abstract
MCS lock and its variants provide scalability on many-core architectures, using lists of lock requests to reduce access contention on the mutex data. Most recent variants have adopted a two-stage design, allowing requests to be allocated from stack memory rather than heap. However, this design still produces mutex access contention and limits fairness in the presence of fast paths. This paper proposes the Freezer mechanism and its optimization methods, which extend the list structure operations of MCS lock, to reduce mutex access without using heap memory and enable an independent choice of fairness policies and fast paths. Additionally, we propose four optimization methods for queue-based reader-writer locks. Our evaluation using three benchmarks demonstrated the effectiveness of the proposed fair reader-writer locks. They achieved up to 3.5× higher throughput and improved tail latency by up to 2.7× compared to the baselines.
Takashi Hoshino 0002, Kenjiro Taura
PPoPP2
2025 Leveraging LLM for Detecting and Explaining LLM-generated Code in Python Programming Courses
Jeonghun Baek, Tetsuro Yamazaki, Akimasa Morihata, Junichiro Mori, Yoko Yamakata, Kenjiro Taura, Shigeru Chiba
SIGCSE (2)6
2024 Efficiently Adapting Stateless Model Checking for C11/C++11 to Mixed-Size Accesses
abstract
Abstract Stateless model checking (SMC) is crucial for productivity in verified concurrent programming, and its recent developments for C/C++ and weak memory models are remarkable. The state-of-the-art SMC for C, GenMC, efficiently verifies C programs based on C11 atomics and pthreads. However, it does not support mixed-size accesses, accesses to the same memory region with different-sized types, even though they are ubiquitous in C/C++, particularly the code for memory management. As a result, GenMC does not work for C/C++ programs containing memory management. To resolve this problem, we develop a method of adapting GenMC to mixed-size accesses preserving its optimality. We experimentally evaluate the efficiency of our extended implementation of GenMC and its efficacy for memory management programs.
Shigeyuki Sato 0001, Taiyo Mizuhashi, Genki Kimura, Kenjiro Taura
APLAS4
2024 ARIM-mdx Data System: Towards a Nationwide Data Platform for Materials Science
abstract
In modern materials science, effective and high-volume data management across leading-edge experimental facilities and world-class supercomputers is indispensable for cutting-edge research. However, existing integrated systems that handle data from these resources have primarily focused just on smaller-scale cross-institutional or single-domain operations. As a result, they often lack the scalability, efficiency, agility, and interdisciplinarity, needed for handling substantial volumes of data from various researchersIn this paper, we introduce ARIM-mdx data system1, aiming at a nationwide data platform for materials science in Japan. Currently in its trial phase, the platform has been involving 11 universities and institutes all over Japan, and it is utilized by over 800 researchers from around 140 organizations in academia and industry, being intended to gradually expand its reach. The ARIM-mdx data system, as a pioneering nationwide data platform, has the potential to contribute to the creation of new research communities and accelerate innovations.
Masatoshi Hanai, Ryo Ishikawa, Mitsuaki Kawamura, Masato Ohnishi, Norio Takenaka, Kou Nakamura, Daiju Matsumura, Seiji Fujikawa, Hiroki Sakamoto, Yukinori Ochiai, Tetsuo Okane, Shin-Ichiro Kuroki, Atsuo Yamada, Toyotaro Suzumura, Kenjiro Taura, Yoshio Mita, Naoya Shibata, Yuichi Ikuhara
IEEE Big Data15
2024 Enhancing Privacy of Spatiotemporal Federated Learning Against Gradient Inversion Attacks
Lele Zheng, Yang Cao 0011, Renhe Jiang, Kenjiro Taura, Yulong Shen 0001, Sheng Li 0010, Masatoshi Yoshikawa
DASFAA (1)4
2023 Associative Operator Precedence Parsing: A Method To Increase Data Parsing Parallelism
abstract
Many data often come with a high volume in textual format (JSON, XML, CSV). Because parsing can easily dominate data analysis time, researchers have been working on parallelizing parsing. Operator Precedence Parsing (OPP), among candidate parsing methods, is amenable to parallelization, with a practical algorithm proposed. The “locally parsable” property allows the parser to deduce if a reduction is safe with limited context. However, when the grammar has productions that tend to produce a highly skewed parse tree, OPP raises reductions mostly in serial, and the parsing still suffers from a long critical path. In pactice, OPP has little or even no speedup when parsing data because data often contain high percentage of parallel elements (e.g., JSON array elements separated by commas) produced from such productions, a situation that frequently occurs when processing big data.
Kenjiro Taura
HPC Asia2
2023 Itoyori: Reconciling Global Address Space and Global Fork-Join Task Parallelism
abstract
This paper introduces Itoyori, a task-parallel runtime system designed to tackle the challenge of scaling task parallelism (more specifically, nested fork-join parallelism) beyond a single node. The partitioned global address space (PGAS) model is often employed in task-parallel systems, but naively combining them can lead to poor performance due to fine-grained and redundant remote memory accesses. Itoyori addresses this issue by automatically caching global memory accesses at runtime, enabling efficient cache sharing among parallel tasks running on the same processor. As a real-world case study, we ported an existing task-parallel implementation of the Fast Multipole Method (FMM) to distributed memory with Itoyori and achieved a 7.5× speedup when scaled from a single node to 12 nodes and up to 6.0× faster performance than without caching. This study demonstrates that global-view fork-join programming can be made practical and scalable, while requiring minimal changes to the shared-memory code.
Shumpei Shiina, Kenjiro Taura
SC2
2022 Distributed Continuation Stealing is More Scalable than You Might Think
abstract
The need for load balancing in applications with irregular parallelism has motivated research on work stealing. An important choice in work-stealing schedulers is between child stealing or continuation stealing. In child stealing, a newly created task is made stealable by other processors, whereas in continuation stealing, the caller's continuation is made stealable by executing the newly created task first, which preserves the serial execution order. Although the benefits of continuation stealing have been demonstrated on shared memory by Cilk and other runtime systems, it is rarely employed on distributed memory, presumably because it has been thought to be difficult to implement and inefficient as it involves migration of call stacks across nodes. Akiyama and Taura recently introduced efficient RDMA-based continuation stealing, but the practicality of distributed continuation stealing is still unclear because a comparison of its performance with that of child stealing has not previously been performed. This paper presents the results of a comparative performance analysis of continuation stealing and child stealing on distributed memory. To clarify the full potential of continuation stealing, we first investigated various RDMA-based synchronization (task join) implementations, which had not previously been fully inves-tigated. The results revealed that, when the task synchronization pattern was complicated, continuation stealing performed better than child stealing despite its relatively long steal latency due to stack migration. Notably, our runtime system achieved almost perfect scaling on 110,592 cores in an unbalanced tree search (UTS) benchmark. This scalability is comparable to or even better than that of state-of-the-art bag-of-tasks counterparts.
Shumpei Shiina, Kenjiro Taura
CLUSTER2
2022 SimdFSM: An Adaptive Vectorization of Finite State Machines for Speculative Execution
Kenjiro Taura
PDCAT2
2022 Improving Cache Utilization of Nested Parallel Programs by Almost Deterministic Work Stealing
abstract
Nested (fork-join) parallelism eases parallel programming by enabling high-level expression of parallelism and leaving the mapping between parallel tasks and hardware to the runtime scheduler. A challenge in dynamic scheduling of nested parallelism is how to exploit data locality, which has become more demanding in the deep cache hierarchies of modern processors with a large number of cores. This paper introducesalmost deterministic work stealing (ADWS), which efficiently exploits data locality by deterministically planning a cache-hierarchy-aware schedule, while allowing a little scheduling variety to facilitate dynamic load balancing. Furthermore, we propose an extension of our prior work on ADWS to achieve better shared cache utilization. The improved version of the scheduler is calledmulti-level ADWS. The idea is that only part of a computation whose working set size is small enough to fit into a shared cache is scheduled by ADWS within the cache recursively, thus avoiding excessive capacity misses. Our evaluation on a benchmark of parallel decision tree construction demonstrated that multi-level ADWS outperformed the conventional random work stealing of Cilk Plus by 61% and it showed a 40% performance improvement over the previous ADWS design.
Shumpei Shiina, Kenjiro Taura
IEEE Trans. Parallel Distributed Syst.2
2021 Plex: Scaling Parallel Lexing with Backtrack-Free Prescanning
abstract
Lexical analysis, which converts input text into a list of tokens, plays an important role in many applications, including compilation and data extraction from texts. To recognize token patterns, a lexer incorporates a sequential computation model - automaton as its basic building component. As such, it is considered difficult to parallelize due to the inherent data dependency. Much work has been done to accelerate lexical analysis through parallel techniques. Unfortunately, existing attempts mainly rely on language-specific remedies for input segmentation, which makes it not only tricky for language extension, but also challenging for automatic lexer generation. This paper presents Plex - an automated tool for generating parallel lexers from user-defined grammars. To overcome the inherent sequentiality, Plex applies a fast prescanning phase to collect context information prior to scanning. To reduce the overheads brought by prescanning, Plex adopts a special automaton, which is derived from that of the scanner, to avoid backtracking behavior and exploits data-parallel techniques. The evaluation under several languages shows that the prescanning overhead is small, and consequently Plex is scalable and achieves 9.8-11.5X speedups using 18 threads.
Shigeyuki Sato 0001, Qiheng Liu, Kenjiro Taura
IPDPS4
2021 Automatic Graph Partitioning for Very Large-scale Deep Learning
abstract
This work proposes RaNNC (Rapid Neural Network Connector) as middleware for automatic hybrid parallelism. In recent deep learning research, as exemplified by T5 and GPT-3, the size of neural network models continues to grow. Since such models do not fit into the memory of accelerator devices, they need to be partitioned by model parallelism techniques. Moreover, to accelerate training for huge training data, we need a combination of model and data parallelisms, i.e., hybrid parallelism. Given a model description for PyTorch without any specification for model parallelism, RaNNC automatically partitions the model into a set of subcomponents so that (1) each subcomponent fits a device memory and (2) a high training throughput for pipeline parallelism is achieved by balancing the computation times of the subcomponents. Since the search space for partitioning models can be extremely large, RaNNC partitions a model through the following three phases. First, it identifies atomic subcomponents using simple heuristic rules. Next it groups them into coarser-grained blocks while balancing their computation times. Finally, it uses a novel dynamic programming-based algorithm to efficiently search for combinations of blocks to determine the final partitions. In our experiments, we compared RaNNC with two popular frameworks, Megatron-LM (hybrid parallelism) and GPipe (originally proposed for model parallelism, but a version allowing hybrid parallelism also exists), for training models with increasingly greater numbers of parameters. In the pre-training of enlarged BERT models, RaNNC successfully trained models five times larger than those Megatron-LM could, and RaNNC's training throughputs were comparable to Megatron-LM's when pre-training the same models. RaNNC also achieved better training throughputs than GPipe on both the enlarged BERT model pre-training (GPipe with hybrid parallelism) and the enlarged ResNet models (GPipe with model parallelism) in all of the settings we tried. These results are remarkable, since RaNNC automatically partitions models without any modification to their descriptions; Megatron-LM and GPipe require users to manually rewrite the models' descriptions.
Masahiro Tanaka, Kenjiro Taura, Toshihiro Hanawa, Kentaro Torisawa
IPDPS2
2021 Pitfalls of InfiniBand with On-Demand Paging
abstract
InfiniBand is a popular high-performance interconnect and offers Remote Direct Memory Access (RDMA), which enables low-latency communication based on kernel bypassing. Although the conventional RDMA technology necessitates manual physical memory management, an emerging extension, On-Demand Paging (ODP), implements automatic memory management based on RDMA-triggered page faults, which benefits productivity. Although the existing studies said the overhead of a page fault of ODP to be small enough, an in-depth investigation in various network situations including retransmission and timeout is missing. In this work, we conduct a comprehensive analysis of the actual behaviors of ODP on different devices and reveal two awful performance pitfalls, which incur longer latencies 3-4 orders of magnitude than a common-case page fault does. We also experimentally demonstrate that the revealed pitfalls are harmful to existing software systems. This paper presents our experimental analysis and lessons learned therefrom.
Takuya Fukuoka, Shigeyuki Sato 0001, Kenjiro Taura
ISPASS3
2021 Lightweight preemptive user-level threads
abstract
Many-to-many mapping models for user- to kernel-level threads (or "M:N threads") have been extensively studied for decades as a lightweight substitute for current Pthreads implementations that provide a simple one-to-one mapping ("1:1 threads"). M:N threads derive performance from their ability to allow users to context switch between threads and control their scheduling entirely in user space with no kernel involvement. This same ability, however, causes M:N threads to lose the kernel-provided ability of implicit OS preemption---threads have to explicitly yield control for other threads to be scheduled. Hence, programs over nonpreemptive M:N threads can cause core starvation, loss of prioritization, and, sometimes, deadlock unless programs are written to explicitly yield in proper places. This paper explores two techniques for M:N threads to efficiently achieve implicit preemption similar to 1:1 threads: signal-yield and KLT-switching. Overheads of these techniques, with our optimizations, can be less than 1% compared with nonpreemptive M:N threads. Our evaluation with three applications demonstrates that our preemption techniques for M:N threads improve core utilization and enhance the performance by utilizing lightweight context switching and flexible scheduling of M:N threads.
Shumpei Shiina, Shintaro Iwasaki, Kenjiro Taura, Pavan Balaji
PPoPP3
2020 On the Correct Measurement of Application Memory Bandwidth and Memory Access Latency
abstract
Diagnosing if an application suffers from DRAM contention can be a challenging task. One method is to compare the hardware memory bandwidth limit with the measured memory bandwidth of an application. Another method is based on memory access latency. The latency of a DRAM access in an uncontended state is a hardware characteristic. If an application shows higher DRAM access latency, the increase comes from queuing delays and the application is limited by DRAM bandwidth. Hardware-based measurement of the application's latency and bandwidth can be done with low-overhead and is agnostic of the application's implementation.
Christian Helm, Kenjiro Taura
HPC Asia2
2020 Automatic Identification and Precise Attribution of DRAM Bandwidth Contention
abstract
The limited DRAM bandwidth of today’s computing systems is a bottleneck for many applications. But the identification of DRAM bandwidth contention in applications is difficult. The measured bandwidth consumption of an application can not identify bandwidth contention. In theory, NUMA systems can provide higher memory bandwidth. But applications often make poor use of the resources. To address these challenges, we introduce a novel method to identify DRAM bandwidth contention and bad usage of NUMA resources. It consists of metrics to judge the severity of bandwidth contention and the degree of imbalanced resource usage in NUMA systems. Our tool can automatically scan an application for DRAM contention problems. Together with the precise location of the origin, intuitive optimization guidance is given. This approach for finding DRAM contention is based on the memory access latency. By comparing the experienced latency of an application with the uncontended hardware latency, we can find the contention. Hardware instruction sampling enables precise identification of the origin. It also provides information about accessed memories, which we use to calculate a NUMA imbalance metric. In a detailed evaluation with several micro-benchmarks, we show that our new method does indeed quantify the severity of DRAM contention beyond the possibilities of simple bandwidth consumption measurement and existing tools. We also apply our approach to real applications and confirm that it gives useful optimization advice to users.
Christian Helm, Kenjiro Taura
ICPP2
2020 Reliable Reverse Engineering of Intel DRAM Addressing Using Performance Counters
abstract
The memory controller of a processor translates the physical memory address to hardware components such as memory channels, ranks, and banks. This DRAM address mapping is of interest to many researchers in the fields of IT security, hardware architecture, system software, and performance tuning. However, Intel processors are using a complex and undocumented DRAM addressing. The addressing can be different for every system because it depends on many aspects such as the processor model, DIMM population on the motherboard, and BIOS settings. Thus an analysis for every individual system is necessary. In this paper, we introduce an automatic and reliable method for reverse engineering the DRAM addressing of Intel server-class processors. In contrast to existing approaches, it is reliable, measurement errors are unlikely to occur, and can be detected if they occur. Our method mainly relies on CPU hardware performance counters to precisely locate the accessed DRAM component. It eliminates the problem of wrong attribution that is common in timing based approaches. We validated our method by reversing engineering the DRAM addressing of a diverse set of Intel processors. This set includes Broadwell, Haswell, and Skylake micro-architectures, with various core counts, DIMM arrangements, and BIOS settings. We show the correctness of the determined addressing functions using micro-benchmarks that access specific DRAM components.
Christian Helm, Soramichi Akiyama, Kenjiro Taura
MASCOTS3
2020 Parallelizing and optimizing neural Encoder-Decoder models without padding on multi-core architecture
Yuchen Qiao, Kazuma Hashimoto, Akiko Eriguchi, Haixia Wang 0001, Dongsheng Wang 0002, Yoshimasa Tsuruoka, Kenjiro Taura
Future Gener. Comput. Syst.7
2020 Analyzing the Performance Trade-Off in Implementing User-Level Threads
abstract
User-level threads have been widely adopted as a means of achieving lightweight concurrent execution without the costs of OS-level threads. Nevertheless, the costs of managing user-level threads represent a performance barrier that dictates how fine grained the concurrency exposed by an application can be without incurring significant overheads; this in turn may translate into insufficient parallelism to exploit highly parallel systems. This article is a deep dive into the fundamental costs in implementing user-level threads. We first identify that one of the highest sources of fork-join overheads stems from deviations, events that incur context switching during the execution of a thread and disrupt a run-to-completion execution. We then conduct an in-depth investigation of a wide spectrum of methods with respect to how they handle deviations while covering both parent- and child-first scheduling policies. Our methodology involves a comprehensive instruction- and cache-level analysis of all methods on several modern CPU architectures. The primary finding of our evaluation is that dynamic promotion methods that assume the absence of deviation and dynamically provide context-switching support offer the best trade-off between performance and capability when the likelihood of deviation is low.
Shintaro Iwasaki, Abdelhalim Amer, Kenjiro Taura, Pavan Balaji
IEEE Trans. Parallel Distributed Syst.3
2019 BOLT: Optimizing OpenMP Parallel Regions with User-Level Threads
abstract
OpenMP is widely used by a number of applications, computational libraries, and runtime systems. As a result, multiple levels of the software stack use OpenMP independently of one another, often leading to nested parallel regions. Although exploiting such nested parallelism is a potential opportunity for performance improvement, it often causes destructive performance with leading OpenMP runtimes because of their reliance on heavyweight OS-level threads. User-level threads (ULTs) are more lightweight alternatives but existing ULT-based runtimes suffer from several shortcomings: 1) thread management costs remain significant and outweigh the benefits from additional parallelism; 2) the shift to ULTs often hurts the more common flat parallelism case; and 3) absence of user control over thread-to-CPU binding, a critical feature on modern systems. This paper presents BOLT, a practical ULT-based OpenMP runtime system that efficiently supports both flat and nested parallelism. This is accomplished on three fronts: 1) advanced data reuse and thread synchronization strategies; 2) thread coordination that adapts to the level of oversubscription; and 3) an implementation of the modern OpenMP thread-to-CPU binding interface tailored to ULT-based runtimes. The result is a highly optimized runtime that transparently achieves similar performance compared with leading state-of-the-art widely used OpenMP runtimes under flat parallelism, while outperforming all existing runtimes under nested parallelism.
Shintaro Iwasaki, Abdelhalim Amer, Kenjiro Taura, Pavan Balaji
PACT3
2019 Software combining to mitigate multithreaded MPI contention
abstract
Efforts to mitigate lock contention from concurrent threaded accesses to MPI have reduced contention through fine-grained locking, avoided locking altogether by offloading communication to dedicated threads, or alleviated negative side effects from contention by using better lock management protocols. The blocking nature of lock-based methods, however, wastes the asynchrony benefits of nonblocking MPI operations, and the offloading model sacrifices CPU resources and incurs unnecessary software offloading overheads under low contention.
Abdelhalim Amer, Charles Archer, Michael Blocksome, Chongxiao Cao, Michael Chuvelev, Hajime Fujita 0002, María Jesús Garzarán, Yanfei Guo, Jeff R. Hammond, Shintaro Iwasaki, Kenneth Raffenetti, Mikhail Shiryaev, Min Si, Kenjiro Taura, Sagar Thapaliya, Pavan Balaji
ICS14
2019 Almost deterministic work stealing
abstract
With task parallel models, programmers can easily parallelize divide-and-conquer algorithms by using nested fork-join structures. Work stealing, which is a popular scheduling strategy for task parallel programs, can efficiently perform dynamic load balancing; however, it tends to damage data locality and does not scale well with memory-bound applications. This paper introduces Almost Deterministic Work Stealing (ADWS), which addresses the issue of data locality of traditional work stealing by making the scheduling almost deterministic. Specifically, ADWS consists of two parts: (i) deterministic task allocation, which deterministically distributes tasks to workers based on the amount of work for each task, and (ii) hierarchical localized work stealing, which dynamically compensates load imbalance in a locality-aware manner. Experimental results show that ADWS is up to nearly 6 times faster than a traditional work stealing scheduler with memory-bound applications, and that dynamic load balancing works well while maintaining good data locality.
Shumpei Shiina, Kenjiro Taura
SC2
2018 Parallelized Software Offloading of Low-Level Communication with User-Level Threads
abstract
Although recent HPC interconnects are assumed to achieve low latency and high bandwidth communication, in practical terms, their performance is often bounded by the network software stacks rather than the underlying hardware because message processing requires a certain amount of computation in CPUs. To exploit the hardware capacity, some existing communication libraries provide an interface for parallelizing accesses to network endpoints with manual hints. However, with growing core counts per node in modern clusters, it is increasingly difficult for users to efficiently handle communication resources in multi-threading environments.
Wataru Endo, Kenjiro Taura
HPC Asia2
2018 Effectiveness of Moldable and Malleable Scheduling in Deep Learning Tasks
abstract
Research and development of deep learning (DL) applications often involves exhaustive trial-and-error, which demands that shared computational resources, especially GPUs, be efficiently allocated. Most DL tasks are moldable or malleable (i.e., the number of allocated GPUs can be changed before or during execution). However, conventional batch schedulers do not take advantage of DL tasks' moldability/malleability, inhibiting speedup when some GPU resources are unallocated. Another opportunity for speedup is to run multiple tasks concurrently on one GPU, which may improve the overall throughput because a single task does not always fully utilize the GPU's computational resources. We propose designing a batch scheduling system that exploits these opportunities to accelerate DL tasks. As a first step, this study conducts an extensive case study to evaluate the speedup of DL tasks when a scheduler treats them as moldable or malleable. That is, the scheduler adjusts the number of GPUs to be (or already) allocated to a task in response to the fluctuating availability of GPUs. Simulations using our real workload trace show that if the scheduler can allocate 1-4 GPUs to a task or assign 1-4 tasks to a GPU, then the average flow time of moldable/malleable DL tasks is shortened by at least 15.1 %/42.5 %, respectively, compared to a Rigid FCFS schedule in which one GPU is allocated to each task.
Ikki Fujiwara, Masahiro Tanaka, Kenjiro Taura, Kentaro Torisawa
ICPADS3
2018 Lessons learned from analyzing dynamic promotion for user-level threading
Shintaro Iwasaki, Abdelhalim Amer, Kenjiro Taura, Pavan Balaji
SC3
2018 Argobots: A Lightweight Low-Level Threading and Tasking Framework
abstract
In the past few decades, a number of user-level threading and tasking models have been proposed in the literature to address the shortcomings of OS-level threads, primarily with respect to cost and flexibility. Current state-of-the-art user-level threading and tasking models, however, either are too specific to applications or architectures or are not as powerful or flexible. In this paper, we present Argobots, a lightweight, low-level threading and tasking framework that is designed as a portable and performant substrate for high-level programming models or runtime systems. Argobots offers a carefully designed execution model that balances generality of functionality with providing a rich set of controls to allow specialization by end users or high-level programming models. We describe the design, implementation, and performance characterization of Argobots and present integrations with three high-level models: OpenMP, MPI, and colocated I/O services. Evaluations show that (1) Argobots, while providing richer capabilities, is competitive with existing simpler generic threading runtimes; (2) our OpenMP runtime offers more efficient interoperability capabilities than production OpenMP runtimes do; (3) when MPI interoperates with Argobots instead of Pthreads, it enjoys reduced synchronization costs and better latency-hiding capabilities; and (4) I/O services with Argobots reduce interference with colocated applications while achieving performance competitive with that of a Pthreads approach.
Abdelhalim Amer, Pavan Balaji, Cyril Bordage, George Bosilca, Alex Brooks, Philip H. Carns, Adrián Castelló 0001, Damien Genet, Thomas Hérault, Shintaro Iwasaki, Prateek Jindal, Laxmikant V. Kalé, Sriram Krishnamoorthy, Jonathan Lifflander, Huiwei Lu, Esteban Meneses, Marc Snir, Yanhua Sun, Kenjiro Taura, Pete Beckman
IEEE Trans. Parallel Distributed Syst.20
2017 Delay Spotter: A Tool for Spotting Scheduler-Caused Delays in Task Parallel Runtime Systems
abstract
Modern task parallel programming models provide sophisticated runtime task schedulers for handling the scheduling of logical tasks on a large and varying number of hardware parallel resources at runtime. The performance of these programming models increasingly rely on how fast their runtime schedulers do their job. The more delay a scheduler incurs in matching a ready task to a free processor core at any point in time, the more impact it causes to the program's parallel execution. We have developed a tool that is able to detect these delayed intervals caused by the scheduler in a parallel execution, and spot them specifically on two kinds of visualizations: the logical task graph captured at runtime (DAG visualizations) and time-series visualizations of threads (timelines). By further analyzing positions of these delays on those visualizations the tool could identify possible scheduling issues in the scheduler that causes these delays, yielding improvement insights for the development of task parallel programming models. From an application programmer's perspective, our tool is useful by being able to contrast differences of various task parallel programming models executing the same program, helping users choose the right model for their application. We demonstrate that usefulness by using the tool to analyze 10 applications in BOTS benchmark suite in our case studies.
An Huynh, Kenjiro Taura
CLUSTER2
2017 Autonomic Resource Management for Program Orchestration in Large-Scale Data Analysis
abstract
Large-scale data analysis applications are becoming more and more prevalent in a wide variety of areas. These applications are composed of many currently available programs called analysis components. Thousands of analysis component processes are orchestrated on many compute nodes. This paper proposes a novel self-tuning framework for optimizing an application's throughput in large-scale data analysis. One challenge is developing efficient orchestration that takes into account the diversity of analysis components and the varying performances of compute nodes. In our previous work, we achieved such an orchestration to a certain degree by introducing our own middleware, which wraps each analysis component as a remote procedure call (RPC) service. The middleware also pools the processes to reduce startup overhead, which is a serious obstacle to achieving high throughput. This work tackles the remaining task of tuning the size of the analysis components' process pools to maximize the application's throughput. This is challenging because analysis components differ drastically in turnaround times and memory footprints. The size of the process pool for each type of analysis component should be set by giving consideration to these properties as well as the constraints on both the memory capacity and the processor core counts. In this work, we formulate this task as a linear programming problem and obtain the optimal pool sizes by solving it. Compared to our previous work, we significantly improved the scalability of our framework by reformulating the performance model to work on hundreds of heterogeneous nodes. We also extended the service allocation mechanism to manage the computational load on each compute node and reduce communication overhead. The experimental results show that our approach is scalable to thousands of analysis component processes running on 200 compute nodes across three clusters. Moreover, our approach significantly reduces memory footprint.
Masahiro Tanaka, Kenjiro Taura, Kentaro Torisawa
IPDPS2
2017 SDAC: Porting Scientific Data to Spark RDDs
Kenjiro Taura, Liu Chao
NPC2
2016 A Static Cut-off for Task Parallel Programs
abstract
Task parallel models supporting dynamic and hierarchical parallelism are believed to offer a promising direction to achieving higher performance and programmability. Divide-and-conquer is the most frequently used idiom in task parallel models, which decomposes the problem instance into smaller ones until they become "trivial" to solve. However, it incurs a high tasking overhead if a task is created for each subproblem. In order to reduce this overhead, a "cut-off" is commonly used, which eliminates task creations where they are unlikely to be beneficial. The manual cut-off typically enlarges leaf tasks by stopping task creations when a subproblem becomes smaller than a threshold, and possibly transforms the enlarged leaf tasks into specialized versions for solving small instances (e.g., use loops instead of recursive calls); it duplicates the coding work and hinders productivity.
Shintaro Iwasaki, Kenjiro Taura
PACT2
2016 Low Latency and Resource-Aware Program Composition for Large-Scale Data Analysis
abstract
The importance of large-scale data analysis has shown a recent increase in a wide variety of areas, such as natural language processing, sensor data analysis, and scientific computing. Such an analysis application typically reuses existing programs as components and is often required to continuously process new data with low latency while processing large-scale data on distributed computation nodes. However, existing frameworks for combining programs into a parallel data analysis pipeline (e.g., workflow) are plagued by the following issues: (1) Most frameworks are oriented toward high-throughput batch processing, which leads to high latency. (2) A specific language is often imposed for the composition and/or such a specific structure as a simple unidirectional dataflow among constituting tasks. (3) A program used as a component often takes a long time to start up due to the heavy load at initialization, which is referred to as the startup overhead. Our solution to these problems is a remote procedure call (RPC)-based composition, which is achieved by our middleware Rapid Service Connector (RaSC). RaSC can easily wrap an ordinary program and make it accessible as an RPC service, called a RaSC service. Using such component programs as RaSC services enables us to integrate them into one program with low latency without being restricted to a specific workflow language or dataflow structure. In addition, a RaSC service masks the startup overhead of a component program by keeping the processes of the component program alive across RPC requests. We also proposed architecture that automatically manages the number of processes to maximize the throughput. The experimental results showed that our approach excels in overall throughput as well as latency, despite its RPC overhead. We also showed that our approach can adapt to runtime changes in the throughput requirements.
Masahiro Tanaka, Kenjiro Taura, Kentaro Torisawa
CCGrid2
2016 Tapas: An Implicitly Parallel Programming Framework for Hierarchical N-Body Algorithms
abstract
Tapas is our new C++ programming framework for hierarchical algorithms such as N-body, on large scale heterogeneous supercomputers. Although N-body and their variants are widely used in scientific applications, their correct implementations are often difficult on such modern machines, as the algorithms are irregular, complex, and involve explicit task parallel programming over distributed nodes. Encapsulating the complexities in a library or a framework has been challenging due to irregular data access over massively distributed memory. Tapas solves this by converting the users clean implicit-style parallel program into an inspector-executor style code on heterogeneous multi-core, multi-node environment solely by the use of C++ template metaprogramming. A prototype implementation of the Fast Multipole Method on Tapas demonstrates a comparable performance and scaling as ExaFMM, the fastest hand-tuned implementation of FMM, as well as efficient usage of hundreds of GPUs. Specifically, the serial performance is 95% of ExaFMM, whereas the distributed-memory strong-scaling evaluation using up to 1500 CPU cores demonstrates 64% to 81% of the ExaFMM performance. The multi-GPU version of the Tapas-based FMM achieves a 5.15x speedup when executed on 100 nodes of TSUBAME2.5 with 300 GPUs.
Keisuke Fukuda, Motohiko Matsuda, Naoya Maruyama, Rio Yokota, Kenjiro Taura, Satoshi Matsuoka
ICPADS5
2016 Fragmented BWT: An Extended BWT for Full-Text Indexing
Masaru Ito, Hiroshi Inoue, Kenjiro Taura
SPIRE3
2015 Uni-Address Threads: Scalable Thread Management for RDMA-Based Work Stealing
abstract
Task-parallel systems have been widely used to parallelize programs. They provide automatic load balancing and programmers can easily parallelize sequential programs, including irregular ones, without considering task placement to physical processors.
Shigeki Akiyama, Kenjiro Taura
HPDC2
2015 SIMD- and Cache-Friendly Algorithm for Sorting an Array of Structures
abstract
This paper describes our new algorithm for sorting an array of structures by efficiently exploiting the SIMD instructions and cache memory of today's processors. Recently, multiway mergesort implemented with SIMD instructions has been used as a high-performance in-memory sorting algorithm for sorting integer values. For sorting an array of structures with SIMD instructions, a frequently used approach is to first pack the key and index for each record into an integer value, sort the key-index pairs using SIMD instructions, then rearrange the records based on the sorted key-index pairs. This approach can efficiently exploit SIMD instructions because it sorts the key-index pairs while packed into integer values; hence, it can use existing high-performance sorting implementations of the SIMD-based multiway mergesort for integers. However, this approach has frequent cache misses in the final rearranging phase due to its random and scattered memory accesses so that this phase limits both single-thread performance and scalability with multiple cores. Our approach is also based on multiway mergesort, but it can avoid costly random accesses for rearranging the records while still efficiently exploiting the SIMD instructions. Our results showed that our approach exhibited up to 2.1x better single-thread performance than the key-index approach implemented with SIMD instructions when sorting 512M 16-byte records on one core. Our approach also yielded better performance when we used multiple cores. Compared to an optimized radix sort, our vectorized multiway mergesort achieved better performance when the each record is large. Our vectorized multiway mergesort also yielded higher scalability with multiple cores than the radix sort.
Hiroshi Inoue, Kenjiro Taura
Proc. VLDB Endow.2
2014 Scalable analysis of multicore data reuse and sharing
abstract
The performance and energy efficiency of multicore systems are increasingly dominated by the costs of communication. As hardware parallelism grows, developers require more powerful tools to assess the data sharing and reuse properties of their algorithms. The reuse distance is an effective metric to study the temporal locality of programs and model private and shared caches. But the application of this method is challenging. First, generating memory traces is very expensive in storage and very intrusive on execution, possibly distorting the parallel schedule. And second, the algorithm is computationally very expensive, limiting the length, memory size and parallelism of analyzable programs.
Miquel Pericàs, Kenjiro Taura, Satoshi Matsuoka
ICS2
2014 Faster Set Intersection with SIMD instructions by Reducing Branch Mispredictions
abstract
Set intersection is one of the most important operations for many applications such as Web search engines or database management systems. This paper describes our new algorithm to efficiently find set intersections with sorted arrays on modern processors with SIMD instructions and high branch misprediction penalties. Our algorithm efficiently exploits SIMD instructions and can drastically reduce branch mispredictions. Our algorithm extends a merge-based algorithm by reading multiple elements, instead of just one element, from each of two input arrays and compares all of the pairs of elements from the two arrays to find the elements with the same values. The key insight for our improvement is that we can reduce the number of costly hard-to-predict conditional branches by advancing a pointer by more than one element at a time. Although this algorithm increases the total number of comparisons, we can execute these comparisons more efficiently using the SIMD instructions and gain the benefits of the reduced branch misprediction overhead. Our algorithm is suitable to replace existing standard library functions, such as std::set_intersection in C++, thus accelerating many applications, because the algorithm is simple and requires no preprocessing to generate additional data structures. We implemented our algorithm on Xeon and POWER7+. The experimental results show our algorithm outperforms the std::set_intersection implementation delivered with gcc by up to 5.2x using SIMD instructions and by up to 2.1x even without using SIMD instructions for 32-bit and 64-bit integer datasets. Our SIMD algorithm also outperformed an existing algorithm that can leverage SIMD instructions.
Hiroshi Inoue, Moriyoshi Ohara, Kenjiro Taura
Proc. VLDB Endow.3
2013 A selective checkpointing mechanism for query plans in a parallel database system
abstract
Most existing parallel database systems achieve fault tolerance by aborting unfinished queries upon a failure and restart the entire from the beginning. This is inefficient for long running queries of OLAP workloads. To solve this problem, this paper presents a selective checkpointing mechanism which materializes the outputs of some necessary operators, enabling to resume queries from middle of the execution upon failures. Each query is represented by a DAG of relational operators in which data are typically pipelined between operators. The goal of the mechanism is to find a set of operators whose outputs are worth being checkpointed to minimize the expected runtime of the whole query. It firstly provides a cost model to estimate the expected runtime of a whole query plan under a given failure probability for each operator. Then a divide-and-conquer algorithm is proposed to find a close-to-optimal solution to the problem. The algorithm divides the query plan into subplans with smaller search spaces. For a given query plan with n operators, the algorithm runs in O(n) time. The mechanism is implemented in a shared-nothing parallel database system called ParaLite which provides a coordination layer to glue many SQLite instances together, and parallelizes SQL queries across them. The experimental results indicate that different fault-tolerant strategies affect the overall runtimes of queries. Our selective checkpointing mechanism can choose reasonable operators to be checkpointed and outperforms other fault-tolerant strategies. In addition, the divide-and-conquer algorithm taken by our mechanism has a smaller overhead than brute-force approach while keeping a similar effectiveness.
Kenjiro Taura
IEEE BigData2
2013 Parallel and memory-efficient Burrows-Wheeler transform
abstract
In large scale applications of genome analysis, compressed indexes are very useful to avoid the memory usage imposed by suffix arrays. FM-Index is one of the most important such compressed indexes. While it is very memory efficient, its construction requires much larger working memory during construction than the final index itself. Specifically, the Burrows-Wheeler transform (BWT), which is computed on its way to constructing FM-Index, requires a suffix array, which is larger than BWT and FM-Index. To address this, there have been many efforts to compute the Burrows-Wheeler in a smaller space, one of the fastest of which achieves a linear time complexity. Since it is unlikely that any sequential algorithm achieves a sublinear time complexity, the next logical step is a parallel algorithm that achieves a sublinear critical path with a space requirement smaller than suffix arrays. In this paper, we propose such an algorithm. It is based on a divide-and-conquer algorithm, which can be efficiently executed with task parallel models. We evaluated our algorithm on human genome data and obtained almost the same speed with BWT-IS, which is the known fastest algorithm.
Shinya Hayashi, Kenjiro Taura
IEEE BigData2
2013 Design and implementation of GXP make - A workflow system based on make
Kenjiro Taura, Takuya Matsuzaki, Makoto Miwa, Yoshikazu Kamoshida, Daisaku Yokoyama, Nan Dun, Takeshi Shibata, Choi Sung Jun, Jun'ichi Tsujii
Future Gener. Comput. Syst.1
2012 ParaLite: Supporting Collective Queries in Database System to Parallelize User-Defined Executable
abstract
This paper proposes extensions to parallel database systems called collective queries and User-Defined eXecutables (UDX). A collective query is an SQL query whose results are distributed to multiple clients and then processed by them in parallel, using arbitrary external programs (user-defined executables). The intended applications are data intensive work-flows, typically built out of various independently developed executables and scripts. Collective queries facilitate description of such workflows by making data parallel execution of external programs on big data easy and streamlined. It also provides the workflow developers with a familiar and powerful language SQL, for flexible data filtering and stereotypical data processing tasks. We implement this concept in a system "ParaLite", a parallel database system based on a popular lightweight database SQ Lite. It equips with data transfer optimization algorithms that distribute query results to multiple clients, taking both communication cost and compute loads into account. We verified the correctness and performance of Para Lite and the experimental results show that Para Lite has good performance on SQL processing and achieves good scalability for the parallelization of UDX.
Kenjiro Taura
CCGRID2
2010 Fine-Grained Profiling for Data-Intensive Workflows
abstract
Profiling is an effective dynamic analysis approach to investigate complex applications. ParaTrac is a user-level profiler using file system and process tracing techniques for data-intensive workflow applications. In two respects ParaTrac helps users refine the orchestration of workflows. First, the profiles of I/O characteristics enable users to quickly identify bottlenecks of underlying I/O subsystems. Second, ParaTrac can exploit fine-grained data-processes interactions in workflow execution to help users understand, characterize, and manage realistic data-intensive workflows. Experiments on thoroughly profiling Montage workflow demonstrate that ParaTrac is scalable to tracing events of thousands of processes and effective in guiding fine-grained workflow scheduling or workflow management systems improvements.
Nan Dun, Kenjiro Taura, Akinori Yonezawa
CCGRID2
2010 File-Access Characteristics of Data-Intensive Workflow Applications
abstract
This paper studies five real-world data intensive workflow applications in the fields of natural language processing, astronomy image analysis, and web data analysis. Data intensive workflows are increasingly becoming important applications for cluster and Grid environments. They open new challenges to various components of workflow execution environments including job dispatchers, schedulers, file systems, and file staging tools. Their impacts on real workloads are largely unknown. Under- standing characteristics of real-world workflow applications is a required step to promote research in this area. To this end, we analyse real-world workflow applications focusing on their file access patterns and summarize their implications to schedulers and file system/staging designs.
Takeshi Shibata, SungJun Choi, Kenjiro Taura
CCGRID3
2010 Design and Implementation of GXP Make - A Workflow System Based on Make
abstract
This paper describes the rational behind designing workflow systems based on the Unix make by showing a number of idioms useful for workflows comprising many tasks. It also demonstrates a specific design and implementation of such a workflow system called GXP make. GXP make supports all the features of GNU make and extends its platforms from single node systems to clusters, clouds, supercomputers, and distributed systems. Interestingly, it is achieved by a very small code base that does not modify GNU make implementation at all. While being not ideal for performance, it achieved a useful performance and scalability of dispatching one million tasks in approximately 16,000 seconds (60 tasks per second, including dependence analysis) on an 8 core Intel Nehalem node. For real applications, recognition and classification of protein-protein interactions from biomedical texts on a supercomputer with more than 8,000 cores are described.
Kenjiro Taura, Takuya Matsuzaki, Makoto Miwa, Yoshikazu Kamoshida, Daisaku Yokoyama, Nan Dun, Takeshi Shibata, Choi Sung Jun, Jun'ichi Tsujii
eScience1
2010 ParaTrac: a fine-grained profiler for data-intensive workflows
abstract
The realistic characteristics of data-intensive workflows are critical to optimal workflow orchestration and profiling is an effective approach to investigate the behaviors of such complex applications. ParaTrac is a fine-grained profiler for data-intensive workflows by using user-level file system and process tracing techniques. First, ParaTrac enables users to quickly understand the I/O characteristics of from entire application to specific processes or files by examining low-level I/O profiles. Second, ParaTrac automatically exploits fine-grained data-processes interactions in workflow to help users intuitively and quantitatively investigate realistic execution of data-intensive workflows. Experiments on thoroughly profiling Montage workflow demonstrate both the scalability and effectiveness of ParaTrac. The overhead of tracing thousands of processes is around 16%. We use low-level I/O profiles and informative workflow DAGs to illustrate the vantage of fine-grained profiling by helping users comprehensively understand the application behaviors and refine the scheduling for complex workflows. Our study also suggests that current workflow management systems may use fine-grained profiles to provide more flexible control for optimal workflow execution.
Nan Dun, Kenjiro Taura, Akinori Yonezawa
HPDC2
2010 A global address space framework for irregular applications
abstract
Practical parallel scientific applications with domain decompositions, such as finite element methods, require irregular domain decompositions of complicated-shaped objects. However, existing PGAS frameworks, such as Global Arrays and XcalableMP, have supported the productive description of exchanging ghost points only for regular domain decompositions. With these backgrounds, we propose, implement and evaluate Distributed Memory Interface (DMI), a global address space framework for irregular applications. DMI provides highly productive APIs called read-write-set for irregular domain decompositions and complicated orderings in practical scientific applications.
Kentaro Hara, Kenjiro Taura
HPDC2
2010 File-access patterns of data-intensive workflow applications and their implications to distributed filesystems
abstract
This paper studies five real-world data intensive workflow applications in the fields of natural language processing, astronomy image analysis, and web data analysis. Data intensive workflows are increasingly becoming important applications for cluster and Grid environments. They open new challenges to various components of workflow execution environments including job dispatchers, schedulers, file systems, and file staging tools. The keys to achieving high performance are efficient data sharing among executing hosts and locality-aware scheduling that reduces the amount of data transfer. While much work has been done on scheduling workflows, many of them use synthetic or random workload. As such, their impacts on real workloads are largely unknown. Understanding characteristics of real-world workflow applications is a required step to promote research in this area. To this end, we analyse real-world workflow applications focusing on their file access patterns and summarize their implications to schedulers and file system/staging designs.
Takeshi Shibata, SungJun Choi, Kenjiro Taura
HPDC3
2009 GMount: An Ad Hoc and Locality-Aware Distributed File System by Using SSH and FUSE
abstract
Developing and deploying distributed file system has been important for the Grid computing. By GMount, non-privileged users can instantaneously and effortlessly build a distributed file system on arbitrary machines that are reachable via SSH. It is scalable to hundreds of nodes in the wide-area Grid environments and adapts to NAT/Firewall. Unlike conventional distributed file systems, GMount can directly harness local file systems of each node without importing/exporting application data and utilize the network topology to make the metadata operations locality-aware. In this paper, we present the design and implementation of GMount by using two popular modules: SSH and FUSE. We demonstrate its viability and locality-aware metadata operation performance in a large scale Grid with over 320 nodes spreading across 12 clusters that are connected by heterogeneous wide-area links.
Nan Dun, Kenjiro Taura, Akinori Yonezawa
CCGRID2
2009 High performance wide-area overlay using deadlock-free routing
abstract
Overlay networks as the communication medium in parallel and distributed applications have gained prominence, especially in Grid environments. However, providing both throughput performance and reliable communication on overlays have been given little attention. The core of this problem is that intermediate nodes have limited buffer memory, while the forwarding throughput must yield Gbps. Yet, implementing a naive flow control can deadlock the overlay. Thus, high performance flow control on overlays is a critical concern in heterogeneous wide-area networks, where input/output link throughput can vary significantly. We propose an overlay scheme that couples TCP connections and fixed intermediate buffer memory while adapting deadlock-free routing for our overlay routing in heterogeneous wide-area networks. Our scheme eliminates memory overflows at forwarding nodes by fixed buffer memory and deadlocks via a deadlock-free routing algorithm that resolves adaptation challenges for heterogeneous wide-area networks. Our overlay construction and routing optimizations account for underlying network latency and bandwidth information. Simulation on 13 clusters (515 nodes) and evaluation on 7 clusters (170 nodes) show that our deadlock-free routing poses negligible overhead in comparison to deadlock-unaware routing, and comparably with direct communication. We further demonstrate that for certain collective communications, our overlay even out-performs direct communication by mitigating or completely avoiding network contention. We show this on systems ranging from a single-switch cluster with 36 nodes to a Grid environment with 4 clusters and 291 nodes.
Ken Hironaka, Hideo Saito 0002, Kenjiro Taura
HPDC3
2008 Scalable Data Gathering for Real-Time Monitoring Systems on Distributed Computing
abstract
Real-time monitoring is increasingly becoming important in various scenes of large scale, multi-site distributed/parallel computing, e.g, understanding behavior of systems, scheduling resources, and debugging applications. Dedicated networks on inter-site communications are rarely available for the monitoring purposes. Therefore, for realtime monitoring systems, reducing communication cost is important to handle a large number of nodes with limited network resources. We implemented a real-time Grid monitoring system called VGXP, with techniques for low cost data gathering. It tries to send only diffs to recent data, and adapts to the requested data freshness and tolerable errors to minimize required communication. We evaluate monitoring overheads of the proposed method on a distributed environment consisting of 8-sites with 500 nodes. In a realistic setting where the sampling interval is set to 0.5 seconds and the tolerable error to 2%, the CPU usage of the server to gather data from all nodes was 0.2% and the transfer rate was less than 5 kbps. The transfer rate did not exceed 50 kbps even if we gather a detailed per-process statistics.
Yoshikazu Kamoshida, Kenjiro Taura
CCGRID2
2008 A Stable Broadcast Algorithm
abstract
Distributing large data to many nodes, known as a broadcast or a multicast, is an important operation in parallel and distributed computing. Most previous broadcast algorithms explicitly or implicitly try to deliver data to all nodes in the same rate. This assumption is reasonable for homogeneous environments where all nodes have similar receiving capabilities. However, when nodes have various receiving capabilities, nodes with slow-receiving capabilities slow down the entire receiving bandwidth in these algorithms. In such settings, each node desires to receive data at its largest possible bandwidth and to start computation as soon as it receives the data. In this paper, we propose to say a broadcast is stable when the bandwidth to a node is never sacrificed by the presence of other, possibly slow, receiving nodes, and proposes the stability as a desired property of broadcast algorithms. In addition, we show a simple and efficient stable broadcast when the topology among nodes is a tree and each link has a symmetric bandwidth. This work improves upon previously proposed algorithms such as FPFR and Balanced Multicasting. For general graphs, it outperforms them when the network is heterogeneous and for trees, our algorithm is proved to be stable and optimal. In a real environment with 100 machines in 4 clusters, our scheme achieved 2.1 to 2.6 times aggregate bandwidth compared to the best result in the other algorithms. We also demonstrated the stability by adding a slow node to a broadcast. Some simulations also showed that our algorithm also performs well in many bandwidth distributions.
Kei Takahashi, Hideo Saito 0002, Takeshi Shibata, Kenjiro Taura
CCGRID4
2007 Locality-aware Connection Management and Rank Assignment forWide-area MPI
abstract
We propose a connection management scheme that limits the number of inter-cluster connections and forwards messages for processes that cannot communicate directly. We also propose a rank assignment scheme that finds rank-process mappings with low communication overhead by solving the quadratic assignment problem. Our proposed methods perform locality-aware communication optimizations, and do so without tedious manual configuration by obtaining latency and traffic information from a short profiling run of the environment and the application. Using these methods, we implemented a wide-area-enabled MPI library called MC-MPI, and evaluated its performance by running the NAS parallel benchmarks on 256 real nodes distributed across 4 clusters. MC-MPI was able to limit the number of process pairs that established connections to just 10% without suffering a performance penalty. Moreover, MC-MPI was able to find rank assignments that resulted in up to 160% better performance than locality-unaware assignments.
Hideo Saito 0002, Kenjiro Taura
CCGRID2
2007 A fast topology inference: a building block for network-aware parallel processing
abstract
Adapting to the network is the key to achieving high performance for communication-intensive applications, including scientific computing,data intensive computing, and multicast, especially in Grid environments. This paper investigates an approach of representing network as a tree of participating hosts and switches matching or approximating their physical topology, and describes a fast, non-intrusive, and portable algorithm for inferring such a topology. This representation and the proposed inference algorithm serves as a key to building network-aware applications in a portable manner. The algorithm is based solely on RTTs of small packets between end hosts; it does not rely on popular but not universally available protocols such as trace route and SNMP. Another benefit is that it can handle all layers of network uniformly without any a priori knowledge of cluster configurations. The required number of measurements is O(Nd) in certain idealizing assumptions made for the purpose of analysis, where N is the number of participating processes and d the diameter of the network, which is usually small in real networks. In our experimental environment, the inference algorithm built a topology of 64 hosts in a single cluster in 4 seconds and and that of 256 hosts across 4 clusters in 15 seconds. It is able to not only identify clusters within a Grid, but also to partially identify the Layer 2 topology within a cluster. This is important for optimizing bandwidth-limited operations such as broadcast. We built several network-aware applications upon the inference system, including efficient bandwidth measurements and long message broadcasts. The topology is used to schedule as many measurements as possible in parallel without competing on shared links. We were able to build a bandwidth map of 256 hosts across 4 clusters in 27 seconds.
Tatsuya Shirai, Hideo Saito 0002, Kenjiro Taura
HPDC3
2007 Locality-aware connection management and rank assignment for wide-area MPI
abstract
We propose locality and application-aware connection management and rank assignment schemes for wide area message passing systems. Using our connection management scheme on 256 processors in 4 clusters, the NAS Parallel Benchmarks ran with no performance penalty with as few as 10% of the connections. Moreover, our rank assignment scheme wasable to find assignments that resulted in up to 160% better performance than locality and application-unaware assignments.
Hideo Saito 0002, Kenjiro Taura
PPoPP2
2006 Monte Carlo Go Has a Way to Go
Haruhiro Yoshimoto, Kazuki Yoshizoe, Tomoyuki Kaneko, Akihiro Kishimoto, Kenjiro Taura
AAAI5
2004 High performance LU factorization for non-dedicated clusters
abstract
This paper describes an implementation of parallel LU factorization. The focus is to achieve high performance on non-dedicated clusters, where the number of available computing resources may be arbitrary and even dynamically changing. We accommodate joining/leaving processes by describing the algorithm in the Phoenix programming model. We achieve high performance in this setting by a combination of techniques including a latency tolerant communication and data partitioning that achieves both load balance and small communication volume for arbitrary and dynamically changing number of processors. We observed 130 GFlops with 128 processes on a 70-node dual 2.4GHz Xeon cluster, at matrix size = 46080. This performance is comparable to that of the High Performance Linpack (HPL). When cluster nodes are loaded by background processes, our implementation surpasses HPL.
Toshio Endo, Kenji Kaneda, Kenjiro Taura, Akinori Yonezawa
CCGRID3
2004 Routing and resource discovery in Phoenix Grid-enabled message passing library
abstract
We describe the design and implementation of a "Grid-enabled" message passing library, in the context of the Phoenix message passing model. It supports: (1) message routing between nodes not directly reachable due to firewalls and/or NAT; (2) resource discovery facilitating ease of configuration that allows nodes without static names; (e.g., DHCP nodes) to participate in computation without specific efforts; and (3) nodes dynamically joining/leaving computation at runtime. We argue that, in future Grid environments, all of the above functions, not just routing across firewalls, will become important issues of Grid-enabled message passing systems including MPI. Unlike solutions commonly proposed by previous work on a Grid-enabled MPI, our system runs a distributed resource discovery and routing table construction algorithm, rather than assuming all such pieces of information are available in a static configuration file or alike. Experimental results using 400 nodes in three LAN indicate that our algorithm is able to dynamically discover participating peers, connect them to each other and calculate a routing table. The elapsed time of our algorithm is only approximately twice as long as that of offline route calculation that just connects nodes based on a fully given configuration.
Kenji Kaneda, Kenjiro Taura, Akinori Yonezawa
CCGRID2
2003 Phoenix: a parallel programming model for accommodating dynamically joining/leaving resources
abstract
This paper proposes Phoenix, a programming model for writing parallel and distributed applications that accommodate dynamically joining/leaving compute resources. In the proposed model, nodes involved in an application see a large and fixed virtual node name space. They communicate via messages, whose destinations are specified by virtual node names, rather than names bound to a physical resource. We describe Phoenix API and show how it allows a transparent migration of application states, as well as dynamically joining/leaving nodes as its by-product. We also demonstrate through several application studies that Phoenix model is close enough to regular message passing, thus it is a general programming model that facilitates porting many parallel applications/algorithms to more dynamic environments. Experimental results indicate applications that have a small task migration cost can quickly take advantage of dynamically joining resources using Phoenix. Divide-and-conq! uer algorithms written in Phoenix achieved a good speedup with a large number of nodes across multiple LANs (120 times speedup using 169 CPUs across three LANs). We believe Phoenix provides a useful programming abstraction and platform for emerging parallel applications that must be deployed across multiple LANs and/or shared clusters having dynamically varying resource conditions.
Kenjiro Taura, Kenji Kaneda, Toshio Endo, Akinori Yonezawa
PPoPP1
2003 Virtual private grid: a command shell for utilizing hundreds of machines efficiently
Kenji Kaneda, Kenjiro Taura, Akinori Yonezawa
Future Gener. Comput. Syst.2
2002 Virtual Private Grid: A Command Shell for Utilizing Hundreds of Machines Efficiently
abstract
We design and implement Virtual Private Grid (VPG), a shell that can easily and securely utilize a large number of machines distributed over multiple administrative domains. Today, many people have an access to a large number of machines across multiple subnets or geographically distributed places. These machines are managed by different administrators, and for the sake of security and administration cost, they impose various restrictions on their use. Methods to work around these restrictions are found on a case-by-case basis and require human intervention. There-fore, it increases the user's cost to utilize remote machines significantly, and consequently decreases the utilization of computational resources. VPG works around these restrictions automatically and can easily utilize a large number of machines in multiple administrative domains. We run VPG on approximately 100 nodes (270 CPUs). Experimental results show that VPG utilizes remote machines more efficiently than other job submission tools.
Kenji Kaneda, Kenjiro Taura, Akinori Yonezawa
CCGRID2
2001 Predicting Scalability of Parallel Garbage Collectors on Shared Memory Multiprocessors
abstract
This paper describes a performance prediction model of parallel mark-sweep garbage collectors (GC) on shared memory multiprocessors. The prediction model takes the heap snapshot and memory access cost parameters (latency and occupancy) as inputs, and outputs performance of the parallel marking on any given number of processors. It takes several factors Mat affects performance into account: cache misses costs, memory access contention, and increase of misses by parallelization We evaluate this model by comparing the predicted GC performance and measured performance on two architecturally different shared memory machines: Ultra Enterprise 10000 (crossbar connected SMP) and Origin 2000 (hypercube connected DSM). Our model accurately predicts qualitatively different speedups on the two machines that occurred in one application, which turn out to be due to contentions on a memory node. Lit addition to performance analysis, applications of the proposed model include adaptive GC algorithm to achieve optimal performance based on the prediction. This paper shows the effect of automatic regulation of GC parallelism.
Toshio Endo, Kenjiro Taura, Akinori Yonezawa
IPDPS2
2000 The MicroGrid: a Scientific Tool for Modeling Computational Grids
abstract
The complexity and dynamic nature of the Internet (and the emerging Computational Grid) demand that middleware and applications adapt to the changes in configuration and availability of resources. However, to the best of our knowledge there are no simulation tools which support systematic exploration of dynamic Grid software (or Grid resource) behavior. We describe our vision and initial efforts to build tools to meet these needs. Our MicroGrid simulation tools enable Globus applications to be run in arbitrary virtual grid resource environments, enabling broad experimentation. We describe the design of these tools, and their validation on micro- benchmarks, the NA parallel benchmarks, and an entire Grid application. These validation experiments show that the MicroGrid can match actual experiments within a few percent (2% to 4%).
Hyo Jung Song, Xianan Liu, Dennis Jakobsen, Ranjita Bhagwan, Xingbin Zhang, Kenjiro Taura, Andrew A. Chien
SC6
2000 Extending Java virtual machine with integer-reference conversion
abstract
Java virtual machine (JVM) is an architecture-independent code execution environment. It has recently been used not only for the Java language but also for other languages such as Scheme and ML. On JVM, however, all values are statically typed as either immediate or reference, and types are checked before the execution of a program to prove that invalid memory access will never occur. This property sometimes makes implementation of other languages on JVM inefficient. In particular, implementation of a dynamically typed language is very inefficient because all possible values including frequently used ones such as integers must be represented by instances of a class. In this paper, we introduce a new type into JVM, which is a supertype of reference types and a tagged integer type. This allows a more efficient implementation of dynamically typed language on JVM. It does not require any new instruction, maintains binary-compatibility of existing bytecode, and retains the safety of the original JVM. We modified an existing Scheme system running on JVM to exploit this extension and got a factor of 20 speedup for simple integer functions. Our extension imposes little performance penalty on existing JVM code generated from Java; we observed essentially no penalty for Spec JVM benchmarks. Copyright © 2000 John Wiley & Sons, Ltd.
Yutaka Oiwa, Kenjiro Taura, Akinori Yonezawa
Concurr. Pract. Exp.2
1999 StackThreads/MP: Integrating Futures into Calling Standards
abstract
An implementation scheme of fine-grain multithreading that needs no changes to current calling standards for sequential languages and modest extensions to sequential compilers is described. Like previous similar systems, it performs an asynchronous call as if it were an ordinary procedure call, and detaches the callee from the caller when the callee suspends or either of them migrates to another processor. Unlike previous similar systems, it detaches and connects arbitrary frames generated by off-the-shelf sequential compilers obeying calling standards. As a consequence, it requires neither a frontend preprocessor nor a native code generator that has a builtin notion of parallelism. The system practically works with unmodified GNU C compiler (GCC). Desirable extensions to sequential compilers for guaranteeing portability and correctness of the scheme are clarified and claimed modest. Experiments indicate that sequential performance is not sacrificed for practical applications and both sequential and parallel performance are comparable to Cilk[8], whose current implementation requires a fairly sophisticated preprocessor to C. These results show that efficient asynchronous calls (a.k.a. future calls) can be integrated into current calling standard with a very small impact both on sequential performance and compiler engineering.
Kenjiro Taura, Kunio Tabata, Akinori Yonezawa
PPoPP1
1997 An Efficient Compilation Framework for Languages Based on a Concurrent Process Calculus
Yoshihiro Oyama, Kenjiro Taura, Akinori Yonezawa
Euro-Par2
1997 Fine-grain Multithreading with Minimal Compiler Support - A Cost Effective Approach to Implementing Efficient Multithreading Languages
abstract
It is difficult to map the execution model of multithreading languages (languages which support fine-grain dynamic thread creation) onto the single stack execution model of C. Consequently, previous work on efficient multithreading uses elaborate frame formats and allocation strategy, with compilers customized for them. This paper presents an alternative cost-effective implementation strategy for multithreading languages which can maximally exploit current sequential C compilers. We identify a set of primitives whereby efficient dynamic thread creation and switch can be achieved and clarify implementation issues and solutions which work under the stack frame layout and calling conventions of current C compilers. The primitives are implemented as a C library and named StackThreads. In StackThreads, a thread creation is done just by a C procedure call, maximizing thread creation performance. When a procedure suspends an execution, the context of the procedure, which is roughly a stack frame of the procedure, is saved into heap and resumed later. With StackThreads, the compiler writer can straightforwardly translate sequential constructs of the source language into corresponding C statements or expressions, while using StackThreads primitives as a blackbox mechanism which switches execution between C procedures.
Kenjiro Taura, Akinori Yonezawa
PLDI1
1997 An Effective Garbage Collection Strategy for Parallel Programming Languages on Large Scale Distributed-Memory Machines
abstract
This paper describes the design and implementation of a garbage collection scheme on large-scale distributed-memory computers and reports various experimental results. The collector is based on the conservative GC library by Boehm & Weiser. Each processor traces local pointers using the GC library while traversing remote pointers by exchanging "mark messages" between processors. It exhibits a promising performance---in the most space-intensive settings we tested, the total collection overhead ranges from 5% up to 15% of the application running time (excluding idle time). We not only examine basic performance figures such as the total overhead or latency of a global collection, but also demonstrate how local collection scheduling strategies affect application performance. In our collector, a local collection is scheduled either independently or synchronously. Experimental results show that the benefit of independent local collections has been overstated in the literature. Independent local collections slowed down application performance to 40%, by increasing the average communication latency. Synchronized local collections exhibit much more robust performance characteristics than independent local collections and the overhead for global synchronization is not significant. Furthermore, we show that an adaptive collection scheduler can select the appropriate local collection strategy based on the application's behavior. The collector has been used in a concurrent object-oriented language ABCL/f and the performance is measured on a large-scale parallel computer (256 processors) using four non-trivial applications written in ABCL/f.
Kenjiro Taura, Akinori Yonezawa
PPoPP1
1997 A Scalable Mark-Sweep Garbage Collector on Large-Scale Shared-Memory Machines
abstract
This work describes implementation of a mark-sweep garbage collector (GC) for shared-memory machines and reports its performance. It is a simple ''parallel'' collector in which all processors cooperatively traverse objects in the global shared heap. The collector stops the application program during a collection and assumes a uniform access cost to all locations in the shared heap. Implementation is based on the Boehm-Demers-Weiser conservative GC (Boehm GC). Experiments have been done on Ultra Enterprise 10000 (Ultra Sparc processor 250 MHz, 64 processors). We wrote two applications, BH (an N-body problem solver) and CKY (a context free grammar parser) in a parallel extension to C++.Through the experiments, We observe that load balancing is the key to achieving scalability. A naive collector without load redistribution hardly exhibits speed-up (at most fourfold speed-up on 64 processors). Performance can be improved by dynamic load balancing, which exchanges objects to be scanned between processors, but we still observe that straightforward implementation severely limits performance. First, large objects become a source of significant load imbalance, because the unit of load redistribution is a single object. Performance is improved by splitting a large object into small pieces before pushing it onto the mark stack. Next, processors spend a significant amount of time uselessly because of serializing method for termination detection using a shared counter. This problem suddenly appeared on more than 32 processors. By implementing non-serializing method for termination detection, the idle time is eliminated and performance is improved. With all these careful implementation, we achieved average speed-up of 28.0 in BH and 28.6 in CKY on 64 processors.
Toshio Endo, Kenjiro Taura, Akinori Yonezawa
SC2
1996 Visualization of RNA secondary structures using highly parallel computers
abstract
Results of RNA secondary structure prediction algorithm are usually given as a set of hydrogen bonds between bases. However, we cannot know the precise structure of an RNA molecule by only knowing which bases form hydrogen bonds. One way to understand the structure of an RNA molecule is to visualize it using a planar graph so that we can easily know the geometric relations among the substructures such as stacking regions and loops. To do this, we consider bases to be particles on a plane and introduce a repulsive force and an attractive force among these particles and determine their positions according to these forces. A naive algorithm requires O(N2) time but we can reduce it to O(NlogN) with an approximation algorithm which is often used in the area of N-body simulation. Our program is written in parallel object-oriented language 'Schematic' which is recently developed. Efficiency of our implementation on a parallel computer and results of visualization of secondary structures are presented using cadang-cadang coconut viroid as an example.
Akihiro Nakaya, Kenjiro Taura, Kenji Yamamoto, Akinori Yonezawa
Comput. Appl. Biosci.2
1993 Highly Efficient and Encapsulated Re-use of Synchronization Code in Concurrent Object-Oriented Languages
abstract
Re-use of synchronization code in concurrent OOlanguages has been considered difficult due to inheritance anomaly, which we minimize with our new proposal. Designed with high practicality in mind, we propose language primitives (plus their implementation) with the following characteristics: (1) it allows multiple synchronization schemes---the language schemes for programming synchronization---to coexist and be integrated, (2) re-use of synchronization code is done similarly to sequential OO-languages for user familiarity, (3) it offers high degree of encapsulation---even synchronization schemes could be encapsulated in superclasses in many cases, and (4) it can be efficiently implemented on conventional MPPs. We demonstrate the effectiveness of our proposal with solutions to the example inheritance anomaly cases from [16]. We also give an overview of the implementation architecture, along with preliminary benchmarks. The proposed language primitives are being incorporated into our AB...
Satoshi Matsuoka, Kenjiro Taura, Akinori Yonezawa
OOPSLA2
1993 An Efficient Implementation Scheme of Concurrent Object-Oriented Languages on Stock Multicomputers
abstract
Several novel techniques for efficient implementtion of concurrent object-oriented languages on general purpose, stock multicomputers are presented. These techniques have been developed in implementing our concurrent object-oriented language ABCL on a Fujitsu Laboratory's experimental multicomputer AP1000 consisting of 512 SPARC chips. The propsed intra-node scheduling mechanism reduces the cost of local message passing. The cost of intra-node asynchronous message passing is about 20 SPARC instructions in the bst case, including locality checking, dynamic method lookup, and scheduling. The minimum latency of asynchronous internode message passing is about 9μs, or about 120 instructions, employing the self-dispatching mechanism independently proposed by Eicken et al. A large scale benchmark which involves 9,000,000 message passings shows 440 times speedup on the 512 nodes system compared to the sequential version of the same algorithm. We rely on simple hardware support for message passing and use no specialized architectural supports for object-oriented computing. Thus, we are able to enjoy the benefits of future progress in standard processor technology. Our result shows that concurrent object-oriented languages can be implemented efficiently on conventional multicomputers.
Kenjiro Taura, Satoshi Matsuoka, Akinori Yonezawa
PPoPP1