Yogesh L. Simmhan

dblp:s/YogeshSimmhan · also Yogesh Simmhan · DBLP profile ↗
← Back
98ranked-venue papers
16as first author
40since 2021 · last 2026
0000-0003-4140-7774ORCID · verified

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

Systems, architecture and hardware · 60 · 7 first-author · 29 since 2021Applied, interdisciplinary, general and emerging computing · 19 · 6 first-author · 3 since 2021Software engineering, systems software and programming languages · 16 · 6 first-author · 3 since 2021Databases, data management, data science and information retrieval · 10 · 1 first-author · 4 since 2021Artificial intelligence and machine learning · 5 · 1 first-author · 2 since 2021Computer networks · 2 · 2 since 2021Human-computer interaction and ubiquitous computing · 1 · 1 since 2021
YearPublicationVenuePosition
2026 Xfagent: Automating Multi-Cloud Deployment of Agentic Workflows on Faas Platforms
Varad Kulkarni, Vaibhav Jha, Nikhil Reddy, Anand Eswaran, Praveen Jayachandran, Yogesh L. Simmhan
CCGrid6
2026 ATLAS: Efficient Out-of-Core Inference for Billion-Scale Graph Neural Networks
abstract
Graph Neural Network (GNN) inference on billion-scale graphs is critical for domains like fintech and recommendation systems. Full-graph inference on these large graphs can be challenging due to high communication costs in distributed settings and high I/O costs in disk-backed Out-of-Core (OOC) settings. Existing OOC systems, operating across disk and memory, primarily focus on GNN training and perform poorly for full-graph inference due to massive read amplification, irregular I/O and memory pressure. We present ATLAS, a disk-based GNN inference framework that enables efficient full-graph, layer-wise inference on graphs whose topologies, features and intermediate embeddings exceed the available memory on single machines. ATLAS replaces gather-based execution with a broadcast-based model that enables sequential, single-pass streaming reads of features and embeddings per layer. A tiered memory–disk hierarchy with minimum-pending-message eviction, graph reordering and a GPU-accelerated pipeline sustains high throughput within 128 GiB RAM and 2 TiB SSD. Across out-of-core graphs with up to 4B edges and 550 GiB features and multiple GNN architectures, ATLAS improves end-to-end inference time by ≈ 12–30 × over State-of-the-Art (SOTA) OOC baselines on a single workstation, while remaining within \(\approx 5\%\) when features fit in memory.
Pranjal Naman, Yogesh L. Simmhan
HPDC2
2026 Billion-Scale Fintech Analytics: Scalable Data Management and Anomaly Detection at NPCI
Bharadwaj Dasari, Turaga Sai Dhiraj, Ganesh Jambhrunkar, Thirumalai Kailasam, Charu Vikram, Saurav Singla, Pranjal Naman, Yogesh L. Simmhan
ICDE8
2026 AeroResQ: Edge-accelerated UAV framework for scalable, resilient and collaborative escape route planning in wildfire scenarios
Suman Raj, Radhika Mittal, Rajiv Mayani, Pawel Zuk, Anirban Mandal, Michael Zink, Yogesh L. Simmhan, Ewa Deelman
Future Gener. Comput. Syst.7
2026 OptimES: Optimizing federated learning using remote embeddings for graph neural networks
Pranjal Naman, Yogesh L. Simmhan
J. Parallel Distributed Comput.2
2026 A Resource-centric Analysis and Optimization of NoSQL Workloads using Distressed Resource Volume Metric
Gunika Verma, Aashutosh A V, Pooja Srinivas, Yogesh L. Simmhan, Ayush Choure, Harshit Shah, Mayukh Das, Prashant Sasatte, Chetan Bansal, Abhijit Pai, Suraj Dixit, Achint Agrawal
Proc. VLDB Endow.4
2026 Characterizing FaaS Workflows on Public Clouds: The Good, the Bad and the Ugly
abstract
Function-as-a-service (FaaS) is a popular serverless computing paradigm for event-driven functions that elastically scale on public clouds. FaaS workflows (e.g.,AWS Step FunctionsandAzure Durable Functions), are composed from FaaS functions (e.g., AWS Lambda and Azure Functions) to build practical applications. But, the complex interactions between functions in the workflow and limited visibility into the internals of proprietary FaaS platforms are major impediments to analyzing a FaaS workflow's performance. While several works characterize FaaS platforms to derive such insights, or offer FaaS Workflow benchmarks, there is a lack of a principled of FaaS workflow platforms, which have unique scaling, performance and costing behavior influenced by the platform design, dataflow and workloads. In this article, we perform extensive evaluations of three popular FaaS workflow platforms from AWS and Azure, running 25 micro-benchmark and application workflows over$139k$invocations. Our detailed analysis confirms some conventional wisdom but also uncovers unique insights on the function execution, workflow orchestration, inter-function interactions, cold-start scaling and monetary costs. Our observations help developers better configure and program these platforms, set performance and scalability expectations, and identify research gaps on enhancing the platforms.
Varad Kulkarni, Nikhil Reddy, Tuhin Khare, Abhinandan S. Prasad, Chitra Babu, Yogesh L. Simmhan
IEEE Trans. Parallel Distributed Syst.6
2025 Choreography and Profiling of Quantum-Classical FaaS Workflows on Hybrid Clouds
abstract
Quantum computing is entering the mainstream as part of cloud offerings, where it serves as a special-purpose accelerator in larger applications. However, it is still challenging for developers and researchers to design, build and manage the resources for Hybrid Quantum-Classical (HQC) applications that include both classical logic (x86, ARM) and quantum circuit blocks, and run across both traditional and quantum processors. Further, quantum hardware is available on public cloud and even on-premise as private clouds, with varying capabilities, costs and queue times. Further, such quantum circuits also expose optimization methods that offer cost, time, accuracy and parallelism trade-offs. So, there is a compelling for easy composition of HQC applications that can be effortlessly and efficiently deployed on hybrid clouds. In this paper, we propose a framework to intuitively compose and deploy “zero-touch” quantum-classical Function-as-a-Service (FaaS) workflows through various workflow patterns that leverage diverse cloud system (workflow partitioning, adaptive polling) and quantum (circuit cutting, qubit reuse) optimizations. These utilize our XFaaS FaaS workflow framework for hybrid cloud deployments on AWS and Azure, and IBM Qiskit SDK for the quantum circuit toolchain. We also offer detailed experimental profiling of these optimizations for realistic and synthetic HQC applications on real clouds, and on real and simulated quantum hardware, and analyze the benefits of cloud system and quantum circuit optimizations. Our results demonstrate up to 53 % improvement in time and 80 % in cost when quantum circuit optimization on hardware is used in conjunction with dynamic fan-out.
Vaibhav Jha, Shikhar Srivastava 0003, Tarun Harishchandra Pal, Vaishnav Manoj Kavitha, Ritajit Majumdar, Tuhin Khare, Padmanabha Venkatagiri Seshadri, Varad Kulkarni, Anupama Ray, Yogesh L. Simmhan
CCGrid10
2025 Optimizing Federated Learning for Scalable Power-demand Forecasting in Microgrids
abstract
Real-time monitoring of power consumption in cities and micro-grids through the Internet of Things (IoT) can help forecast future demand and optimize grid operations. But moving all consumer-level usage data to the cloud for predictions and analysis at fine time scales can expose activity patterns. Federated Learning (FL) is a privacy-sensitive collaborative DNN training approach that retains data on edge devices, trains the models on private data locally, and aggregates the local models in the cloud. But key challenges exist: (i) clients can have non-independently identically distributed (non-IID) data, and (ii) the learning should be computationally cheap while scaling to 1000s of (unseen) clients. In this paper, we develop and evaluate several optimizations to FL training across edge and cloud for time-series demand forecasting in micro-grids and city-scale utilities using DNNs to achieve a high prediction accuracy while minimizing the training cost. We showcase the benefit of using exponentially weighted loss while training and show that it further improves the prediction of the final model. Finally, we evaluate these strategies by validating over 1000s of clients for three states in the US from the OpenEIA corpus, and performing FL both in a pseudo-distributed setting and a Pi edge cluster. The results highlight the benefits of the proposed methods over baselines like ARIMA and DNNs trained for individual consumers, which are not scalable.
Roopkatha Banerjee, Sampath Koti, Gyanendra Singh, Anirban Chakraborty 0001, Gurunath Gurrala, Bhushan Jagyasi, Yogesh L. Simmhan
eScience7
2025 Federated Learning Within Global Energy Budget over Heterogeneous Edge Accelerators
Roopkatha Banerjee, Tejus Chandrashekar, Ananth Eswar, Yogesh L. Simmhan
Euro-Par (1)4
2025 Ripple: Scalable Incremental GNN Inferencing on Large Streaming Graphs
abstract
Most real-world graphs are dynamic in nature, with continuous and rapid updates to the graph topology, and vertex and edge properties. Such frequent updates pose significant challenges for inferencing over Graph Neural Networks (GNNs). Current approaches that perform vertex-wise and layer-wise inferencing are impractical for dynamic graphs as they cause redundant computations, expand to large neighborhoods, and incur high communication costs for distributed setups, resulting in slow update propagation that often exceeds real-time latency requirements. This motivates the need for streaming GNN inference frameworks that are efficient and accurate over large, dynamic graphs. We propose Ripple, a framework that performs fast incremental updates of embeddings arising due to updates to the graph topology or vertex features. Ripple provides a generalized incremental programming model, leveraging the properties of the underlying aggregation functions employed by GNNs to efficiently propagate updates to the affected neighborhood and compute the exact new embeddings. Besides a single-machine design, we also extend this execution model to distributed inferencing, to support large graphs that do not fit in a single machine’s memory. Ripple on a single machine achieves up to ≈ 28000 updates/sec for sparse graphs like Arxiv and ≈1200 updates/sec for larger and denser graphs like Products, with latencies of 0.1ms–1s that are required for near-realtime applications. The distributed version of Ripple offers up to ≈30× better throughput over the baselines, due to 70× lower communication costs during updates.
Pranjal Naman, Yogesh L. Simmhan
ICDCS2
2025 AeroDaaS: Towards an Application Programming Framework for Drones-as-a-Service
abstract
The increasing adoption of UAVs with advanced sensors and GPU-accelerated edge computing has enabled real-time AI-driven applications in fields such as precision agriculture, wildfire monitoring, and environmental conservation. However, integrating deep learning on UAVs remains challenging due to platform heterogeneity, real-time constraints, and the need for seamless cloud-edge coordination. To address these challenges, we introduce AeroDaaS, a service-oriented framework that abstracts UAV-based sensing complexities and provides a Drone-as-a-Service (DaaS) model for intelligent decision-making. AeroDaaS offers modular service primitives for on-demand UAV sensing, navigation, and analytics as composable microservices, ensuring cross-platform compatibility and scalability across heterogeneous UAV and edge-cloud infrastructures. We implement and evaluate a preliminary version of AeroDaaS for two real-world DaaS applications. We require$\leq 40$lines of code for the applications and see minimal platform overhead of$\leq 20 ~\text{ms}$per frame and$\leq 0.5$GB memory usage on Orin Nano. These early results are promising for AeroDaaS as an efficient, flexible and scalable UAV programming framework for autonomous aerial analytics.
Suman Raj, Rajdeep Singh, Kautuk Astu, Yogesh L. Simmhan
ICWS4
2025 Adaptive heuristics for scheduling DNN inferencing on edge and cloud for personalized UAV fleets
Suman Raj, Radhika Mittal, Harshil Gupta, Yogesh L. Simmhan
Future Gener. Comput. Syst.4
2025 Flotilla: A scalable, modular and resilient federated learning framework for heterogeneous resources
Roopkatha Banerjee, Prince Modi, Jinal Vyas, Chunduru Sri Abhijit, Tejus Chandrashekar, Harsha Varun Marisetty, Manik Gupta, Yogesh L. Simmhan
J. Parallel Distributed Comput.8
2025 AerialDB: A federated peer-to-peer spatio-temporal edge datastore for drone fleets
Shashwat Jaiswal, Suman Raj, Subhajit Sidhanta, Yogesh L. Simmhan
Pervasive Mob. Comput.4
2025 Triparts: Scalable Streaming Graph Partitioning to Enhance Community Structure
abstract
k-way edge based partitioning algorithms for processing large streaming graphs, such as social networks and web crawls, assign each arriving edge to one of the k partitions. This can result in vertices being replicated on multiple partitions. Typically, such partitioning algorithms aim to balance the edge counts across partitions while minimizing the vertex replication. However, such objectives ignore the community structure inherently embedded in the graph, which is an important quality metric for clustering and graph mining applications that subsequently operate on the partitions. To address this gap, we propose a novel optimization goal to maximize the number of local triangles in the partitions as an additional objective. Triangle count is an effective metric to measure the conservation of community structure. Further, we propose TriParts a family of heuristics for online partitioning over an edge stream. They use three complementary state data structures: Bloom Filters, Triangle Map and High degree Map. Each state adds tangible value to meet our objectives. We validate TriParts on six diverse real world graphs with up to 1.6B edges and varying triangle densities. Our best heuristic outperforms the state-of-the-art DBH and HDRF streaming graph partitioners on the triangle-count metric by up to 4–8.3x while maintaining competitive vertex replication factor and edge-balancing. We achieve an ingest rate of 500k edges/sec on a 16 node cluster. We also offer detailed results on the configuration parameters, scalability and overheads of TriParts, and its practical benefits for distributed graph analytics.
Ruchi Bhoot, Tuhin Khare, Manoj Agarwal, Siddharth D. Jaiswal, Yogesh L. Simmhan
Proc. VLDB Endow.5
2025 CARL: Cost-Optimized Online Container Placement on VMs Using Adversarial Reinforcement Learning
abstract
Containerization has become popular for the deployment of applications on public clouds. Large enterprises may host 100 s of applications on 1000 s containers that are placed onto Virtual Machines (VMs). Such placement decisions happen continuously as applications are updated by DevOps pipelines that deploy the containers. Managing the placement of container resource requests onto the available capacities of VMs needs to be cost-efficient. This is well-studied, and usually modelled as a multi-dimensional Vector Bin-packing Problem (VBP). Many heuristics, and recently machine learning approaches, have been developed to solve this NP-hard problem for real-time decisions. We propose CARL, a novel approach to solve VBP through Adversarial Reinforcement Learning (RL) for cost minimization. It mimics the placement behavior of an offline semi-optimal VBP solver (teacher), while automatically learning a reward function for reducing the VM costs which out-performs the teacher. It requires limited historical container workload traces to train, and is resilient to changes in the workload distribution during inferencing. We extensively evaluate CARL on workloads derived from realistic traces from Google and Alibaba for the placement of 5 k–10 k container requests onto 2 k–8 k VMs, and compare it with classic heuristics and state-of-the-art RL methods. (1) CARL isfast, e.g., making placement decisions at$\approx 1900$requests/sec onto 8,900 candidate VMs. (2) It isefficient, achieving$\approx 16\%$lower VM costs than classic and contemporary RL methods. (3) It isrobustto changes in the workload, offering competitive results even when the resource needs or inter-arrival time of the container requests skew from the training workload.
Prathamesh Saraf Vinayak, Saswat Subhajyoti Mallick, Lakshmi Jagarlamudi, Anirban Chakraborty 0001, Yogesh L. Simmhan
IEEE Trans. Cloud Comput.5
2024 XFBench: A Cross-Cloud Benchmark Suite for Evaluating FaaS Workflow Platforms
abstract
Functions-as-a-Service (FaaS) is a widely used serverless computing abstraction that helps developers build applications using event-driven, stateless functions that execute on the cloud. Commercial FaaS platforms such as AWS Lambda and Azure Functions offer elastic auto-scaling and invocation-level billing to ease operations. Applications are often composed as a dataflow of FaaS functions that are orchestrated by FaaS workflow platforms such as AWS Step Functions or Azure Durable Functions. However, the proprietary nature of FaaS platforms on public clouds means that their internals are less understood. While benchmarks to characterize FaaS platforms exist, none are available for a principled evaluation of FaaS workflow platforms. Further, they are less configurable and often limited to simple workloads and a single cloud provider. We address this limitation by proposing XFBench, an end-to-end automated benchmarking framework for FaaS workflows that works across clouds, and an accompanying function, workflow, and workload suite. The user provides a generic definition of the workflow and workload for benchmarking, and XFBench automatically deploys the workflows across multiple cloud platforms, generates client requests, and profiles the execution. We validate XFBench with realistic workflows and workloads on AWS and Azure platforms in different global regions to offer early insights into understanding the inter-function communication, function execution time, and cold start scaling.
Varad Kulkarni, Nikhil Reddy, Tuhin Khare, Harini Mohan, Jahnavi Murali, Mohith A, Ragul B, Sanjai Balajee, Sanjjit S, Swathika D, Vaishnavi S, Yashasvee V, Chitra Babu, Abhinandan S. Prasad, Yogesh L. Simmhan
CCGrid15
2024 Building AI Agents for Autonomous Clouds: Challenges and Design Principles
abstract
The rapid growth in the use of Large Language Models (LLMs) and AI Agents as part of software development and deployment is revolutionizing the information technology landscape. While code generation receives significant attention, a higher-impact application lies in using agents for the operational resilience of cloud services, which currently require significant human effort and domain knowledge. There is a growing interest in AI for IT Operations (AIOps), which aims to automate complex operational tasks, like fault localization and root cause analysis, reducing human intervention and customer impact. However, achieving the vision of autonomous and self-healing clouds through AIOps is hampered by the lack of standardized frameworks for building, evaluating, and improving AIOps agents. This vision paper lays the groundwork for such a framework by framing the requirements and then discussing design decisions that satisfy them. We also propose AIOpsLab, a prototype implementation leveraging agent-cloud-interface that orchestrates an application, injects real-time faults using chaos engineering, and interfaces with an agent to localize and resolve the faults. We report promising results and lay the groundwork to build a modular and robust framework for building, evaluating, and improving agents for autonomous clouds.
Manish Shetty, Yinfang Chen, Gagan Somashekar, Minghua Ma, Yogesh L. Simmhan, Xuchao Zhang, Jonathan Mace, Dax Vandevoorde, Pedro Henrique B. Las-Casas, Shachee Mishra Gupta, Suman Nath, Chetan Bansal, Saravan Rajmohan
SoCC5
2024 Optimizing Federated Learning Using Remote Embeddings for Graph Neural Networks
Pranjal Naman, Yogesh L. Simmhan
Euro-Par (2)2
2024 PowerTrain: Fast, generalizable time and power prediction models to optimize DNN training on accelerated edges
Prashanthi S. K, Saisamarth Taluri, Beautlin S, Lakshya Karwa, Yogesh L. Simmhan
Future Gener. Comput. Syst.5
2024 Improved Algorithms for Co-Scheduling of Edge Analytics and Routes for UAV Fleet Missions
abstract
Unmanned Aerial Vehicles (UAVs) or drones are increasingly used for urban applications like traffic monitoring and construction surveys. Autonomous navigation allows drones to visitwaypointsand accomplishactivitiesas part of theirmission. A common activity is to hover and observe a location using on-board cameras. Advances in Deep Neural Networks (DNNs) allow such videos to be analyzed for automated decision making. UAVs also host edge computing capability for on-board inferencing by such DNNs. To this end, for a fleet of drones, we propose a novelMission Scheduling Problem ()that co-schedules the flight routes to visit and record video at waypoints, and their subsequent on-board edge analytics. The proposed schedule maximizes the data capture and computing utilities from the activities while meeting the activity deadlines, and the energy and computing constraints. We first prove that is NP-hard and then optimally solve it by formulating a mixed integer linear programming (MILP) problem. Next, we design five time-efficient heuristic algorithms that provide sub-optimal but fast solutions that are empirically competitive with the optimal solution. Evaluation of these five schedulers using real drone traces demonstrate utility–runtime trade-offs under diverse workloads.
Aakash Khochare, Francesco Betti Sorbelli, Yogesh L. Simmhan, Sajal K. Das 0001
IEEE/ACM Trans. Netw.3
2024 TARIS: Scalable Incremental Processing of Time-Respecting Algorithms on Streaming Graphs
abstract
Temporal graphs change with time and have a lifespan associated with each vertex and edge. These graphs are suitable to process time-respecting algorithms where the traversed edges must have monotonic timestamps. Interval-centric Computing Model (ICM) is a distributed programming abstraction to design such temporal algorithms. There has been little work on supporting time-respecting algorithms at large scales for streaming graphs, which are updated continuously at high rates (Millions/s), such as in financial and social networks. In this article, we extend the windowed-variant of ICM for incremental computing over streaming graph updates. We formalize the properties of temporal graph algorithms and prove that our model of incremental computing over streaming updates is equivalent to batch execution of ICM. We design TARIS, a novel distributed graph platform that implements these incremental computing features. We use efficient data structures to reduce memory access and enhance locality during graph updates. We also propose scheduling strategies to interleave updates with computing, and streaming strategies to adapt the execution window for incremental computing to the variable input rates. Our detailed and rigorous evaluation of temporal algorithms on large-scale graphs with up to$2\,\text{B}$edges show that TARIS out-performs contemporary baselines, Tink and Gradoop, by 3–4 orders of magnitude, and handles a high input rate of$ 83k$–$ 587\,\text{M}$Mutations/s with latencies in the order of seconds–minutes.
Ruchi Bhoot, Suved Sanjay Ghanmode, Yogesh L. Simmhan
IEEE Trans. Parallel Distributed Syst.3
2023 XFaaS: Cross-platform Orchestration of FaaS Workflows on Hybrid Clouds
abstract
Functions as a Service (FaaS) have gained popularity for programming public clouds due to their simple abstraction, ease of deployment, effortless scaling and granular billing. Cloud providers also offer basic capabilities to compose these functions into workflows. FaaS and FaaS workflow models, however, are proprietary to each cloud provider. This prevents their portability across cloud providers, and requires effort to design workflows that run on different cloud providers or data centers. Such requirements are increasingly important to meet regulatory requirements, leverage cost arbitrage and avoid vendor lock-in. Further, the FaaS execution models are also different, and the overheads of FaaS workflows due to message indirection and cold-starts need custom optimizations for different platforms. In this paper, we propose XFaaS, a cross-platform deployment and orchestration engine for FaaS workflows to operate on multiple clouds. XFaaS allows “zero touch” deployment of functions and workflows across AWS and Azure clouds by automatically generating the necessary code wrappers, cloud queues, and coordinating with the native FaaS engine of the cloud providers. It also uses intelligent function fusion and placement logic to reduce the workflow execution latency in a hybrid cloud while mitigating costs, using performance and billing models specific to the providers based in detailed benchmarks. Our empirical results indicate that fusion offers up to ≈75 % benefits in latency and ≈57% reduction in cost, while placement strategies reduce the latency by ≈ 24%, compared to baselines in the best cases.
Aakash Khochare, Tuhin Khare, Varad Kulkarni, Yogesh L. Simmhan
CCGrid4
2023 CCGRID 2023: A Holistic Approach to Inclusion and Belonging
abstract
“CCGRID will act with responsibility as its primary consideration; with equity, diversity, and inclusion as its central goals.” from the CCGRID 2023 web site [1]
Beth Plale, Preeti Malakar, Meenakshi D'Souza, Hemangee K. Kapoor, Yogesh L. Simmhan, Ilkay Altintas, S. Manohar 0001
CCGrid5
2023 Scheduling DNN Inferencing on Edge and Cloud for Personalized UAV Fleets
abstract
Drone fleets with onboard cameras coupled with DNN inferencing models can support diverse applications, from infrastructure monitoring to package deliveries. Here, we propose to use one or more “buddy” drones to help Visually Impaired People (VIPs) lead an active lifestyle. Video inferencing tasks from such drones are used to navigate the drone and alert the VIP to threats, and hence have strict execution deadlines. They have a choice to execute either on an accelerated edge like Nvidia Jetson linked to the drone, or on a cloud INFerencing-as-a-Service (INFaaS). However, making this decision is a challenge given the latency and cost trade-offs, and network variability in outdoor environments. We propose a deadline-driven heuristic to schedule a stream of diverse DNN inferencing tasks executing over video segments generated by multiple drones linked to an edge, with the option to execute on the cloud. We use strategies like task dropping, work stealing and migration, and dynamic adaptation to cloud variability, to fully utilize the captive edge with intelligent offloading to the cloud, to maximize the utility and the number of tasks completed. We evaluate our strategies using a setup that emulates a fleet of > 50 drones within city conditions supporting> 25 VIPs, with real DNN models executing on drone video streams, using Jetson Nano edges and AWS Lambda cloud functions. Our detailed comparison of our strategy exhibits a task completion rate of up to 91 %, up to 2.5× higher utility compared to the baselines and 68% higher utility with network variability.
Suman Raj, Harshil Gupta, Yogesh L. Simmhan
CCGrid3
2023 A Lossless Compression Pipeline for Petabyte-Scale Whole Genome Sequencing Data
abstract
Whole genome sequencing (WGS) technologies have enabled high-throughput cost-effective genome sequencing at the population scale. A single WGS instrument can sequence millions of DNA molecules simultaneously, leading to the generation of massive datasets. GenomeIndia is an ongoing national project aimed at sequencing the genomes of 10,000 Indian individuals. The GenomeIndia sequencing centers are completing the generation of petabyte-scale genomic data. This has raised an urgent need for scalable lossless compression software to facilitate cost-effective storage and exchange of data. By default, each WGS file produced in the GenomeIndia project is stored in the standard unmapped BAM (uBAM) format. A uBAM file saves the DNA sequences as well as metadata associated with the sequencing experiment. We have developed an open-source software pipeline that enables parallel lossless compression and decompression of uBAM files. It produces compressed output that is approximately 5 x smaller than the input uBAM files. We carefully engineered the pipeline by integrating different bioinformatics tools such as SPRING, Picard, SAMtools, and PySAM. We evaluated the parallel efficiency of our approach using thorough performance profiling and strong-scaling experiments.
Ajeya Bhat, Sai Manasa Chadalavada, Nagakishore Jammula, Yogesh L. Simmhan
HiPC5
2023 Performance Characterization of Containerized DNN Training and Inference on Edge Accelerators
abstract
Edge devices have typically been used for DNN in-ferencing. The increase in the compute power of accelerated edges is leading to their use in DNN training also. As privacy becomes a concern on multi-tenant edge devices, Docker containers provide a lightweight virtualization mechanism to sandbox models. But their overheads for edge devices are not yet explored. In this work, we study the impact of containerized DNN inference and training workloads on an NVIDIA AGX Orin edge device and contrast it against bare metal execution on running time, CPU, GPU and memory utilization, and energy consumption. Our analysis provides several interesting insights on these overheads.
Prashanthi S. K, Vinayaka Hegde, Keerthana Patchava, Ankita Das, Yogesh L. Simmhan
HiPC5
2023 A Co-Simulation Framework for Communication and Control in Autonomous Multi-Robot Systems
abstract
Multi-Robot Systems (MRS) are transforming diverse domains like logistics, cargo management, and agriculture. However, ensuring that the behavior is correct under various network conditions within such complex environments is challenging, and meeting the desired automation goals is difficult. We propose the CORNET 2.0 co-simulation framework to jointly and accurately simulate multi-agent robotic systems within physical environments and the communication network models within such environments. Our modular framework allows diverse robot and network models to seamlessly integrate to simulate the robot's autonomy, physical space, and network features, such as latency, throughput, and loss intrinsic to the network topology and communication technology. A key novelty of CORNET 2.0 is its accurate synchronizing of mobility and time, which ensures that the physical location of a robot at a point in time, and the network properties and packets that flow from that location, are aligned. This is vital to model and validate MRS coordination algorithms that rely on network interactions. We provide a detailed evaluation of CORNET 2.0 in modeling real-world MRS use cases, such as leader-follower and warehouse environments, that help highlights the benefits.
Srikrishna Acharya, Mukunda Bharatheesha, Yogesh L. Simmhan, Bharadwaj S. Amrutur
IROS3
2022 Resilient Execution of Data-triggered Applications on Edge, Fog and Cloud Resources
abstract
Internet of Things (loT) is leading to the pervasive availability of streaming data about the physical world, coupled with edge computing infrastructure deployed as part of smart cities and 5G rollout. These constrained, less reliable but cheap resources are complemented by fog resources that offer feder-ated management and accelerated computing, and pay-as-you-go cloud resources. There is a lack of intuitive means to deploy application pipelines to consume such diverse streams, and to execute them reliably on edge and fog resources. We propose an innovative application model to declaratively specify queries to match streams of micro-batch data from stream sources and trigger the distributed execution of data pipelines. We also design a resilient scheduling strategy using advanced reservation on reliable fogs to guarantee dataflow completion within a deadline while minimizing the execution cost. Our detailed experiments on over 100 virtual loT resources and for$\approx 10k$task executions, with comparison against baseline scheduling strategies, illustrates the cost-effectiveness, resilience and scalability of our framework.
Prateeksha Varshney, Shriram Ramesh, Shayal Chhabra, Aakash Khochare, Yogesh L. Simmhan
CCGRID5
2022 Toward Scientific Workflows in a Serverless World
abstract
Serverless computing and FaaS have gained popularity due to their ease of design, deployment, scaling and billing on clouds. However, when used to compose and orchestrate scientific workflows, they pose limitations due to cold starts, message indirection, vendor lock-in and lack of provenance support. Here, we propose a design for a Ser verless Scientific Workflow Orchestrator that overcomes these challenges using techniques like function fusion, pilot invocations and data fabrics.
Aakash Khochare, Yogesh L. Simmhan, Sameep Mehta, Arvind Agarwal
e-Science2
2022 Optimizing the interval-centric distributed computing model for temporal graph algorithms
abstract
Temporal graphs assign lifespans to their vertices, edges and attributes. Large temporal graphs are common for finding the shortest paths in transit networks and contact tracing for COVID-19. Graph programming abstractions like Interval-centric Computing Model (ICM) extend Google's Pregel model to intuitively compose and execute time-dependent graph algorithms in a distributed environment. However, the benefits of easier algorithmic design in ICM are offset by performance bottlenecks in its TimeWarp shuffle and messaging phases. Here, we design several optimizations to ICM to reduce these overheads. We propose local optimizations within a vertex execution by unrolling messages before TimeWarp (LU), and deferring messaging till all local computations complete (DS). We also temporally partition the interval graph into windows (WICM) to flatten the execution load. We offer a proof of equivalence between ICM and these techniques. Our detailed empirical evaluation for six real-world graphs with up to 133M vertices, 5.5B edges and 365 time-points, for six temporal traversal algorithms executing on a commodity cluster with 8 nodes, shows that LU, DS and WICM together significantly reduce the average algorithm runtime by ≈ 61% (≈ 15 mins) over ICM, and reduce message communication by ≈ 38%(≈ 3.2B) on average.
Animesh Baranawal, Yogesh L. Simmhan
EuroSys2
2022 DiPETrans: A framework for distributed parallel execution of transactions of blocks in blockchains
abstract
Summary Contemporary blockchain such as Bitcoin and Ethereum execute transactions serially by miners and validators and determine the Proof‐of‐Work (PoW). Such serial execution is unable to exploit modern multi‐core resources efficiently, hence limiting the system throughput and increasing the transaction acceptance latency. The objective of this work is to increase the transaction throughput by introducing parallel transaction execution using a static analysis over the transaction dependencies. We propose the DiPETrans framework for distributed execution of transactions in a block. Here, peers in the blockchain network form a community of trusted nodes to execute the transactions and find the PoW in‐parallel, using a leader–follower approach. During mining, the leader statically analyzes the transactions, creates different groups (shards) of independent transactions, and distributes them to followers to execute concurrently. After execution, the community's compute power is utilized to solve the PoW concurrently. When a block is successfully created, the leader broadcasts the proposed block to other peers in the network for validation. On receiving a block, the validators re‐execute the block transactions and accept the block if they reach the same state as shared by the miner. Validation can also be done in parallel, following the same leader–follower approach as mining. We report experiments using over 5 million real transactions from the Ethereum blockchain and execute them using our DiPETrans framework to empirically validate the benefits of our techniques over a traditional sequential execution. We achieve a maximum speedup of 2.2 and 2.0 and an average speedup of 1.6 and 1.5 for the miner and the validator, respectively, with 100–500 transactions per block when using 6 machines in the community. Further, we achieve a peak of 5 end‐to‐end block creation speedup using a parallel miner over a serial miner.
Shrey Baheti, Parwat Singh Anjana, Sathya Peri, Yogesh L. Simmhan
Concurr. Comput. Pract. Exp.4
2022 Guest editorial: Special issue on the 2020 IEEE symposium on real-time distributed computing (ISORC)
Tommaso Cucinotta, Frank Mueller 0001, Yogesh L. Simmhan
J. Syst. Archit.3
2021 Event Related Data Collection from Microblog Streams
Manoj K. Agarwal, Animesh Baranawal, Yogesh L. Simmhan, Manish Gupta 0001
DEXA (2)3
2021 Heuristic Algorithms for Co-scheduling of Edge Analytics and Routes for UAV Fleet Missions
abstract
Unmanned Aerial Vehicles (UAVs) or drones are increasingly used for urban applications like traffic monitoring and construction surveys. Autonomous navigation allows drones to visit waypoints and accomplish activities as part of their mission. A common activity is to hover and observe a location using on-board cameras. Advances in Deep Neural Networks (DNNs) allow such videos to be analyzed for automated decision making. UAVs also host edge computing capability for on-board inferencing by such DNNs. To this end, for a fleet of drones, we propose a novel Mission Scheduling Problem (MSP) that co-schedules the flight routes to visit and record video at waypoints, and their subsequent on-board edge analytics. The proposed schedule maximizes the utility from the activities while meeting activity deadlines as well as energy and computing constraints. We first prove that MSP is NP-hard and then optimally solve it by formulating a mixed integer linear programming (MILP) problem. Next, we design two efficient heuristic algorithms, jsc and vrc, that provide fast sub-optimal solutions. Evaluation of these three schedulers using real drone traces demonstrate utility-runtime trade-offs under diverse workloads.
Aakash Khochare, Yogesh L. Simmhan, Francesco Betti Sorbelli, Sajal K. Das 0001
INFOCOM2
2021 Granite: A distributed engine for scalable path queries over temporal property graphs
Shriram Ramesh, Animesh Baranawal, Yogesh L. Simmhan
J. Parallel Distributed Comput.3
2021 Cost-Effective Sharing of Streaming Dataflows for IoT Applications
abstract
Internet of Things (IoT) applications are often designed as dataflows that analyze sensor data in real-time to make decisions. Stream processing systems likeApache Stormexecute these on Cloud infrastructure. As IoT applications within shared data environments like smart cities grow, they will duplicate tasks like pre-processing and analytics. This offers the opportunity to collaboratively reuse the outputs of overlapping dataflows, improving the resource efficiency on Clouds. We proposedataflow reuse algorithmsthat when given a submitted dataflow, identify the intersection of reusable tasks and streams from existing dataflows to form amerged dataflow, with guaranteed equivalence of their output streams. Algorithms to unmerge dataflows when they are removed, and defragment partially reused dataflows are also proposed. We implement these algorithms for the Storm fast-data platform, and validate their performance and resource savings using 86 real and synthetic dataflows from eScience and IoT domains. Our reuse strategies reduce the number of running tasks by 34–45 percent and the cumulative CPU usage by 29–63 percent. Including defragmentation of incremental dataflows achieves a monetary savings on Cloud resources of 36–44 percent compared to dataflows without reuse, and has limited redeployment overheads.
Shilpa Chaturvedi, Sahil Tyagi, Yogesh L. Simmhan
IEEE Trans. Cloud Comput.3
2021 VIoLET: An Emulation Environment for Validating IoT Deployments at Large Scales
abstract
Internet of Things (IoT) deployments have been growing manifold, encompassing sensors, networks, edge, fog, and cloud resources. Despite the intense interest from researchers and practitioners, most do not have access to large-scale IoT testbeds for validation. Simulation environments that allow analytical modeling are a poor substitute for evaluating software platforms or application workloads in realistic computing environments. Here, we propose a virtual environment for validating Internet of Things at large scales (VIoLET), an emulator for defining and launching large-scale IoT deployments within cloud VMs. It allows users to declaratively specify container-based compute resources that match the performance of native IoT compute devices using Docker. These can be inter-connected by complex topologies on which bandwidth and latency rules are enforced. Users can configure synthetic sensors for data generation as well. We also incorporate models for CPU resource dynamism, and for failure and recovery of the underlying devices. We offer a detailed comparison of VIoLET’s compute and network performance between the virtual and physical deployments, evaluate its scaling with deployments with up to 1, 000 devices and 4, 000 device-cores, and validate its ability to model resource dynamism. Our extensive experiments show that the performance of the virtual IoT environment accurately matches the expected behavior, with deviations levels within what is seen in actual physical devices. It also scales to 1, 000s of devices and at a modest cloud computing costs of under 0.15% of the actual hardware cost, per hour of use, with minimal management effort. This IoT emulation environment fills an essential gap between IoT simulators and real deployments.
Shrey Baheti, Shreyas Badiger, Yogesh L. Simmhan
ACM Trans. Cyber Phys. Syst.3
2021 A Scalable Platform for Distributed Object Tracking Across a Many-Camera Network
abstract
Advances in deep neural networks (DNN) and computer vision (CV) algorithms have made it feasible to extract meaningful insights from large-scale deployments of urban cameras. Tracking an object of interest across the camera network in near real-time is a canonical problem. However, current tracking platforms have two key limitations: 1) They are monolithic, proprietary and lack the ability to rapidly incorporate sophisticated tracking models, and 2) They are less responsive to dynamism across wide-area computing resources that include edge, fog, and cloud abstractions. We address these gaps using Anveshak, a runtime platform for composing and coordinating distributed tracking applications. It provides a domain-specific dataflow programming model to intuitively compose a tracking application, supporting contemporary CV advances like query fusion and re-identification, and enabling dynamic scoping of the camera network's search space to avoid wasted computation. We also offer tunable batching and data-dropping strategies for dataflow blocks deployed on distributed resources to respond to network and compute variability. These balance the tracking accuracy, its real-time performance, and the active camera-set size. We illustrate the concise expressiveness of the programming model for four tracking applications. Our detailed experiments for a network of 1000 camera-feeds on modest resources exhibit the tunable scalability, performance, and quality trade-offs enabled by our dynamic tracking, batching, and dropping strategies.
Aakash Khochare, Aravindhan Krishnan, Yogesh L. Simmhan
IEEE Trans. Parallel Distributed Syst.3
2020 A Distributed Path Query Engine for Temporal Property Graphs
abstract
Property graphs are a common form of linked data, with path queries used to traverse and explore them for enterprise transactions and mining. Temporal property graphs are a recent variant where time is a first-class entity to be queried over, and their properties and structure vary over time. These are seen in social, telecom and transit networks. However, current graph databases and query engines have limited support for temporal relations among graph entities, no support for time-varying entities and/or do not scale on distributed resources. We address this gap by extending a linear path query model over property graphs to include intuitive temporal predicates that operate over temporal graphs. We design a distributed execution model for these temporal path queries using the interval-centric computing model, and develop a novel cost model to select an efficient execution plan from several. We perform detailed experiments of our $\mathcal{G}ranite$ distributed query engine using temporal property graphs as large as 52M vertices, 218M edges and 118M properties, and an 800-query workload, derived from the LDBC benchmark. We offer sub-second query latencies in most cases, which is 149×-1140× faster compared to industry-leading Neo4J shared- memory graph database and the JanusGraph/Spark distributed graph query engine. Further, our cost model selects a query plan that is within 10% of the optimal execution time in 90% of the cases. We also scale well, and complete 100% of the queries for all graphs, compared to only 32-92% by baseline systems.
Shriram Ramesh, Animesh Baranawal, Yogesh L. Simmhan
CCGRID3
2020 TorqueDB: Distributed Querying of Time-Series Data from Edge-local Storage
Dhruv Garg, Prathik Shirolkar, Anshu Shukla, Yogesh L. Simmhan
Euro-Par4
2020 An Interval-centric Model for Distributed Computing over Temporal Graphs
abstract
Algorithms for temporal property graphs may be time-dependent (TD), navigating the structure and time concurrently, or time-independent (TI), operating separately on different snapshots. Currently, there is no unified and scalable programming abstraction to design TI and TD algorithms over large temporal graphs. We propose an interval-centric computing model (ICM) for distributed and iterative processing of temporal graphs, where a vertex's time-interval is a unit of data-parallel computation. It introduces a unique time-warp operator for temporal partitioning and grouping of messages that hides the complexity of designing temporal algorithms, while avoiding redundancy in user logic calls and messages sent. GRAPHITE is our implementation of ICM over Apache Giraph, and we use it to design 12 TI and TD algorithms from literature. We rigorously evaluate its performance for diverse real-world temporal graphs - as large as 131M vertices and 5.5B edges, and as long as 219 snapshots. Our comparison with 4 baseline platforms on a 10-node commodity cluster shows that ICM shares compute and messaging across intervals to out-perform them by up to 25×, and matches them even in worst-case scenarios. GRAPHITE also exhibits weak-scaling with near-perfect efficiency.
Swapnil Gandhi, Yogesh L. Simmhan
ICDE2
2020 Characterizing application scheduling on edge, fog, and cloud computing resources
abstract
Summary Cloud computing has grown to become a popular distributed computing service offered by commercial providers. More recently, edge and fog computing resources have emerged on the wide‐area network as part of Internet of things (IoT) deployments. These three resource abstraction layers are complementary, and offer distinctive benefits. Scheduling applications on clouds has been an active area of research, with workflow and data flow models offering a flexible abstraction to specify applications for execution. However, the application programming and scheduling models for edge and fog are still maturing, and can benefit from learnings on cloud resources. At the same time, there is also value in using these resources cohesively for application execution. In this article, we offer a taxonomy of concepts essential for specifying and solving the problem of scheduling applications on edge, fog, and cloud computing resources. We first characterize the resource capabilities and limitations of these infrastructure and offer a taxonomy of application models, quality‐of‐service constraints and goals, and scheduling techniques, based on a literature review. We also tabulate key research prototypes and papers using this taxonomy. This survey benefits developers and researchers on these distributed resources in designing and categorizing their applications, selecting the relevant computing abstraction(s), and developing or selecting the appropriate scheduling algorithm. It also highlights gaps in literature where open problems remain.
Prateeksha Varshney, Yogesh L. Simmhan
Softw. Pract. Exp.2
2019 Adaptive Partition Migration for Irregular Graph Algorithms on Elastic Resources
abstract
Component-centric graph programming models allow distributed graph algorithms to be composed, and executed in an iterative manner on commodity clusters and Clouds. Graphs are partitioned and statically placed on a fixed number of machines before execution. However, many graph algorithms have an irregular execution behavior across partitions in different iterations, which causes resource under-utilization. We propose wo strategies, First Fit Decreasing with Migration Planning (FFDMP) and MinMax, for adaptive partition placement onto an elastic number of Cloud resources for such irregular algorithms. For each iteration, our strategies decide the number of hosts and the placement of partitions on them to balance the compute load, and enact this by migrating partitions between hosts at iteration boundaries. Unlike others, our strategies actively consider the time and cost penalties for moving partitions between hosts, and reduce the overall cost of execution while mitigating any increase in makespan. We implement these strategies on our GoFFish subgraph-centric graph processing platform, and evaluate them for performing Breadth First Search (BFS) on large real-world graphs with 10^7-10^9 edges. Our results show that the proposed strategies reduce the median resource cost by 13-38% when compared to a static placement, increase the median makespan by 1-33%, which is strictly bound by a given time budget, and also out-perform existing baseline scheduling algorithms from literature.
Ravikant Dindokar, Yogesh L. Simmhan
CLOUD2
2019 Dynamic Scaling of Video Analytics for Wide-Area Tracking in Urban Spaces
abstract
Smart City deployments typically have thousands to even hundreds of thousands of Surveillance cameras. Rapid advancements in computer vision techniques due to Deep Neural Networks enable using these camera feeds for performing non-trivial analytics. Tracking a moving object of interest using a large network of cameras, also known as object reidentification, is one such analytic that empowers city administration with capabilities such as finding missing people or prioritizing emergency vehicles. We have built Anveshak, a framework for distributed wide-area tracking. Anveshak fills in the shortcomings of existing Big Data and Deep Learning frameworks by - exposing an intuitive and composable programming model; automating application deployment and orchestration across edge, fog and cloud resources and providing knobs to the user for managing the application performance. The knobs lend the application the ability to scale potentially to thousands of cameras. In this proposal we have designed two representative applications; missing person tracking and priority signalling for emergency vehicles. We empirically verify that the application scales to 1000 cameras on a Cloud-only deployment of 10 Azure VMs with 8 cores and 32GB RAM each. Alternatively, it scales to 500 cameras on a simulated setup of 100 edge, 30 fog, and 1 Cloud VM. We also highlight the effect of the knobs on the application performance. designed two representative applications; missing person tracking and priority signalling for emergency vehicles. We empirically verify that the application scales to 1000 cameras on a Cloud-only deployment of 10 Azure VMs with 8 cores and 32GB RAM each. Alternatively, it scales to 500 cameras on a simulated setup of 100 edge, 30 fog, and 1 Cloud VM. We also highlight the effect of the knobs on the application performance.
Aakash Khochare, Sheshadri K. R, Shriram R., Yogesh L. Simmhan
CCGRID4
2019 SATVAM: Toward an IoT Cyber-Infrastructure for Low-Cost Urban Air Quality Monitoring
abstract
Air pollution is a public health emergency in large cities. The availability of commodity sensors and the advent of Internet of Things (IoT) enable the deployment of a city-wide network of 1000's of low-cost real-time air quality monitors to help manage this challenge. This needs to be supported by an IoT cyber-infrastructure for reliable and scalable data acquisition from the edge to the Cloud. The low accuracy of such sensors also motivates the need for data-driven calibration models that can accurately predict the science variables from the raw sensor signals. Here, we offer our experiences with designing and deploying such an IoT software platform and calibration models, and validate it through a pilot field deployment at two mega-cities, Delhi and Mumbai. Our edge data service is able to even-out the differential bandwidths from the sensing devices and to the Cloud repository, and recover from transient failures. Our analytical models reduce the errors of the sensors from a best-case of 63% using the factory baseline to as low as 21%, and substantially advances the state-of-the-art in this domain.
Yogesh L. Simmhan, Malati Hegde, Rajesh Zele, Sachchida N. Tripathi, Srijith Nair, Sumit K. Monga, Ravi Sahu, Kuldeep Dixit, Ronak Sutaria, Brijesh Mishra, Anamika Sharma, S. V. R. Anand
eScience1
2019 ElfStore: A Resilient Data Storage Service for Federated Edge and Fog Resources
abstract
Edge and fog computing have grown popular as IoT deployments become wide-spread. While application composition and scheduling on such resources are being explored, there exists a gap in a distributed data storage service on the edge and fog layer, instead depending solely on the cloud for data persistence. Such a service should reliably store and manage data on fog and edge devices, even in the presence of failures, and offer transparent discovery and access to data for use by edge computing applications. Here, we present ElfStore, a first-of-its-kind edge-local federated store for streams of data blocks. It uses reliable fog devices as a super-peer overlay to monitor the edge resources, offers federated metadata indexing using Bloom filters, locates data within 2-hops, and maintains approximate global statistics about the reliability and storage capacity of edges. Edges host the actual data blocks, and we use a unique differential replication scheme to select edges on which to replicate blocks, to guarantee a minimum reliability and to balance storage utilization. Our experiments on two IoT virtual deployments with 20 and 272 devices show that ElfStore has low overheads, is bound only by the network bandwidth, has scalable performance, and offers tunable resilience.
Sumit K. Monga, Sheshadri K. R, Yogesh L. Simmhan
ICWS3
2019 Toward Resilient Stream Processing on Clouds using Moving Target Defense
abstract
Big data platforms have grown popular for real-time stream processing on distributed clusters and clouds. However, execution of sensitive streaming applications on shared computing resources increases their vulnerabilities, and may lead to data leaks and injection of spurious logic that can compromise these applications. Here, we adopt Moving Target Defense (MTD) techniques into Fast Data platforms, and propose MTD strategies by which we can mitigate these attacks. Our strategies target the platform, application and data layers, which make these reusable, rather than the OS, virtual machine, or hardware layers, which are environment specific. We use Apache Storm as the canonical distributed stream processing platform for designing our MTD strategies, and offer a preliminary evaluation that indicates the feasibility and evaluates the performance overheads.
Shilpa Chaturvedi, Yogesh L. Simmhan
ISORC2
2019 AutoBoT: Resilient and Cost-Effective Scheduling of a Bag of Tasks on Spot VMs
abstract
Many data and task parallel applications can be modeled as a Bag of Tasks (BoT), and scheduled on distributed systems such as Grids, Clusters, and Clouds. We propose AutoBoT, a collection of scheduling strategies for BoTs with hard deadlines on Cloud Virtual Machines (VMs), to lower the overall monetary cost - a distinctive factor for Clouds. Besides reliable fixed-price VMs, AutoBoT uniquely reduces costs by including preemptible spot-priced VMs that are much cheaper, but are unreliable and have time-variant pricing. It guarantees timely completion by making active runtime decisions on pricing, number of VMs to acquire/release, and on task placement, checkpointing and migration. Our rigorous simulations of 7 Million BoT runs sampled from the Google cluster workload uses a realistic Cloud model and 6 months of Amazon EC2 pricing data to compare AutoBoT against two baseline algorithms. We analyze the impact of BoT size, data centers, time periods, deadline duration, loss budget and checkpointing strategies. AutoBoT often gives X80% profit and rare but bounded losses, compared to using only fixed-price VMs. Further, its 100 percent completion guarantee is 23-42 percent better than using only spot-priced VMs which offer a similar profit.
Prateeksha Varshney, Yogesh L. Simmhan
IEEE Trans. Parallel Distributed Syst.2
2018 Adaptive Energy-Aware Scheduling of Dynamic Event Analytics Across Edge and Cloud Resources
abstract
The growing deployment of sensors as part of Internet of Things (IoT) is generating thousands of event streams. Complex Event Processing (CEP) queries offer a useful paradigm for rapid decision-making over such data sources. While often centralized in the Cloud, the deployment of capable edge devices on the field motivates the need for cooperative event analytics that span Edge and Cloud computing. Here, we identify a novel problem of query placement on edge and Cloud resources for dynamically arriving and departing analytic dataflows. We define this as an optimization problem to minimize the total makespan for all event analytics, while meeting energy and compute constraints of the resources. We propose 4 adaptive heuristics and 3 rebalancing strategies for such dynamic dataflows, and validate them using detailed simulations for 100 - 1000 edge devices and VMs. The results show that our heuristics offer O(seconds) planning time, give a valid and high quality solution in all cases, and reduce the number of query migrations. Furthermore, rebalance strategies when applied in these heuristics have significantly reduced the makespan by around 20 - 25%.
Rajrup Ghosh, Siva Prakash Reddy Komma, Yogesh L. Simmhan
CCGrid3
2018 VIoLET: A Large-Scale Virtual Environment for Internet of Things
Shreyas Badiger, Shrey Baheti, Yogesh L. Simmhan
Euro-Par3
2018 Toward Reliable and Rapid Elasticity for Streaming Dataflows on Clouds
abstract
The pervasive availability of streaming data is driving Fast Data platforms for low-latency streaming applications. Such applications need to respond to dynamism in the input rates and task behavior using scale-in and -out on elastic Cloud resources. Platforms like Apache Storm do not provide robust means to respond to such dynamism and for rapid task migration across VMs. We propose several dataflow checkpoint and migration approaches that allow a running streaming dataflow to migrate, without any loss of in-flight messages or their internal tasks states, while reducing the time to recover and stabilize. We implement these strategies on Storm and evaluate them using micro and application dataflows for scaling in and out on 2 - 21 Cloud VMs. Our results show that we can migrate large dataflows and catchup with their processing 75% faster than Storm, which takes over 140 secs. We also find that our approaches stabilize the application up to 42% faster, and there is no failure and re-processing of messages.
Anshu Shukla, Yogesh L. Simmhan
ICDCS2
2018 Model-driven scheduling for distributed stream processing systems
Anshu Shukla, Yogesh L. Simmhan
J. Parallel Distributed Comput.2
2018 Cover Image
Yogesh L. Simmhan, Pushkara Ravindra, Shilpa Chaturvedi, Malati Hegde, Rashmi Ballamajalu
Softw. Pract. Exp.1
2018 Towards a data-driven IoT software architecture for smart city utilities
abstract
Summary The Internet of things (IoT) is emerging as the next big wave of digital presence for billions of devices on the Internet. Smart cities are a practical manifestation of IoT, with the goal of efficient, reliable, and safe delivery of city utilities like water, power, and transport to residents, through their intelligent management. A data‐driven IoT software platform is essential for realizing manageable and sustainable smart utilities and for novel applications to be developed upon them. Here, we propose such service‐oriented software architecture to address 2 key operational activities in a smart utility: the IoT fabric for resource management and the data and application platform for decision‐making . Our design uses Open Web standards and evolving network protocols, cloud and edge resources, and streaming big data platforms. We motivate our design requirements using the smart water management domain; some of these requirements are unique to developing nations. We also validate the architecture within a campus‐scale IoT testbed at the Indian Institute of Science, Bangalore and present our experiences. Our architecture is scalable to a township or city while also generalizable to other smart utility domains. Our experiences serve as a template for other similar efforts, particularly in emerging markets and highlight the gaps and opportunities for a data‐driven IoT software architecture for smart cities.
Yogesh L. Simmhan, Pushkara Ravindra, Shilpa Chaturvedi, Malati Hegde, Rashmi Ballamajalu
Softw. Pract. Exp.1
2018 Distributed Scheduling of Event Analytics across Edge and Cloud
abstract
Internet of Things (IoT) domains generate large volumes of high-velocity event streams from sensors, which need to be analyzed with low latency to drive decisions. Complex Event Processing (CEP) is a Big Data technique to enable such analytics and is traditionally performed on Cloud Virtual Machines (VM). Leveraging captive IoT edge resources in combination with Cloud VMs can offer better performance, flexibility, and monetary costs for CEP. Here, we formulate an optimization problem for energy-aware placement of CEP queries , composed as an analytics dataflow, across a collection of edge and Cloud resources, with the goal of minimizing the end-to-end latency for the dataflow. We propose a Genetic Algorithm (GA) meta-heuristic to solve this problem and compare it against a brute-force optimal algorithm (BF). We perform detailed real-world benchmarks on the compute, network, and energy capacity of edge and Cloud resources. These results are used to define a realistic and comprehensive simulation study that validates the BF and GA solutions for 45 diverse CEP dataflows, LAN and WAN setup, and different edge resource availability. We compare the GA and BF solutions against random and Cloud-only baselines for different configurations for a total of 1,764 simulation runs. Our study shows that GA is within 97% of the optimal BF solution that takes hours, maps dataflows with 4--50 queries in 1--26s, and only fails to offer a feasible solution ≤20% of the time.
Rajrup Ghosh, Yogesh L. Simmhan
ACM Trans. Cyber Phys. Syst.2
2017 Collaborative Reuse of Streaming Dataflows in IoT Applications
abstract
Distributed Stream Processing Systems (DSPS) like Apache Storm and Spark Streaming enable composition of continuous dataflows that execute persistently over data streams. They are used by Internet of Things (IoT) applications to analyze sensor data from Smart City cyber-infrastructure, and make active utility management decisions. As the ecosystem of such IoT applications that leverage shared urban sensor streams continue to grow, applications will perform duplicate pre-processing and analytics tasks. This offers the opportunity to collaboratively reuse the outputs of overlapping dataflows, thereby improving the resource efficiency. In this paper, we propose dataflow reuse algorithms that given a submitted dataflow, identifies the intersection of reusable tasks and streams from a collection of running dataflows to form a merged dataflow. Similar algorithms to unmerge dataflows when they are removed are also proposed. We implement these algorithms for the popular Apache Storm DSPS, and validate their performance and resource savings for 35 synthetic dataflows based on public OPMW workflows with diverse arrival and departure distributions, and on 21 real IoT dataflows from RIoTBench. We see that our Reuse algorithms reduce the count of running tasks by 38 - 46% for the two workloads, and a reduction in cumulative CPU usage of 36-51%, that can result in real cost savings on Cloud resources.
Shilpa Chaturvedi, Sahil Tyagi, Yogesh L. Simmhan
eScience3
2017 ARM Wrestling with Big Data: A Study of Commodity ARM64 Server for Big Data Workloads
abstract
ARM processors have dominated the mobile device market in the last decade due to their favorable computing to energy ratio. In this age of Cloud data centers and Big Data analytics, the focus is increasingly on power efficient processing, rather than just high throughput computing. ARM's first commodity server-grade processor is the recent AMD A1100-series processor, based on a 64-bit ARM Cortex A57 architecture. In this paper, we study the performance and energy efficiency of a server based on this ARM64 CPU, relative to a comparable server running an AMD Opteron 3300-series x64 CPU, for Big Data workloads. Specifically, we study these for Intel's HiBench suite of web, query and machine learning benchmarks on Apache Hadoop v2.7 in a pseudo-distributed setup, for data sizes up to 20GB files, 5M web pages and 500M tuples. Our results show that the ARM64 server's runtime performance is comparable to the x64 server for integer-based workloads like Sort and Hive queries, and only lags behind for floating-point intensive benchmarks like PageRank, when they do not exploit data parallelism adequately. We also see that the ARM64 server takes 1/3rd the energy, and has an Energy Delay Product (EDP) that is 50-71% lower than the x64 server. These results hold promise for ARM64 data centers hosting Big Data workloads to reduce their operational costs, while opening up opportunities for further analysis.
Jayanth Kalyanasundaram, Yogesh L. Simmhan
HiPC2
2017 Demystifying Fog Computing: Characterizing Architectures, Applications and Abstractions
abstract
Internet of Things (IoT) has accelerated the deployment of millions of sensors at the edge of the network, through Smart City infrastructure and lifestyle devices. Cloud computing platforms are often tasked with handling these large volumes and fast streams of data from the edge. Recently, Fog computing has emerged as a concept for low-latency and resource-rich processing of these observation streams, to complement Edge and Cloud computing. In this paper, we review various dimensions of system architecture, application characteristics and platform abstractions that are manifest in this Edge, Fog and Cloud eco-system. We highlight novel capabilities of the Edge and Fog layers, such as physical and application mobility, privacy sensitivity, and a nascent runtime environment. IoT application case studies based on first-hand experiences across diverse domains drive this categorization. We also highlight the gap between the potential and the reality of Fog computing, and identify challenges that need to be overcome for the solution to be sustainable. Taken together, our article can help platform and application developers bridge the gap that remains in making Fog computing viable.
Prateeksha Varshney, Yogesh L. Simmhan
ICFEC2
2017 \mathbb ECHO : An Adaptive Orchestration Platform for Hybrid Dataflows across Cloud and Edge
Pushkara Ravindra, Aakash Khochare, Sivaprakash Reddy, Sarthak Sharma, Prateeksha Varshney, Yogesh L. Simmhan
ICSOC6
2017 Introducing distributed dynamic data-intensive (D3) science: Understanding applications and infrastructure
abstract
Summary A common feature across many science and engineering applications is the amount and diversity of data and computation that must be integrated to yield insights. Datasets are growing larger and becoming distributed; their location, availability, and properties are often time‐dependent. Collectively, these characteristics give rise to dynamic distributed data‐intensive applications. While “static” data applications have received significant attention, the characteristics, requirements, and software systems for the analysis of large volumes of dynamic, distributed data, and data‐intensive applications have received relatively less attention. This paper surveys several representative dynamic distributed data‐intensive application scenarios, provides a common conceptual framework to understand them, and examines the infrastructure used in support of applications.
Shantenu Jha, Daniel S. Katz, André Luckow, Neil P. Chue Hong, Omer F. Rana, Yogesh L. Simmhan
Concurr. Comput. Pract. Exp.6
2017 RIoTBench: An IoT benchmark for distributed stream processing systems
abstract
Summary The Internet of Things (IoT) is an emerging technology paradigm where millions of sensors and actuators help monitor and manage physical, environmental, and human systems in real time. The inherent closed‐loop responsiveness and decision making of IoT applications make them ideal candidates for using low latency and scalable stream processing platforms. Distributed stream processing systems (DSPS) hosted in cloud data centers are becoming the vital engine for real‐time data processing and analytics in any IoT software architecture. But the efficacy and performance of contemporary DSPS have not been rigorously studied for IoT applications and data streams. Here, we propose RIoTBench , a real‐time IoT benchmark suite, along with performance metrics, to evaluate DSPS for streaming IoT applications. The benchmark includes 27 common IoT tasks classified across various functional categories and implemented as modular microbenchmarks. Further, we define four IoT application benchmarks composed from these tasks based on common patterns of data preprocessing, statistical summarization, and predictive analytics that are intrinsic to the closed‐loop IoT decision‐making life cycle. These are coupled with four stream workloads sourced from real IoT observations on smart cities and smart health, with peak streams rates that range from 500 to 10 000 messages/second from up to 3 million sensors. We validate the RIoTBench suite for the popular Apache Storm DSPS on the Microsoft Azure public cloud and present empirical observations. This suite can be used by DSPS researchers for performance analysis and resource scheduling, by IoT practitioners to evaluate DSPS platforms, and even reused within IoT solutions.
Anshu Shukla, Shilpa Chaturvedi, Yogesh L. Simmhan
Concurr. Comput. Pract. Exp.3
2017 Knowledge-infused and consistent Complex Event Processing over real-time and persistent streams
Qunzhi Zhou, Yogesh L. Simmhan, Viktor Prasanna 0001
Future Gener. Comput. Syst.2
2016 A meta-graph approach to analyze subgraph-centric distributed programming models
abstract
Component-centric distributed graph processing models that use bulk synchronous parallel (BSP) execution have grown popular. These overcome short-comings of Big Data platforms like Hadoop for processing large graphs. However, literature on formal analysis of these component-centric abstractions for different graphs, graph partitioning, and graph algorithms is lacking. Here, we propose an coarse-grained analytical approach based on a meta-graph sketch to examine the characteristics of component-centric graph programming models. We apply this sketch to subgraph- and block-centric abstractions, and draw a comparison with vertex-centric models like Google's Pregel. We explore the impact of various graph partitioning techniques on the meta-graph, and the impact of the meta-graph on graph algorithms. This decouples large unwieldy graphs and their partitioning artifacts from their algorithmic analysis. We evaluate our approach for five spatial and powerlaw graphs, four different partitioning strategies, and for PageRank and Breadth First Search algorithms. We show that this novel analytical technique is simple, scalable and yet gives a reliable estimate of the number of supersteps, and the communication and computational complexities of the algorithms for various graphs.
Ravikant Dindokar, Neel Choudhury, Yogesh L. Simmhan
IEEE BigData3
2016 Elastic Partition Placement for Non-stationary Graph Algorithms
abstract
Distributed graph platforms like Pregel have usedvertex-centric programming models to process the growing cor-pus of graph datasets using commodity clusters. However, theirregular structure of graphs causes load imbalances acrossmachines, and this is exacerbated for non-stationary graphalgorithms where not all parts of the graph are active at thesame time. As a result, such graph platforms do not make efficientuse of distributed resources. In this paper, we decouple graphpartitioning from placement on hosts, and introduce strategiesfor elastic placement of graph partitions on Cloud VMs to reducethe cost of execution compared to a static placement, even aswe minimize the increase in makespan. These strategies areinnovative in modeling the graph algorithm's non-stationarybehavior a priori using a metagraph sketch. We validate ourstrategies for several real-world graphs, using runtime tracesfor approximate Betweenness Centrality (BC) algorithm on oursubgraph-centric GoFFish graph platform. Our strategies areable to reduce the cost of execution by up to 54%, comparedto a static placement, while achieving a makespan that is within25% of the optimal.
Ravikant Dindokar, Yogesh L. Simmhan
CCGrid2
2016 GoDB: From Batch Processing to Distributed Querying over Property Graphs
abstract
Property Graphs with rich attributes over vertices and edges are becoming common. Querying and mining such linked Big Data is important for knowledge discovery and mining. Distributed graph platforms like Pregel focus on batch execution on commodity clusters. But exploratory analytics requires platforms that are both responsive and scalable. We propose Graph-oriented Database (GoDB), a distributed graph database that supports declarative queries over large property graphs. GoDB builds upon our GoFFish subgraph-centric batch processing platform, leveraging its scalability while using execution heuristics to offer responsiveness. The GoDB declarative query model supports vertex, edge, path and reachability queries, and this is translated to a distributed execution plan on GoFFish. We also propose a novel cost model to choose a query plan that minimizes the execution latency. We evaluate GoDB deployed on the Azure IaaS Cloud, over real-world property graphs and for a diverse workload of 500 queries. These show that the cost model selects the optimal execution plan at least 80% of the time, and helps GoDB weakly scale with the graph size. A comparative study with Titan, a leading open-source graph database, shows that we complete all queries, each in ≤ 1.6 secs, while Titan cannot complete up to 42% of some query workloads.
Nitin Jamadagni, Yogesh L. Simmhan
CCGrid2
2016 Cloud computing for data-driven science and engineering
abstract
Cloud computing for data-driven science and engineering* During the past decade, data-driven science and engineering have emerged as a key paradigm for performing scientific research, enabling innovations through new kinds of experiments that were earlier impossible.Today's science has access to advanced instruments like next generation genome sequencers, gigapixel survey telescopes, and networks of sensors that monitor cyber-physical systems, and these are generating datasets that are growing exponentially in complexity and data volume.Big Data, across all dimensions of volume, velocity, variety, and veracity, are offering unique opportunities to enable scientific discovery as well as novel challenges to scientific platforms.Such dynamic, distributed, and data-intensive applications hold the solutions to vital scientific and societal problems of the 21st century.In order to achieve breakthrough in new knowledge, there is a need to develop data-driven system models, perform analytics at large scales, manage data from instruments and analyses, and share and visualize the results with scientific peers and the society at large.To this end, cloud computing offers a computing model for running such data-intensive scientific and engineering applications.Clouds have democratized resource access to underserved disciplines, making it possible to perform nontrivial scientific explorations for just a few hundred dollars.Clouds are particularly cost-effective for Big Data applications due to their co-location of elastic compute resources with data, and their use of commodity hardware, which economizes on costs for non-high performance computing (HPC) workloads.Many contemporary Big Data platforms that have emerged from online enterprises such as Google and Twitter are also optimized for such commodity hardware, as found in their own data centers.Of course, there are costs associated with data transfer and storage, in keeping with the pay-as-you-go model, that may not be well suited for applications requiring frequent transfer of and long-term storage of terabytes of data.Likewise, it is valuable to understand how data-intensive or even HPC applications that have been developed for computing grids, at one end, and applications developed for workstation tools like MATLAB and R, at the other end, can be effectively run on clouds.These are some of the practical realities that are worth exploring on the relevance of clouds for data-driven scientific applications.In this special issue, we have compiled a set of articles that discusses new research, development, and deployment efforts in running eScience and eEngineering workloads on cloud infrastructures and platforms.The open solicitation, which followed the 3rd Workshop on Scientific Cloud Computing (ScienceCloud), invited research and case studies on a variety of topics relevant to data-driven scientific computing on clouds: use of cloud-based technologies to address innovative compute and data-driven scientific problems that are not well served by current HPC clusters and grids, programming platforms for elastic and Big Data applications, performance and cost-effective computing on clouds, and gaps in diverse cloud fabrics and service offerings, among others.In all, the special issue received 28 articles, of which six were selected for publication after multiple rounds of reviews and revisions.The special issue starts with two articles that explore the runtime platform support required for executing Big Data science on clouds.In TomusBlobs: Scalable Data-Intensive Processing on Azure Clouds [1], the authors address the limitations of data storage within IaaS clouds such as Amazon S3 and Microsoft Azure BLOBs that are, while co-located in the data center, not present in the virtual machines (VMs) and need to be accessed over the network.Their distributed storage on VMs is optimized for concurrent access and elastic scaling, even *Corrections added on 5 February 2016, after first online publication: references to the paper "Pilot-abstractions for distributed data-intensive cloud applications" have been removed.
Yogesh L. Simmhan, Lavanya Ramakrishnan, Gabriel Antoniu, Carole A. Goble
Concurr. Comput. Pract. Exp.1
2015 Fault-Tolerant and Elastic Streaming MapReduce with Decentralized Coordination
abstract
The MapReduce programming model, due to its simplicity and scalability, has become an essential tool for processing large data volumes in distributed environments. Recent Stream Processing Systems (SPS) this model to provide low-latency analysis of high-velocity continuous data streams. However, integrating MapReduce with streaming poses challenges: first, the runtime variations in data characteristics such as data-rates and key-distribution cause resource overload, that in-turn leads to fluctuations in the Quality of the Service (QoS), and second, the stateful reducers, whose state depends on the complete tuple history, necessitates efficient fault-recovery mechanisms to maintain the desired QoS in the presence of resource failures. We propose an integrated streaming MapReduce architecture leveraging the concept of consistent hashing to support runtime elasticity along with locality-aware data and state replication to provide efficient load-balancing with low-overhead fault-tolerance and parallel fault-recovery from multiple simultaneous failures. Our evaluation on a private cloud shows up to 2.8× improvement in peak throughput compared to Apache Storm SPS, and a low recovery latency of 700 - 1500 ms from multiple failures.
Alok Gautam Kumbhare, Marc Frîncu, Yogesh L. Simmhan, Viktor Prasanna 0001
ICDCS3
2015 Distributed Programming over Time-Series Graphs
abstract
Graphs are a key form of Big Data, and performing scalable analytics over them is invaluable to many domains. There is an emerging class of inter-connected data which accumulates or varies over time, and on which novel algorithms both over the network structure and across the time-variant attribute values is necessary. We formalize the notion of time-series graphs and propose a Temporally Iterative BSP programming abstraction to develop algorithms on such datasets using several design patterns. Our abstractions leverage a sub-graph centric programming model and extend it to the temporal dimension. We present three time-series graph algorithms based on these design patterns and abstractions, and analyze their performance using the Offish distributed platform on Amazon AWS Cloud. Our results demonstrate the efficacy of the abstractions to develop practical time-series graph algorithms, and scale them on commodity hardware.
Yogesh L. Simmhan, Neel Choudhury, Charith Wickramaarachchi, Alok Gautam Kumbhare, Marc Frîncu, Cauligi S. Raghavendra, Viktor Prasanna 0001
IPDPS1
2015 Editorial: Scalable Systems for Big Data Management and Analytics
Srinivas Aluru, Yogesh L. Simmhan
J. Parallel Distributed Comput.2
2015 Reactive Resource Provisioning Heuristics for Dynamic Dataflows on Cloud Infrastructure
abstract
The need for low latency analysis over high-velocity data streams motivates the need for distributed continuous dataflow systems. Contemporary stream processing systems use simple techniques to scale on elastic cloud resources to handle variable data rates. However, application QoS is also impacted by variability in resource performance exhibited by clouds and hence necessitates autonomic methods of provisioning elastic resources to support such applications on cloud infrastructure. We develop the concept of “dynamic dataflows” which utilize alternate tasks as additional control over the dataflow's cost and QoS. Further, we formalize an optimization problem to represent deployment and runtime resource provisioning that allows us to balance the application's QoS, value, and the resource cost. We propose two greedy heuristics, centralized and sharded, based on the variable-sized bin packing algorithm and compare against a Genetic Algorithm (GA) based heuristic that gives a near-optimal solution. A large-scale simulation study, using the linear road benchmark and VM performance traces from the AWS public cloud, shows that while GA-based heuristic provides a better quality schedule, the greedy heuristics are more practical, and can intelligently utilize cloud elasticity to mitigate the effect of variability, both in input data rates and cloud resource performance, to meet the QoS of fast data applications.
Alok Gautam Kumbhare, Yogesh L. Simmhan, Marc Frîncu, Viktor Prasanna 0001
IEEE Trans. Cloud Comput.2
2015 Holistic Measures for Evaluating Prediction Models in Smart Grids
abstract
The performance of prediction models is often based on “abstract metrics” that estimate the model's ability to limit residual errors between the observed and predicted values. However, meaningful evaluation and selection of prediction models for end-user domains requires holistic and application-sensitive performance measures. Inspired by energy consumption prediction models used in the emerging “big data” domain of Smart Power Grids, we propose a suite of performance measures to rationally compare models along the dimensions of scale independence, reliability, volatility and cost. We include both application independent and dependent measures, the latter parameterized to allow customization by domain experts to fit their scenario. While our measures are generalizable to other domains, we offer an empirical analysis using real energy use data for three Smart Grid applications: planning, customer education and demand response, which are relevant for energy sustainability. Our results underscore the value of the proposed measures to offer a deeper insight into models' behavior and their impact on real applications, which benefit both data mining researchers and practitioners.
Saima Aman, Yogesh L. Simmhan, Viktor Prasanna 0001
IEEE Trans. Knowl. Data Eng.2
2014 PLAStiCC: Predictive Look-Ahead Scheduling for Continuous Dataflows on Clouds
abstract
Scalable stream processing and continuous dataflow systems are gaining traction with the rise of big data due to the need for processing high velocity data in near real time. Unlike batch processing systems such as MapReduce and workflows, static scheduling strategies fall short for continuous data flows due to the variations in the input data rates and the need for sustained throughput. The elastic resource provisioning of cloud infrastructure is valuable to meet the changing resource needs of such continuous applications. However, multi-tenant cloud resources introduce yet another dimension of performance variability that impacts the application's throughput. In this paper we propose Plastic, an adaptive scheduling algorithm that balances resource cost and application throughput using a prediction-based look-ahead approach. It not only addresses variations in the input data rates but also the underlying cloud infrastructure. In addition, we also propose several simpler static scheduling heuristics that operate in the absence of accurate performance prediction model. These static and adaptive heuristics are evaluated through extensive simulations using performance traces obtained from Amazon AWS IaaS public cloud. Our results show an improvement of up to 20% in the overall profit as compared to the reactive adaptation algorithm.
Alok Gautam Kumbhare, Yogesh L. Simmhan, Viktor Prasanna 0001
CCGRID2
2014 GoFFish: A Sub-graph Centric Framework for Large-Scale Graph Analytics
Yogesh L. Simmhan, Alok Gautam Kumbhare, Charith Wickramaarachchi, Soonil Nagarkar, Santosh Ravi, Cauligi S. Raghavendra, Viktor Prasanna 0001
Euro-Par1
2014 Cost-Efficient and Resilient Job Life-Cycle Management on Hybrid Clouds
abstract
Cloud infrastructure offers democratized access to on-demand computing resources for scaling applications beyond captive local servers. While on-demand, fixed-price Virtual Machines (VMs) are popular, the availability of cheaper, but less reliable, spot VMs from cloud providers presents an opportunity to reduce the cost of hosting cloud applications. Our work addresses the issue of effective and economic use of hybrid cloud resources for planning job executions with deadline constraints. We propose strategies to manage a job's life-cycle on spot and on on-demand VMs to minimize the total dollar cost while assuring completion. With the foundation of stochastic optimization, our reusable table-based algorithm (RTBA) decides when to instantiate VMs, at what bid prices, when to use local machines, and when to checkpoint and migrate the job between these resources, with the goal of completing the job on time and with the minimum cost. In addition, three simpler heuristics are proposed as comparison. Our evaluation using historical spot prices for the Amazon EC2 market shows that RTBA on an average reduces the cost by 72%, compared to running on only on-demand VMs. It is also robust to fluctuations in spot prices. The heuristic, H3, often approaches RTBA in performance and may prove adequate for ad hoc jobs due to its simplicity.
Hsuan-Yi Chu, Yogesh L. Simmhan
IPDPS2
2013 Scalable prediction of energy consumption using incremental time series clustering
abstract
Time series datasets are a canonical form of high velocity Big Data, and often generated by pervasive sensors, such as found in smart infrastructure. Performing predictive analytics on time series data can be computationally complex, and requires approximation techniques. In this paper, we motivate this problem using a real application from the smart grid domain. We propose an incremental clustering technique, along with a novel affinity score for determining cluster similarity, which help reduce the prediction error for cumulative time series within a cluster. We evaluate this technique, along with optimizations, using real datasets from smart meters, totaling ~700,000 data points, and show the efficacy of our techniques in improving the prediction error of time series data within polynomial time.
Yogesh L. Simmhan, Muhammad Usman Noor
IEEE BigData1
2013 Towards hybrid online on-demand querying of realtime data with stateful complex event processing
abstract
Emerging Big Data applications in areas like ecommerce and energy industry require both online and on-demand queries to be performed over vast and fast data arriving as streams. These present novel challenges to Big Data management systems. Complex Event Processing (CEP) is recognized as a high performance online query scheme which in particular deals with the velocity aspect of the 3-V's of Big Data. However, traditional CEP systems do not consider data variety and lack the capability to embed ad hoc queries over the volume of data streams. In this paper, we propose H2O, a stateful complex event processing framework, to support hybrid online and on-demand queries over realtime data. We propose a semantically enriched event and query model to address data variety. A formal query algebra is developed to precisely capture the stateful and containment semantics of online and on-demand queries. We describe techniques to achieve the interactive query processing over realtime data featured by efficient online querying, dynamic stream data persistence and on-demand access. The system architecture is presented and the current implementation status reported.
Qunzhi Zhou, Yogesh L. Simmhan, Viktor Prasanna 0001
IEEE BigData2
2013 Continuous Dataflow Update Strategies for Mission-Critical Applications
abstract
Continuous data flows complement scientific work-flows by allowing composition of real time data ingest and analytics pipelines to process data streams from pervasive sensors and "always-on" scientific instruments. Such data flows are mission-critical applications that cannot suffer downtime, need to operate consistently, and are long running, but may need to be updated to fix bugs or add features. This poses the problem: How do we update the continuous dataflow application with minimal disruption? In this paper, we formalize different types of dataflow update models for continuous dataflow applications, and identify the qualitative and quantitative metrics to be considered when choosing an update strategy. We propose five dataflow update strategies, and analytically characterize their performance trade-offs. We validate one of these consistent, low-latency update strategies using the Floe dataflow engine for an eEngineering application from the Smart Power Grid domain, and show its relative performance benefits against a naïve update strategy.
Charith Wickramaarachchi, Yogesh L. Simmhan
e-Science2
2013 Optimizations and Analysis of BSP Graph Processing Models on Public Clouds
abstract
Large-scale graph analytics is a central tool in many fields, and exemplifies the size and complexity of Big Data applications. Recent distributed graph processing frameworks utilize the venerable Bulk Synchronous Parallel (BSP) model and promise scalability for large graph analytics. This has been made popular by Google's Pregel, which provides an architecture design for BSP graph processing. Public clouds offer democratized access to medium-sized compute infrastructure with the promise of rapid provisioning with no capital investment. Evaluating BSP graph frameworks on cloud platforms with their unique constraints is less explored. Here, we present optimizations and analyses for computationally complex graph analysis algorithms such as betweenness-centrality and all-pairs shortest paths on a native BSP framework we have developed for the Microsoft Azure Cloud, modeled on the Pregel graph processing model. We propose novel heuristics for scheduling graph vertex processing in swaths to maximize resource utilization on cloud VMs that lead to a 3.5x performance improvement. We explore the effects of graph partitioning in the context of BSP, and show that even a well partitioned graph may not lead to performance improvements due to BSP's barrier synchronization. We end with a discussion on leveraging cloud elasticity for dynamically scaling the number of BSP workers to achieve a better performance than a static deployment, and at a significantly lower cost.
Mark Redekopp, Yogesh L. Simmhan, Viktor Prasanna 0001
IPDPS2
2013 Exploiting application dynamism and cloud elasticity for continuous dataflows
abstract
Contemporary continuous dataflow systems use elastic scaling on distributed cloud resources to handle variable data rates and to meet applications' needs while attempting to maximize resource utilization. However, virtualized clouds present an added challenge due to the variability in resource performance -- over time and space -- thereby impacting the application's QoS. Elastic use of cloud resources and their allocation to continuous dataflow tasks need to adapt to such infrastructure dynamism. In this paper, we develop the concept of "dynamic dataflows" as an extension to continuous dataflows that utilizes alternate tasks and allows additional control over the dataflow's cost and QoS. We formalize an optimization problem to perform both deployment and runtime cloud resource management for such dataflows, and define an objective function that allows trade-off between the application's value against resource cost. We present two novel heuristics, local and global, based on the variable sized bin packing heuristics to solve this NP-hard problem. We evaluate the heuristics against a static allocation policy for a dataflow with different data rate profiles that is simulated using VM performance traces from a private cloud data center. The results show that the heuristics are effective in intelligently utilizing cloud elasticity to mitigate the effect of both input data rate and cloud resource performance variabilities on QoS.
Alok Gautam Kumbhare, Yogesh L. Simmhan, Viktor Prasanna 0001
SC2
2012 Cryptonite: A Secure and Performant Data Repository on Public Clouds
abstract
Cloud storage has become immensely popular for maintaining synchronized copies of files and for sharing documents with collaborators. However, there is heightened concern about the security and privacy of Cloud-hosted data due to the shared infrastructure model and an implicit trust in the service providers. Emerging needs of secure data storage and sharing for domains like Smart Power Grids, which deal with sensitive consumer data, require the persistence and availability of Cloud storage but with client-controlled security and encryption, low key management overhead, and minimal performance costs. Cryptonite is a secure Cloud storage repository that addresses these requirements using a Strongbox model for shared key management. We describe the Cryptonite service and desktop client, discuss performance optimizations, and provide an empirical analysis of the improvements. Our experiments shows that Cryptonite clients achieve a 40% improvement in file upload bandwidth over plaintext storage using the Azure Storage Client API despite the added security benefits, while our file download performance is 5 times faster than the baseline for files greater than 100MB.
Alok Gautam Kumbhare, Yogesh L. Simmhan, Viktor Prasanna 0001
IEEE CLOUD2
2012 Incorporating Semantic Knowledge into Dynamic Data Processing for Smart Power Grids
Qunzhi Zhou, Yogesh L. Simmhan, Viktor Prasanna 0001
ISWC (2)2
2011 An Analysis of Security and Privacy Issues in Smart Grid Software Architectures on Clouds
abstract
Power utilities globally are increasingly upgrading to Smart Grids that use bi-directional communication with the consumer to enable an information-driven approach to distributed energy management. Clouds offer features well suited for Smart Grid software platforms and applications, such as elastic resources and shared services. However, the security and privacy concerns inherent in an information-rich Smart Grid environment are further exacerbated by their deployment on Clouds. Here, we present an analysis of security and privacy issues in a Smart Grids software architecture operating on different Cloud environments, in the form of a taxonomy. We use the Los Angeles Smart Grid Project that is underway in the largest U.S. municipal utility to drive this analysis that will benefit both Cloud practitioners targeting Smart Grid applications, and Cloud researchers investigating security and privacy.
Yogesh L. Simmhan, Alok Gautam Kumbhare, Baohua Cao, Viktor Prasanna 0001
IEEE CLOUD1
2011 Towards Reliable, Performant Workflows for Streaming-Applications on Cloud Platforms
abstract
Scientific workflows are commonplace in eScience applications. Yet, the lack of integrated support for data models, including streaming data, structured collections and files, is limiting the ability of workflows to support emerging applications in energy informatics that are stream oriented. This is compounded by the absence of Cloud data services that support reliable and performant streams. In this paper, we propose and present a scientific workflow framework that supports streams as first-class data, and is optimized for performant and reliable execution across desktop and Cloud platforms. The workflow framework features and its empirical evaluation on a private Eucalyptus cloud are presented.
Daniel Zinn, Quinn J. Hart, Timothy M. McPhillips, Bertram Ludäscher, Yogesh L. Simmhan, Michail Giakkoupis, Viktor Prasanna 0001
CCGRID5
2011 The Open Provenance Model core specification (v1.1)
Luc Moreau 0001, Ben Clifford, Juliana Freire, Joe Futrelle, Yolanda Gil, Paul Groth, Natalia Kwasnikowska, Simon Miles, Paolo Missier, James D. Myers, Beth Plale, Yogesh L. Simmhan, Eric G. Stephan, Jan Van den Bussche
Future Gener. Comput. Syst.12
2011 Analysis of approaches for supporting the Open Provenance Model: A case study of the Trident workflow workbench
Yogesh L. Simmhan, Roger S. Barga
Future Gener. Comput. Syst.1
2011 Special Section: The third provenance challenge on using the open provenance model for interoperability
Yogesh L. Simmhan, Paul Groth, Luc Moreau 0001
Future Gener. Comput. Syst.1
2010 Bridging the Gap between Desktop and the Cloud for eScience Applications
abstract
The widely discussed scientific data deluge creates a need to computationally scale out eScience applications beyond the local desktop and cope with variable loads over time. Cloud computing offers a scalable, economic, on-demand model well matched to these needs. Yet cloud computing creates gaps that must be crossed to move existing science applications to the cloud. In this article, we propose a Generic Worker framework to deploy and invoke science applications in the cloud with minimal user effort and predictable cost-effective performance. Our framework addresses three distinct challenges posed by the cloud: the complexity of application deployment, invocation of cloud applications from desktop clients, and efficient transparent data transfers across desktop and the cloud. We present an implementation of the Generic Worker for the Microsoft Azure Cloud and evaluate its use for a genomics application. Our evaluation shows that the user complexity to port and scale the application is substantially reduced while introducing a negligible performance overhead of of <; 5% for the genomics application when scaling to 20 VM instances.
Yogesh L. Simmhan, Catharine van Ingen, Girish Subramanian
IEEE CLOUD1
2010 Comparison of resource platform selection approaches for scientific workflows
abstract
Cloud computing is increasingly considered as an additional computational resource platform for scientific workflows. The cloud offers opportunity to scale-out applications from desktops and local cluster resources. Each platform has different properties (e.g., queue wait times in high performance systems, virtual machine startup overhead in clouds) and characteristics (e.g., custom environments in cloud) that makes choosing from these diverse resource platforms for a workflow execution a challenge for scientists. Scientists are often faced with deciding resource platform selection trade-offs with limited information on the actual workflows. While many workflow planning methods have explored resource selection or task scheduling, these methods often require fine-scale characterization of the workflow that is onerous for a scientist. In this paper, we describe our early exploratory work in using blackbox characteristics for a cost-benefit analysis of using different resource platforms. In our blackbox method, we use only limited high-level information on the workflow length, width, and data sizes. The length and width are indicative of the workflow duration and parallelism. We compare the effectiveness of this approach to other resource selection models using two exemplar scientific workflows on desktop, local cluster, HPC center, and cloud platforms. Early results suggest that the blackbox model often makes the same resource selections as a more fine-grained whitebox model. We believe the simplicity of the blackbox model can help inform a scientist on the applicability of a new resource platform, such as cloud resources, even before porting an existing workflow.
Yogesh L. Simmhan, Lavanya Ramakrishnan
HPDC1
2009 Building Reliable Data Pipelines for Managing Community Data Using Scientific Workflows
abstract
The growing amount of scientific data from sensors and field observations is posing a challenge to ¿data valets¿ responsible for managing them in data repositories. These repositories built on commodity clusters need to reliably ingest data continuously and ensure its availability to a wide user community. Workflows provide several benefits to modeling data-intensive science applications and many of these benefits can help manage the data ingest pipelines too. But using workflows is not panacea in itself and data valets need to consider several issues when designing workflows that behave reliably on fault prone hardware while retaining the consistency of the scientific data. In this paper, we propose workflow designs for reliable data ingest in a distributed environment and identify workflow framework features to support resilience. We illustrate these using the data pipeline for the Pan-STARRS repository, one of the largest digital surveys that accumulates 100TB of data annually to support 300 astronomers.
Yogesh L. Simmhan, Catharine van Ingen, Alex Szalay, Roger S. Barga, Jim Heasley
eScience1
2008 The Trident Scientific Workflow Workbench
abstract
In our demonstration we present Trident, a scientific workflow workbench built on top of a commercial workflow system to leverage existing functionality to the extent possible. Trident is being developed in collaboration with the scientific computing community for use in a number of ongoing eScience projects that make use of scientific workflows, in particular the Pan-STARRS sky survey project and the Ocean Observatory Initiative. In our demonstration of Trident we will illustrate the ability to utilize both local and cloud resources for storage and execution, as well as services such as provenance, monitoring, logging and scheduling workflows over clusters. Our goal is to release Trident in early 2009 as an open source accelerator for others to use for eScience projects and to continue extending with support for new workflow features and services.
Roger S. Barga, Jared Jackson, Nelson Araujo, Dean Guo, Nitin Gautam, Yogesh L. Simmhan
eScience6
2008 On Building Scientific Workflow Systems for Data Management in the Cloud
abstract
Scientific workflows have become an archetype to model in silico experiments in the Cloud by scientists. There is a class of workflows that are used to by "data valets" to prepare raw data from scientific instruments into a science-ready form for use by scientists. These share data-intensive traits with traditional scientific workflows, yet differ significantly, for example, in the required degree of reliability and the type of provenance collected. We compare and contrast science application and data valet workflows through exemplar eScience projects to drive shared and unique requirements for scientific workflows across diverse users in a Science Cloud.
Yogesh L. Simmhan, Roger S. Barga, Catharine van Ingen, Edward D. Lazowska, Alex Szalay
eScience1
2008 Special Issue: The First Provenance Challenge
abstract
Abstract The first Provenance Challenge was set up in order to provide a forum for the community to understand the capabilities of different provenance systems and the expressiveness of their provenance representations. To this end, a functional magnetic resonance imaging workflow was defined, which participants had to either simulate or run in order to produce some provenance representation, from which a set of identified queries had to be implemented and executed. Sixteen teams responded to the challenge, and submitted their inputs. In this paper, we present the challenge workflow and queries, and summarize the participants' contributions. Copyright © 2007 John Wiley & Sons, Ltd.
Luc Moreau 0001, Bertram Ludäscher, Ilkay Altintas, Roger S. Barga, Shawn Bowers, Steven P. Callahan, George Chin, Ben Clifford, Shirley Cohen, Sarah Cohen Boulakia, Susan B. Davidson, Ewa Deelman, Luciano A. Digiampietri, Ian T. Foster, Juliana Freire, James Frew, Joe Futrelle, Tara Gibson, Yolanda Gil, Carole A. Goble, Jennifer Golbeck, Paul Groth, David A. Holland, Jihie Kim, David Koop, Ales Krenek, Timothy M. McPhillips, Gaurang Mehta, Simon Miles, Dominic Metzger, Steve Munroe, James D. Myers, Beth Plale, Norbert Podhorszki, Varun Ratnakar, Emanuele Santos, Carlos Scheidegger, Karen Schuchardt, Margo I. Seltzer, Yogesh L. Simmhan, Cláudio T. Silva, Peter Slaughter, Eric G. Stephan, Robert Stevens 0001, Daniele Turi, Huy T. Vo, Michael Wilde, Jun Zhao 0003, Yong Zhao 0009
Concurr. Comput. Pract. Exp.41
2008 Query capabilities of the Karma provenance framework
abstract
Abstract Provenance metadata in e‐Science captures the derivation history of data products generated from scientific workflows. Provenance forms a glue linking workflow execution with associated data products, and finds use in determining the quality of derived data, tracking resource usage, and for verifying and validating scientific experiments. In this article, we discuss the scope of provenance collected in the Karma provenance framework used in the LEAD Cyberinfrastructure project, distinguishing provenance metadata from generic annotations. We further describe our approaches to querying for different forms of provenance in Karma in the context of queries in the first provenance challenge. We use an incremental, building‐block method to construct provenance queries based on the fundamental querying capabilities provided by the Karma service centered on the provenance data model. This has the advantage of keeping the Karma service generic and simple, and yet supports a wide range of queries. Karma successfully answers all but one challenge query. Copyright © 2007 John Wiley & Sons, Ltd.
Yogesh L. Simmhan, Beth Plale, Dennis Gannon
Concurr. Comput. Pract. Exp.1
2006 A Framework for Collecting Provenance in Data-Centric Scientific Workflows
abstract
The increasing ability for the Earth sciences to sense the world around us is resulting in a growing need for data-driven applications that are under the control of data-centric workflows composed of grid- and Web-services. The focus of our work is on provenance collection/or these workflows, necessary to validate the workflow and to determine quality of generated data products. The challenge we address is to record uniform and usable provenance metadata that meets the domain needs while minimizing the modification burden on the service authors and the performance overhead on the workflow engine and the services. The framework, based on a loosely-coupled publish-subscribe architecture for propagating provenance activities, satisfies the needs of detailed provenance collection while a performance evaluation of a prototype finds a minimal performance overhead (in the range of 1% for an eight service workflow using 271 data products)
Yogesh L. Simmhan, Beth Plale, Dennis Gannon
ICWS1
2005 Service Oriented Architectures for Science Gateways on Grid Systems
Dennis Gannon, Beth Plale, Marcus Christie, Scott Jensen, Gopi Kandaswamy, Suresh Marru, Sangmi Lee Pallickara, Satoshi Shirasuna, Yogesh L. Simmhan, Aleksander Slominski, Yiming Sun 0001
ICSOC11
2005 Building Grid Portal Applications From a Web Service Component Architecture
abstract
This work describes an approach to building Grid applications based on the premise that users who wish to access and run these applications prefer to do so without becoming experts on Grid technology. We describe an application architecture based on wrapping user applications and application workflows as Web services and Web service resources. These services are visible to the users and to resource providers through a family of Grid portal components that can be used to configure, launch, and monitor complex applications in the scientific language of the end user. The applications in this model are instantiated by an application factory service. The layered design of the architecture makes it possible for an expert to configure an application factory service with a custom user interface client that may be dynamically loaded into the portal.
Dennis Gannon, Jay Alameda, Octav Chipara, Marcus Christie, Vinayak Dukle, Matthew Farrellee, Gopi Kandaswamy, Deepti Kodeboyina, Sriram Krishnan, Charles W. Moad, Marlon E. Pierce, Beth Plale, Albert L. Rossi, Yogesh L. Simmhan, Anuraag Sarangi, Aleksander Slominski, Satoshi Shirasuna, Thomas Thomas
Proc. IEEE15