Tekin Bicer

dblp:30/9083 · also Tekin Biçer · DBLP profile ↗
← Back
23ranked-venue papers
7as first author
11since 2021 · last 2026
0000-0002-8428-5159ORCID · verified

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

Systems, architecture and hardware · 19 · 6 first-author · 8 since 2021Software engineering, systems software and programming languages · 2 · 1 first-author · 1 since 2021Graphics, computer vision, multimedia, augmented reality and games · 2 · 2 since 2021Applied, interdisciplinary, general and emerging computing · 2 · 1 first-author · 1 since 2021
YearPublicationVenuePosition
2026 Towards Transparent Checkpointing with AI-driven Code Generation
abstract
Adding reliable checkpoint/restart support to an MPI scientific application is a time-consuming expert effort that requires deep knowledge of both the application and resilience. We ask whether a frontier large language model can perform this work end-to-end without human intervention. We assemble a benchmark suite of MPI applications spanning diverse domains and computation patterns, and drive an iterative code-generation loop for each application using Anthropic’s Claude Opus 4.7 invoked through the OpenCode CLI. Across six scientific applications, the LLM generates working checkpoint/restart code in 50 minutes on average while consuming 3.4 M tokens per application. The generated code adds negligible overhead during normal failure-free execution on five of six applications and recovers from injected process failures with efficiency comparable to human-engineered checkpoint/restart implementations. These results suggest that automated end-to-end LLM-driven resilience engineering is technically viable today for a meaningful fraction of HPC applications.
Hai Nguyen 0005, Tekin Bicer, Kyle Chard, Ian T. Foster, Bogdan Nicolae
HPDC2
2026 StreamGuard: Low-Overhead Resilience for Real-time HPC Data Streams
abstract
Real-time scientific workflows operate on continuous data streams and must produce timely, high-quality results despite executing on complex, failure-prone infrastructure. Hardware faults, network disruptions, and performance anomalies caused by resource contention or system heterogeneity can severely degrade performance and violate real-time constraints. We focus on strengthening the resilience of the producer–consumer streaming pattern, a fundamental building block of scientific streaming workflows. We present two complementary techniques: (i) a dynamic, asynchronous, non-blocking checkpointing mechanism that preserves progress without interrupting computation, and (ii) a progress-aware load redistribution strategy that detects slow workers and proactively rebalances tasks. Together, these mechanisms maintain forward progress and balanced execution even in highly error-prone environments. Experimental results show that our approach reduces the impact of failures and performance anomalies by up to 6 ×, while introducing less than 1% overhead in failure-free execution.
Hai Nguyen 0005, Bogdan Nicolae, Tekin Bicer, Amal Gueroudji, Matthieu Dorier, Kyle Chard, Ian T. Foster
ICS3
2025 Sinogram Inpainting with Physics-Guided Latent Diffusion Model for Synchrotron Light Sources
abstract
X-ray Computed Tomography (XCT) is widely used for imaging materials at microscopic or sub-microscopic lengths in synchrotrons. During experiments, often limited XCT data is collected, compromising the reconstruction quality. Here, we develop a foundation model for XCT with inpainting as downstream task. Our model integrates a Generative AI-based Latent Diffusion Model (LDM) with physics domain knowledge. We incorporate additional loss functions into the autoencoder of the LDM to accurately capture the physical properties of the XCT data acquisition. This loss function and a pre-training step improve the autoencoder’s performance. Lastly, we introduce a novel image blending method to combine the LDM’s output with the original, extremely sparse sinogram data. On real-world test dataset, with 80% randomly masked data, we demonstrate mean SSIM of 0.8699 and 0.7290 for sinogram and reconstructed object respectively. Additionally, we obtain mean PSNR of 36.13 dB and 32.26 dB for sinogram and reconstructed object respectively.
Srutarshi Banerjee, Jiaze E, Bin Ren 0002, Tekin Bicer
ICIP4
2025 mLR: Scalable Laminography Reconstruction based on Memoization
abstract
ADMM-FFT is an iterative method with high reconstruction accuracy for laminography but suffers from excessive computation time and large memory consumption. We introduce mLR, which employs memoization to replace the time-consuming Fast Fourier Transform (FFT) operations based on an unique observation that similar FFT operations appear in iterations of ADMM-FFT. We introduce a series of techniques to make the application of memoization to ADMM-FFT performance-beneficial and scalable. We also introduce variable offloading to save CPU memory and scale ADMM-FFT across GPUs within and across nodes. Using mLR, we are able to scale ADMM-FFT on an input problem of 2K × 2K × 2K, which is the largest input problem laminography reconstruction has ever worked on with the ADMM-FFT solution on limited memory; mLR brings 52.8% performance improvement on average (up to 65.4%), compared to the original ADMM-FFT.
Bin Ma 0025, Viktor Nikitin, Xi Wang 0027, Tekin Bicer, Dong Li 0001
SC4
2025 Efficient distributed continual learning for steering experiments in real-time
Thomas Bouvier, Bogdan Nicolae, Alexandru Costan, Tekin Bicer, Ian T. Foster, Gabriel Antoniu
Future Gener. Comput. Syst.4
2024 Diaspora: Resilience-Enabling Services for Real-Time Distributed Workflows
abstract
The need for real-time processing to enable automated decision making and experimental steering has driven a shift from high-performance computing workflows on a centralized system to a distributed approach that integrates remote data sources, edge devices, and diverse compute facilities. Under this paradigm, data can be processed close to the source where it is generated, thus reducing latency and bandwidth usage. System resilience is thus a key challenge, requiring distributed workflows to survive component failures and to meet stringent quality-of-service requirements, which results in the need to mitigate anomalies such as congestion and low availability of resources. To address these challenges, we propose Diaspora, a unified resilience framework that is inspired by event-driven communication patterns used in public clouds. Specifically, we propose an event fabric that extends across sites, facilities, and computations to provide timely, reliable, and accurate information about data, application, and resource status. On top of the event fabric, we build resilience-enabling services that combine QoS-aware data streaming, resilient data views, resilient compute and data resources, and anomaly detection and prediction, all of which collectively enhance workflow resilience for these scientific cases.
Bogdan Nicolae, Justin M. Wozniak, Tekin Bicer, Hai Nguyen 0005, Haochen Pan, Amal Gueroudji, Maxime Gonthier, Valérie Hayot-Sasson, Eliu A. Huerta, Kyle Chard, Ryan Chard, Matthieu Dorier, Nageswara S. V. Rao, Anees Al-Najjar, Alessandra Corsi, Ian T. Foster
e-Science3
2024 CommBench: Micro-Benchmarking Hierarchical Networks with Multi-GPU, Multi-NIC Nodes
abstract
Modern high-performance computing systems have multiple GPUs and network interface cards (NICs) per node. The resulting network architectures have multilevel hierarchies of subnetworks with different interconnect and software technologies. These systems offer multiple vendor-provided communication capabilities and library implementations (IPC, MPI, NCCL, RCCL, OneCCL) with APIs providing varying levels of performance across the different levels. Understanding this performance is currently difficult because of the wide range of architectures and programming models (CUDA, HIP, OneAPI).
Mert Hidayetoglu, Simon Garcia de Gonzalo, Elliott Slaughter, Yu Li 0041, Christopher Zimmer 0001, Tekin Bicer, Bin Ren 0002, William Gropp, Wen-Mei W. Hwu, Alex Aiken
ICS6
2022 SciStream: Architecture and Toolkit for Data Streaming between Federated Science Instruments
abstract
Modern scientific instruments, such as detectors at synchrotron light sources, generate data at such high rates that online processing is needed for data reduction, feature detection, experiment steering, and other purposes. The same high data rates also demand memory-to-memory streaming from instrument to remote computer, because local computational capacity is limited and data transmissions that engage the file system introduce unacceptable latencies. But efficient and secure memory-to-memory data streaming is challenging to realize in practice, because of a lack of direct external network connectivity for scientific instruments and because of authentication and security requirements. In response, we propose here SciStream, a middlebox-based architecture with control protocols to enable efficient and secure memory-to-memory data streaming between producers and consumers that lack direct network connectivity. We describe the protocols that SciStream uses to establish authenticated and transparent connections between producers and consumers, and we discuss the experiments that we have conducted to evaluate alternative implementation approaches for key SciStream components. Experiments on the Chameleon cloud show that SciStream improves the throughput of a streaming pipeline by an order of magnitude compared with state-of-the-art data transfer methods and adds only ~4μsec latency compared with an ideal scenario in which producers and consumers have direct external connectivity.
Joaquin Chung 0001, Wojciech Zacherek, A. J. Wisniewski, Zhengchun Liu, Tekin Bicer, Rajkumar Kettimuthu, Ian T. Foster
HPDC5
2022 MemXCT: Design, Optimization, Scaling, and Reproducibility of X-Ray Tomography Imaging
abstract
This work extends our previous research entitled “MemXCT: Memory-centric X-ray CT Reconstruction with Massive Parallelization” that was originally published at SC19 conference (Hidayetoğluet al., 2019) with reproducibility of the computational imaging performance. X-ray computed tomography (XCT) is regularly used at synchrotron light sources to study the internal morphology of materials at high resolution. However, experimental constraints, such as radiation sensitivity, can result in noisy or undersampled measurements. Further, depending on the resolution, sample size and data acquisition rates, the resulting noisy dataset can be in the order of terabytes. Advanced iterative reconstruction techniques can produce high-quality images from noisy measurements, but their computational requirements have made their use an exception rather than the rule. We propose a novel memory-centric approach that avoids redundant computations at the expense of additional memory complexity. We develop a memory-centric iterative reconstruction system, MemXCT, that uses an optimized SpMV implementation with two-level pseudo-Hilbert ordering and multi-stage input buffering. We evaluate MemXCT on various supercomputer architectures involving KNL and GPU. MemXCT can reconstruct a large (11K×11K) mouse brain tomogram in 10 seconds using 4096 KNL nodes (256K cores). The results presented in our original article at the SC19 were based on large-scale supercomputing resources. The MemXCT application was selected for the Student Cluster Competition (SCC) Reproducibility Challenge and evaluated on a variety of cloud computing resources by universities around the world in the SC20 conference. We summarize the results of the top-ranked SCC Reproducibility Challenge teams and identify the most pertinent measures for ensuring the reproducibility of our experiments in this article.
Mert Hidayetoglu, Tekin Bicer, Simon Garcia de Gonzalo, Bin Ren 0002, Doga Gürsoy, Rajkumar Kettimuthu, Ian T. Foster, Wen-Mei W. Hwu
IEEE Trans. Parallel Distributed Syst.2
2021 3d Autoencoders For Feature Extraction In X-Ray Tomography
abstract
Real-time steering of time-resolved or in-situ X-ray tomography requires capturing changes in morphological descriptors in a sample (e.g., porosity, particle size, and crack width) during continuous data acquisition. Segmentation of 2D or 3D images followed by quantitative measurement is the conventional method for tracking changes in these descriptors with respect to a previous time-step or a 3D search in a volume. However, image segmentation is expensive. As a faster and unsupervised alternative, we propose a feature-extraction approach using a convolutional autoencoder, where the latent space of the encoder responds to relative changes in morphology without prior knowledge of the morphological descriptors. To test this approach in parametric experiments, we used a digital twin for micro-CT to generate realistic datasets of porous materials. We show that for some configurations of the autoencoder, its n-dimensional latent space encodes a sample's local porosity metrics while disregarding contrast information (relative proportion of absorption and phase contrast) determined by the imaging modality and not the sample. Through dimensionality reduction, the vector's response is visualized in 2D space to find clusters of data with similar porosity metrics. Our approach can extract features from grayscale tomographic data more than 4x faster than a segmentation + pore analysis workflow.
Aniket Tekawade, Zhengchun Liu, Peter Kenesei, Tekin Bicer, Francesco De Carlo, Rajkumar Kettimuthu, Ian T. Foster
ICIP4
2021 Topology-aware optimizations for multi-GPU ptychographic image reconstruction
abstract
Ptychography is an advanced high-resolution X-ray imaging technique that can generate extremely large datasets. Ptychographic reconstruction transforms reciprocal space experimental data to high-resolution 2D real-space images. GPUs have been used extensively to meet the computational requirements of the reconstruction. Generic multi-GPU reconstruction solutions use common communication topologies, such as P2P graph and ring, that are provided by MPI and NCCL libraries, to establish inter-GPU communications. However, these common topologies assume homogeneous physical links between GPUs, resulting in sub-optimal performance on heterogeneous configurations that are composed of both high- (e.g., NVLink) and low-speed (e.g., PCIe) interconnects. This mismatch between application-level communication topology and physical interconnection can cause data transfer congestion, inefficient memory access, and under-utilization of network resources. Here we present topology-aware designs and optimizations to address the aforementioned mismatch and boost end-to-end application performance. We introduce topology-aware data splitting, propose a novel communication topology, and incorporate asynchronous data movement and computation. We evaluate our design and optimizations using real and artificial datasets and compare its performance with that of the direct P2P and NCCL-based approaches. The results show that our optimizations always outperform the counterparts and achieve up to 5.13× and 1.63× communication and end-to-end application speedups, respectively.
Xiaodong Yu 0001, Tekin Bicer, Rajkumar Kettimuthu, Ian T. Foster
ICS2
2020 GPU-Based Static Data-Flow Analysis for Fast and Scalable Android App Vetting
abstract
Many popular vetting tools for Android applications use static code analysis techniques. In particular, Interprocedural Data-Flow Graph (IDFG) construction is the computation at the core of Android static data-flow analysis and consumes most of the analysis time. Many analysis tools use a worklist algorithm, an iterative fixed-point approach, to construct the IDFG. In this paper, we observe that a straightforward GPU parallelization of the worklist algorithm leads to significant underutilization of the GPU resources. We identify four performance bottlenecks, namely, frequent dynamic memory allocations, high branch divergence, workload imbalance, and irregular memory access patterns. Accordingly, we propose GDroid, a GPU-based worklist algorithm implementation with multiple fine-grained optimizations tailored to common characteristics of Android applications. The optimizations considered are: matrix-based data structure, memory access-based node grouping, and worklist merging. Our experimental evaluation, performed on 1000 Android applications, shows that the proposed optimizations are beneficial to performance, and GDroid can achieve up to 128X speedups against a plain GPU implementation.
Xiaodong Yu 0001, Fengguo Wei, Xinming Ou, Michela Becchi, Tekin Bicer, Danfeng Yao
IPDPS5
2020 Petascale XCT: 3D image reconstruction with hierarchical communications on multi-GPU nodes
abstract
X-ray computed tomography is a commonly used technique for noninvasive imaging at synchrotron facilities. Iterative tomographic reconstruction algorithms are often preferred for recovering high quality 3D volumetric images from 2D X-ray images, however, their use has been limited to small/medium datasets due to their computational requirements. In this paper, we propose a high-performance iterative reconstruction system for terabyte(s)-scale 3D volumes. Our design involves three novel optimizations: (1) optimization of (back)projection operators by extending the 2D memory-centric approach to 3D;(2) performing hierarchical communications by exploiting “fat-node” architecture with many GPUs; 3) utilization of mixed-precision types while preserving convergence rate and quality. We extensively evaluate the proposed optimizations and scaling on the Summit supercomputer. Our largest reconstruction is a mouse brain volume with 9×11K×11K voxels, where the total reconstruction time is under three minutes using 24,576 GPUs, reaching 65 PFLOPS: 34% of Summit's peak performance.
Mert Hidayetoglu, Tekin Bicer, Simon Garcia de Gonzalo, Bin Ren 0002, Vincent De Andrade, Doga Gürsoy, Rajkumar Kettimuthu, Ian T. Foster, Wen-Mei W. Hwu
SC2
2019 MemXCT: memory-centric X-ray CT reconstruction with massive parallelization
abstract
X-ray computed tomography (XCT)is used regularly at synchrotron light sources to study the internal morphology of materials at high resolution. However, experimental constraints, such as radiation sensitivity, can result in noisy or undersampled measurements. Further, depending on the resolution, sample size and data acquisition rates, the resulting noisy dataset can be terabyte-scale. Advanced iterative reconstruction techniques can produce high-quality images from noisy measurements, but their computational requirements have made their use exception rather than the rule. We propose here a novel memory-centric approach that avoids redundant computations at the expense of additional memory complexity. We develop a system, MemXCT, that uses an optimized SpMV implementation with two-level pseudo-Hilbert ordering and multi-stage input buffering. We evaluate MemXCT on various supercomputer architectures incolving KNL and GPU. MemXCT can reconstruct a large (11K×11K) mouse brain tomogram in ~10 seconds using 4096 KNL nodes (256K cores), the largest iterative reconstruction achieved in near-real time.
Mert Hidayetoglu, Tekin Bicer, Simon Garcia de Gonzalo, Bin Ren 0002, Doga Gürsoy, Rajkumar Kettimuthu, Ian T. Foster, Wen-Mei W. Hwu
SC2
2018 Graphphi: efficient parallel graph processing on emerging throughput-oriented architectures
abstract
Modern parallel architecture design has increasingly turned to throughput-oriented devices to address concerns about energy efficiency and power consumption. However, graph applications cannot tap into the full potential of such architectures because of highly unstructured computations and irregular memory accesses. In this paper, we present GraphPhi, a new approach to graph processing on emerging Intel Xeon Phi-like architectures, by addressing the restrictions of migrating existing graph processing frameworks on shared-memory multi-core CPUs to this new architecture.
Alexander Powell, Bo Wu 0002, Tekin Bicer, Bin Ren 0002
PACT4
2017 Real-Time Data Analysis and Autonomous Steering of Synchrotron Light Source Experiments
abstract
Modern scientific instruments, such as detectors at synchrotron light sources, can generate data at 10s of GB/sec. Current experimental protocols typically process and validate data only after an experiment has completed, which can lead to undetected errors and prevents online steering. Real-time data analysis can enable both detection of, and recovery from, errors, and optimization of data acquisition. We thus propose an autonomous stream processing system that allows data streamed from beamline computers to be processed in real time on a remote supercomputer, with a control feed-back loop used to make decisions during experimentation. We evaluate our system using two iterative tomographic reconstruction algorithms and varying data generation rates. These experiments are performed in a real-world environment in which data are streamed from a light source to a cluster for analysis and experimental control. We demonstrate that our system can sustain analysis rates of hundreds of projections per second by using up to 1,200 cores, while meeting stringent data quality constraints.
Tekin Bicer, Doga Gürsoy, Rajkumar Kettimuthu, Ian T. Foster, Bin Ren 0002, Vincent De Andrade, Francesco De Carlo
eScience1
2015 Rapid Tomographic Image Reconstruction via Large-Scale Parallelization
Tekin Bicer, Doga Gürsoy, Rajkumar Kettimuthu, Francesco De Carlo, Gagan Agrawal, Ian T. Foster
Euro-Par1
2015 Smart: a MapReduce-like framework for in-situ scientific analytics
abstract
In-situ analytics has lately been shown to be an effective approach to reduce both I/O and storage costs for scientific analytics. Developing an efficient in-situ implementation, however, involves many challenges, including parallelization, data movement or sharing, and resource allocation. Based on the premise that MapReduce can be an appropriate API for specifying scientific analytics applications, we present a novel MapReduce-like framework that supports efficient in-situ scientific analytics, and address several challenges that arise in applying the MapReduce idea for in-situ processing. Specifically, our implementation can load simulated data directly from distributed memory, and it uses a modified API that helps meet the strict memory constraints of in-situ analytics. The framework is designed so that analytics can be launched from the parallel code region of a simulation program. We have developed both time sharing and space sharing modes for maximizing the performance in different scenarios, with the former even avoiding any copying of data from simulation to the analytics program. We demonstrate the functionality, efficiency, and scalability of our system, by using different simulation and analytics programs, executed on clusters with multi-core and many-core nodes.
Gagan Agrawal, Tekin Bicer, Wei Jiang 0037
SC3
2014 Improving I/O Throughput of Scientific Applications Using Transparent Parallel Compression
abstract
Increasing number of cores in parallel computer systems are allowing scientific simulations to be executed with increasing spatial and temporal granularity. However, this also implies that increasing larger-sized datasets need to be output, stored, managed, and then visualized and/or analyzed using a variety of methods. In examining the possibility of using compression to accelerate all of these steps, we focus on two important questions: "Can compression help save time when data is output from, or input into, a parallel program?", and "How can a scientist's effort in using compression with a parallel program be minimized?". We focus on Pnet CDF, and show how transparent compression can be supported, thus allowing an existing simulation program to start outputting and storing data in a compressed fashion, and similarly, allow a data analysis application to read compressed data. We address challenges in supporting compression when parallel writes are being performed. In our experiments, we first analyze the effects of using compression with micro benchmarks, and then, continue our evaluation using a scientific simulation application, and two data analysis applications. While we obtain up to a factor of 2 improvement in performance for micro benchmarks, the execution time of simulation application is improved up to 22%, and the maximum speedup of data analysis applications is 1.83(with an average speedup of 1.36).
Tekin Bicer, Jian Yin 0002, Gagan Agrawal
CCGRID1
2013 Integrating Online Compression to Accelerate Large-Scale Data Analytics Applications
abstract
Compute cycles in high performance systems are increasing at a much faster pace than both storage and wide-area bandwidths. To continue improving the performance of large-scale data analytics applications, compression has therefore become promising approach. In this context, this paper makes the following contributions. First, we develop a new compression methodology, which exploits the similarities between spatial and/or temporal neighbors in a popular climate simulation dataset and enables high compression ratios and low decompression costs. Second, we develop a framework that can be used to incorporate a variety of compression and decompression algorithms. This framework also supports a simple API to allow integration with an existing application or data processing middleware. Once a compression algorithm is implemented, this framework automatically mechanizes multi-threaded retrieval, multi-threaded data decompression, and the use of informed prefetching and caching. By integrating this framework with a data-intensive middleware, we have applied our compression methodology and framework to three applications over two datasets, including the Global Cloud-Resolving Model (GCRM) climate dataset. We obtained an average compression ratio of 51.68%, and up to 53.27% improvement in execution time of data analysis applications by amortizing I/O time by moving compressed data.
Tekin Bicer, Jian Yin 0002, David Chiu 0001, Gagan Agrawal, Karen Schuchardt
IPDPS1
2012 Time and Cost Sensitive Data-Intensive Computing on Hybrid Clouds
abstract
Purpose-built clusters permeate many of today's organizations, providing both large-scale data storage and computing. Within local clusters, competition for resources complicates applications with deadlines. However, given the emergence of the cloud's pay-as-you-go model, users are increasingly storing portions of their data remotely and allocating compute nodes on-demand to meet deadlines. This scenario gives rise to a hybrid cloud, where data stored across local and cloud resources may be processed over both environments. While a hybrid execution environment may be used to meet time constraints, users must now attend to the costs associated with data storage, data transfer, and node allocation time on the cloud. In this paper, we describe a modeling-driven resource allocation framework to support both time and cost sensitive execution for data-intensive applications executed in a hybrid cloud setting. We evaluate our framework using two data-intensive applications and a number of time and cost constraints. Our experimental results show that our system is capable of meeting execution deadlines within a 3.6% margin of error. Similarly, cost constraints are met within a 1.2% margin of error, while minimizing the application's execution time.
Tekin Bicer, David Chiu 0001, Gagan Agrawal
CCGRID1
2011 A Framework for Data-Intensive Computing with Cloud Bursting
abstract
For many organizations, one attractive use of cloud resources can be through what is referred to as cloud bursting or the hybrid cloud. These refer to scenarios where an organization acquires and manages in-house resources to meet its base need, but can use additional resources from a cloud provider to maintain an acceptable response time during workload peaks. Cloud bursting has so far been discussed in the context of using additional computing resources from a cloud provider. However, as next generation applications are expected to see orders of magnitude increase in data set sizes, cloud resources can be used to store additional data after local resources are exhausted. In this paper, we consider the challenge of data analysis in a scenario where data is stored across a local cluster and cloud resources. We describe a software framework to enable data-intensive computing with cloud bursting, i.e., using a combination of compute resources from a local cluster and a cloud environment to perform Map-Reduce type processing on a data set that is geographically distributed. Our evaluation with three different applications shows that data-intensive computing with cloud bursting is feasible and scalable. Particularly, as compared to a situation where the data set is stored at one location and processed using resources at that end, the average slowdown of our system (using distributed but the same aggregate number of compute resources), is only 15.55%. Thus, the overheads due to global reduction, remote data retrieval, and potential load imbalance are quite manageable. Our system scales with an average speedup of 81% when the number of compute resources is doubled.
Tekin Bicer, David Chiu 0001, Gagan Agrawal
CLUSTER1
2010 Supporting fault tolerance in a data-intensive computing middleware
abstract
Over the last 2-3 years, the importance of data-intensive computing has increasingly been recognized, closely coupled with the emergence and popularity of map-reduce for developing this class of applications. Besides programmability and ease of parallelization, fault tolerance is clearly important for data-intensive applications, because of their long running nature, and because of the potential for using a large number of nodes for processing massive amounts of data. Fault-tolerance has been an important attribute of map-reduce as well in its Hadoop implementation, where it is based on replication of data in the file system. Two important goals in supporting fault-tolerance are low overheads and efficient recovery. With these goals, this paper describes a different approach for enabling data-intensive computing with fault-tolerance. Our approach is based on an API for developing data-intensive computations that is a variation of map-reduce, and it involves an explicit programmer-declared reduction object. We show how more efficient fault-tolerance support can be developed using this API. Particularly, as the reduction object represents the state of the computation on a node, we can periodically cache the reduction object from every node at another location and use it to support failure-recovery. We have extensively evaluated our approach using two data-intensive applications. Our results show that the overheads of our scheme are extremely low, and our system outperforms Hadoop both in absence and presence of failures.
Tekin Bicer, Wei Jiang 0037, Gagan Agrawal
IPDPS1