Kun-Lung Wu

dblp:85/3566 · DBLP profile ↗
← Back
118ranked-venue papers
31as first author
1since 2021 · last 2021
0000-0002-9173-6165ORCID · reported

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

Databases, data management, data science and information retrieval · 64 · 18 first-authorSystems, architecture and hardware · 28 · 5 first-authorSoftware engineering, systems software and programming languages · 17 · 4 first-author · 1 since 2021Artificial intelligence and machine learning · 13 · 3 first-authorGraphics, computer vision, multimedia, augmented reality and games · 4 · 2 first-authorApplied, interdisciplinary, general and emerging computing · 4 · 3 first-authorComputer networks · 2 · 1 first-authorSecurity and privacy · 2Human-computer interaction and ubiquitous computing · 2 · 2 first-author

Expertise — from the expertise taxonomy: the topics of the expert's papers under the CCF categories. A weight counts papers with recency: 1 for a paper about the topic, 0.3 when the topic is its context, halved every five years.

Databases, data mining, and information retrieval
35 papers
Data stream processing · 66% Indexing and storage engines · 8% Spatial and temporal data management · 8%
Computer architecture, parallel and distributed computing, and storage systems
29 papers
Distributed systems · 35% Cloud and datacenter computing · 22% Parallel and multicore computing · 20%
Theoretical computer science
6 papers
Graph algorithms and graph theory · 65% Algorithmic game theory and mechanism design · 29% Algorithms and data structures · 5%
Computer networks
9 papers
Wireless networking · 46% Content delivery and video streaming · 32% Cellular and mobile networks · 17%
Software engineering, system software, and programming languages
1 paper
Programming languages and type systems · 67% Software testing · 33%

Topics — the 30 heaviest of 127, each with the papers that count most for it

TopicWeightPapersLastEvidence papers
Distributed systems
stream processing
1.142020
Generalizable Resource Allocation in Stream Processing via Deep Reinforcement Learning · AAAI 2020
Automated multi-dimensional elasticity for streaming runtimes: poster · PPoPP 2019
Safe Data Parallelism for General Streaming · IEEE Trans. Computers 2015
Electronic design automation › high-level synthesis
scheduling
0.722020
Generalizable Resource Allocation in Stream Processing via Deep Reinforcement Learning · AAAI 2020
Low-synchronization, mostly lock-free, elastic scheduling for streaming runtimes · PLDI 2017
Data stream processing
load shedding
0.562011
Load Shedding in Mobile Systems with MobiQual · IEEE Trans. Knowl. Data Eng. 2011
MobiQual: QoS-aware Load Shedding in Mobile CQ Systems · ICDE 2008
Efficient Construction of Compact Shedding Filters for Data Stream Processing · ICDE 2008
Parallel and multicore computing
graph partitioning
0.412020
Generalizable Resource Allocation in Stream Processing via Deep Reinforcement Learning · AAAI 2020
Cloud and datacenter computing
resource allocation
0.412020
Generalizable Resource Allocation in Stream Processing via Deep Reinforcement Learning · AAAI 2020
Graph algorithms and graph theory › graph decomposition
k-core decomposition
0.422016
Incremental k-core decomposition: algorithms and evaluation · VLDB J. 2016
Streaming Algorithms for k-core Decomposition · Proc. VLDB Endow. 2013
Data stream processing
stream processing systems
0.422018
Challenges and Experiences in Building an Efficient Apache Beam Runner For IBM Streams · Proc. VLDB Endow. 2018
SPADE: the system s declarative stream processing engine · SIGMOD Conference 2008
Data stream processing
continuous query processing
0.442011
Load Shedding in Mobile Systems with MobiQual · IEEE Trans. Knowl. Data Eng. 2011
MobiQual: QoS-aware Load Shedding in Mobile CQ Systems · ICDE 2008
Efficient Construction of Compact Shedding Filters for Data Stream Processing · ICDE 2008
Cloud and datacenter computing › big data platform
stream processing engine
0.312018
Challenges and Experiences in Building an Efficient Apache Beam Runner For IBM Streams · Proc. VLDB Endow. 2018
Distributed systems
fault tolerance
0.332016
Consistent Regions: Guaranteed Tuple Processing in IBM Streams · Proc. VLDB Endow. 2016
Building User-defined Runtime Adaptation Routines for Stream Processing Applications · Proc. VLDB Endow. 2012
Recoverable Distributed Shared Virtual Memory · IEEE Trans. Computers 1990
Embedded and real-time systems › real-time scheduling › adaptive real-time scheduling
elastic scheduling
0.312017
Low-synchronization, mostly lock-free, elastic scheduling for streaming runtimes · PLDI 2017
Graph algorithms and graph theory
graph decomposition
0.212016
Incremental k-core decomposition: algorithms and evaluation · VLDB J. 2016
Data stream processing › window aggregation
sliding-window aggregation
0.212015
General Incremental Sliding-Window Aggregation · Proc. VLDB Endow. 2015
Parallel and multicore computing
data parallelism
0.212015
Safe Data Parallelism for General Streaming · IEEE Trans. Computers 2015
Parallel and multicore computing
parallel programming models and runtimes
0.212015
Safe Data Parallelism for General Streaming · IEEE Trans. Computers 2015
Data stream processing
elastic scaling
0.212014
Elastic Scaling for Data Stream Processing · IEEE Trans. Parallel Distributed Syst. 2014
Data stream processing
throughput optimization
0.212014
Elastic Scaling for Data Stream Processing · IEEE Trans. Parallel Distributed Syst. 2014
Graph data management › dynamic graph algorithms
incremental graph algorithms
0.212013
Streaming Algorithms for k-core Decomposition · Proc. VLDB Endow. 2013
Data stream processing
streaming graph
0.212013
Counting and Sampling Triangles from a Graph Stream · Proc. VLDB Endow. 2013
Data stream processing › streaming graph
streaming graph algorithms
0.212013
Streaming Algorithms for k-core Decomposition · Proc. VLDB Endow. 2013
Graph data management › motif counting
triangle counting
0.212013
Counting and Sampling Triangles from a Graph Stream · Proc. VLDB Endow. 2013
Programming languages and type systems › programming paradigms
dataflow language
0.212013
Testing properties of dataflow program operators · ASE 2013
Programming languages and type systems › language semantics › formal semantics
operator semantics
0.212013
Testing properties of dataflow program operators · ASE 2013
Software testing › random testing
property-based testing
0.212013
Testing properties of dataflow program operators · ASE 2013
Algorithmic game theory and mechanism design › matching
bipartite matching
0.112021
Fair Task Allocation in Crowdsourced Delivery · IEEE Trans. Serv. Comput. 2021
Data stream processing › stream join
sliding window join
0.122007
GrubJoin: An Adaptive, Multi-Way, Windowed Stream Join with Time Correlation-Aware CPU Load Shedding · IEEE Trans. Knowl. Data Eng. 2007
A Load Shedding Framework and Optimizations for M-way Windowed Stream Joins · ICDE 2007
Data stream processing
stream join
0.122007
GrubJoin: An Adaptive, Multi-Way, Windowed Stream Join with Time Correlation-Aware CPU Load Shedding · IEEE Trans. Knowl. Data Eng. 2007
A Load Shedding Framework and Optimizations for M-way Windowed Stream Joins · ICDE 2007
Wireless networking
mobile ad hoc networks
0.112012
Handling Selfishness in Replica Allocation over a Mobile Ad Hoc Network · IEEE Trans. Mob. Comput. 2012
Cloud and datacenter computing
cluster resource management and scheduling
0.112012
On the optimization of schedules for MapReduce workloads in the presence of shared scans · VLDB J. 2012
Distributed systems › replication
data replication
0.112012
Handling Selfishness in Replica Allocation over a Mobile Ad Hoc Network · IEEE Trans. Mob. Comput. 2012

Methods — techniques the papers use, named apart from their topics

state indexing · 0.7garbage collection · 0.7online/offline algorithms · 0.5incremental algorithm · 0.5f-aware algorithm · 0.52-phase allocation · 0.5graph embedding · 0.4encoder-decoder · 0.4deep reinforcement learning · 0.4differentiated load shedding · 0.4operator fusion · 0.4elastic control algorithm · 0.4dynamic threading · 0.4trusted measurements · 0.3lock-free scheduling · 0.3data replay · 0.2chandy-lamport snapshot algorithm · 0.2flat-array implementation · 0.2
YearPublicationVenuePosition
2021 Fair Task Allocation in Crowdsourced Delivery
abstract
Faster and more cost-efficient, crowdsourced delivery is needed to meet the growing customer demands of many industries, including online shopping, on-demand local delivery, and on-demand transportation. The power of crowdsourced delivery stems from the large number of workers potentially available to provide services and reduce costs. It has been shown in social psychology literature that fairness is key to ensuring high worker participation. However, existing assignment solutions fall short on modeling the dynamic fairness metric. In this work, we introduce a new assignment strategy for crowdsourced delivery tasks. This strategy takes fairness towards workers into consideration, while maximizing the task allocation ratio. Since redundant assignments are not possible in delivery tasks, we first introduce a 2-phase allocation model that increases the reliability of a worker to complete a given task. To realize the effectiveness of our model in practice, we present both offline and online versions of our proposed algorithm called F-Aware. Given a task-to-worker bipartite graph, F-Aware assigns each task to a worker that minimizes unfairness, while allocating tasks to use worker capacities as much as possible. We present an evaluation of our algorithms with respect to running time, task allocation ratio (TAR), as well as unfairness and assignment ratio. Experiments show that F-Aware runs around 107× faster than the TAR-optimal solution and allocates 96.9 percent of the tasks that can be allocated by it. Moreover, it is shown that, F-Aware is able to provide a much fair distribution of tasks to workers than the best competitor algorithm.
Fuat Basik, Bugra Gedik, Hakan Ferhatosmanoglu, Kun-Lung Wu
IEEE Trans. Serv. Comput.4
2020 Generalizable Resource Allocation in Stream Processing via Deep Reinforcement Learning
abstract
This paper considers the problem of resource allocation in stream processing, where continuous data flows must be processed in real time in a large distributed system. To maximize system throughput, the resource allocation strategy that partitions the computation tasks of a stream processing graph onto computing devices must simultaneously balance workload distribution and minimize communication. Since this problem of graph partitioning is known to be NP-complete yet crucial to practical streaming systems, many heuristic-based algorithms have been developed to find reasonably good solutions. In this paper, we present a graph-aware encoder-decoder framework to learn a generalizable resource allocation strategy that can properly distribute computation tasks of stream processing graphs unobserved from training data. We, for the first time, propose to leverage graph embedding to learn the structural information of the stream processing graphs. Jointly trained with the graph-aware decoder using deep reinforcement learning, our approach can effectively find optimized solutions for unseen graphs. Our experiments show that the proposed model outperforms both METIS, a state-of-the-art graph partitioning algorithm, and an LSTM-based encoder-decoder model, in about 70% of the test cases.
Xiang Ni, Jing Li 0025, Mo Yu, Kun-Lung Wu
AAAI5
2019 Automating Multi-level Performance Elastic Components for IBM Streams
abstract
Streaming applications exhibit abundant opportunities for pipeline parallelism, data parallelism and task parallelism. Prior work in IBM Streams introduced an elastic threading model that sought the best performance by automatically tuning the number of threads. In this paper, we introduce the ability to automatically discover where that threading model is profitable. However this introduces a new challenge: we have separate performance elastic mechanisms that are designed with different objectives, leading to potential negative interactions and unintended performance degradation. We present our experiences in overcoming these challenges by showing how to coordinate separate but interfering elasticity mechanisms to maxmize performance gains with stable and fast parallelism exploitation. We first describe an elastic performance mechanism that automatically adapts different threading models to different regions of an application. We then show a coherent ecosystem for coordinating this threading model elasticty with thread count elasticity. This system is an online, stable multi-level elastic coordination scheme that adapts different regions of a streaming application to different threading models and number of threads. We implemented this multi-level coordination scheme in IBM Streams and demonstrated that it (a) scales to over a hundred threads; (b) can improve performance by an order of magnitude on two different processor architectures when an application can benefit from multiple threading models; and (c) achieves performance comparable to hand-optimized applications but with much fewer threads.
Xiang Ni, Scott Schneider 0001, Raju Pavuluri, Jonathan Kaus, Kun-Lung Wu
Middleware5
2019 Automated multi-dimensional elasticity for streaming runtimes: poster
abstract
We present the multi-dimensional elasticity support in IBM Streams 4.3. Automatic operator fusion and dynamic threading were introduced in IBM Streams 4.2, which made it easier to map distributed stream processing to multicore systems through a low-cost operator scheduler and thread count elasticity. To enable these features, the same threading model was applied to the entire application. However, in practice, we have found that some applications have regions best executed under different threading models. In this poster, we introduce threading model elasticity and design a coherent ecosystem for both threading model elasticity and thread count elasticity. We propose an online, stable multidimensional elastic control algorithm that adapts different regions of a streaming application to different threading models and number of threads.
Xiang Ni, Scott Schneider 0001, Raju Pavuluri, Jonathan Kaus, Kun-Lung Wu
PPoPP5
2018 Work-efficient parallel union-find
abstract
Summary The incremental graph connectivity (IGC) problem is to maintain a data structure that can quickly answer whether two given vertices in a graph are connected, while allowing more edges to be added to the graph. IGC is a fundamental problem and can be solved efficiently in the sequential setting using a solution to the classical union‐find problem. However, sequential solutions are not sufficient to handle modern‐day large, rapidly‐changing graphs where edge updates arrive at a very high rate. We present the first shared‐memory parallel data structure for union‐find (equivalently, IGC) that is both provably work‐efficient (ie, performs no more work than the best sequential counterpart) and has polylogarithmic parallel depth. We also present a simpler algorithm with slightly worse theoretical properties, but which is easier to implement and has good practical performance. Our experiments on large graph streams with various degree distributions show that it has good practical performance, capable of processing hundreds of millions of edges per second using a 20‐core machine.
Natcha Simsiri, Kanat Tangwongsan, Srikanta Tirthapura, Kun-Lung Wu
Concurr. Comput. Pract. Exp.4
2018 Challenges and Experiences in Building an Efficient Apache Beam Runner For IBM Streams
abstract
This paper describes the challenges and experiences in the development of IBM Streams runner for Apache Beam. Apache Beam is emerging as a common stream programming interface for multiple computing engines. Each participating engine implements a runner to translate Beam applications into engine-specific programs. Hence, applications written with the Beam SDK can be executed on different underlying stream computing engines, with negligible migration penalty. IBM Streams is a widely-used enterprise streaming platform. It has a rich set of connectors and toolkits for easy integration of streaming applications with other enterprise applications. It also supports a broad range of programming language interfaces, including Java, C++, Python, Stream Processing Language (SPL) and Apache Beam. This paper focuses on our solutions to efficiently support the Beam programming abstractions in IBM Streams runner. Beam organizes data into discrete event time windows. This design, on the one hand, supports out-of-order data arrivals, but on the other hand, forces runners to maintain more states, which leads to higher space and computation overhead. IBM Streams runner mitigates this problem by efficiently indexing inter-dependent states, garbage-collecting stale keys, and enforcing bundle sizes. We also share performance concerns in Beam that could potentially impact applications. Evaluations show that IBM Streams runner outperforms Flink runner and Spark runner in most scenarios when running the Beam NEXMark benchmarks. IBM Streams runner is available for download from IBM Cloud Streaming Analytics service console.
Paul Gerver, John Macmillan, Daniel Debrunner, William Marshall, Kun-Lung Wu
Proc. VLDB Endow.6
2017 Low-synchronization, mostly lock-free, elastic scheduling for streaming runtimes
abstract
We present the scalable, elastic operator scheduler in IBM Streams 4.2. Streams is a distributed stream processing system used in production at many companies in a wide range of industries. The programming language for Streams, SPL, presents operators, tuples and streams as the primary abstractions. A fundamental SPL optimization is operator fusion, where multiple operators execute in the same process. Streams 4.2 introduces automatic submission-time fusion to simplify application development and deployment. However, potentially thousands of operators could then execute in the same process, with no user guidance for thread placement. We needed a way to automatically figure out how many threads to use, with arbitrarily sized applications on a wide variety of hardware, and without any input from programmers. Our solution has two components. The first is a scalable operator scheduler that minimizes synchronization, locks and global data, while allowing threads to execute any operator and dynamically come and go. The second is an elastic algorithm to dynamically adjust the number of threads to optimize performance, using the principles of trusted measurements to establish trends. We demonstrate our scheduler's ability to scale to over a hundred threads, and our elasticity algorithm's ability to adapt to different workloads on an Intel Xeon system with 176 logical cores, and an IBM Power8 system with 184 logical cores.
Scott Schneider 0001, Kun-Lung Wu
PLDI2
2016 Work-Efficient Parallel Union-Find with Applications to Incremental Graph Connectivity
Natcha Simsiri, Kanat Tangwongsan, Srikanta Tirthapura, Kun-Lung Wu
Euro-Par4
2016 Dynamic Load Balancing for Ordered Data-Parallel Regions in Distributed Streaming Systems
Scott Schneider 0001, Joel L. Wolf, Kirsten Hildrum, Rohit Khandekar, Kun-Lung Wu
Middleware5
2016 SONIC: streaming overlapping community detection
Ahmet Erdem Sariyüce, Bugra Gedik, Gabriela Jacques-Silva, Kun-Lung Wu, Ümit V. Çatalyürek
Data Min. Knowl. Discov.4
2016 Consistent Regions: Guaranteed Tuple Processing in IBM Streams
abstract
Guaranteed tuple processing has become critically important for many streaming applications. This paper describes how we enabled IBM Streams, an enterprise-grade stream processing system, to provide data processing guarantees. Our solution goes from language-level abstractions to a runtime protocol. As a result, with a couple of simple annotations at the source code level, IBM Streams developers can define consistent regions , allowing any subgraph of their streaming application to achieve guaranteed tuple processing. At runtime, a consistent region periodically executes a variation of the Chandy-Lamport snapshot algorithm to establish a consistent global state for that region. The coupling of consistent states with data replay enables guaranteed tuple processing.
Gabriela Jacques-Silva, Fang Zheng 0003, Daniel Debrunner, Kun-Lung Wu, Victor Dogaru, Michael Spicer, Ahmet Erdem Sariyüce
Proc. VLDB Endow.4
2016 Incremental k-core decomposition: algorithms and evaluation
Ahmet Erdem Sariyüce, Bugra Gedik, Gabriela Jacques-Silva, Kun-Lung Wu, Ümit V. Çatalyürek
VLDB J.4
2015 Sliding windows over uncertain data streams
Michele Dallachiesa, Gabriela Jacques-Silva, Bugra Gedik, Kun-Lung Wu, Themis Palpanas
Knowl. Inf. Syst.4
2015 General Incremental Sliding-Window Aggregation
abstract
Stream processing is gaining importance as more data becomes available in the form of continuous streams and companies compete to promptly extract insights from them. In such applications, sliding-window aggregation is a central operator, and incremental aggregation helps avoid the performance penalty of re-aggregating from scratch for each window change. This paper presents Reactive Aggregator (RA), a new framework for incremental sliding-window aggregation. RA is general in that it does not require aggregation functions to be invertible or commutative, and it does not require windows to be FIFO. We implemented RA as a drop-in replacement for the Aggregate operator of a commercial streaming engine. Given m updates on a window of size n , RA has an algorithmic complexity of O ( m + m log ( n/m )), rivaling the best prior algorithms for any m . Furthermore, RA's implementation minimizes overheads from allocation and pointer traversals by using a single flat array.
Kanat Tangwongsan, Martin Hirzel, Scott Schneider 0001, Kun-Lung Wu
Proc. VLDB Endow.4
2015 Safe Data Parallelism for General Streaming
abstract
Streaming applications process possibly infinite streams of data and often have both high throughput and low latency requirements. They are comprised of operator graphs that produce and consume data tuples. General streaming applications use stateful, selective, and user-defined operators. The stream programming model naturally exposes task and pipeline parallelism, enabling it to exploit parallel systems of all kinds, including large clusters. However, data parallelism must either be manually introduced by programmers, or extracted as an optimization by compilers. Previous data parallel optimizations did not apply to selective, stateful and user-defined operators. This article presents a compiler and runtime system that automatically extracts data parallelism for general stream processing. Data-parallelization is safe if the transformed program has the same semantics as the original sequential version. The compiler forms parallel regions while considering operator selectivity, state, partitioning, and graph dependencies. The distributed runtime system ensures that tuples always exit parallel regions in the same order they would without data parallelism, using the most efficient strategy as identified by the compiler. Our experiments using 100 cores across 14 machines show linear scalability for parallel regions that are computation-bound, and near linear scalability when tuples are shuffled across parallel regions.
Scott Schneider 0001, Martin Hirzel, Bugra Gedik, Kun-Lung Wu
IEEE Trans. Computers4
2014 Fast Nearest Neighbor Search on Large Time-Evolving Graphs
Leman Akoglu, Rohit Khandekar, Vibhore Kumar, Srinivasan Parthasarathy 0002, Deepak Rajan, Kun-Lung Wu
ECML/PKDD (1)6
2014 Parallel streaming frequency-based aggregates
abstract
We present efficient parallel streaming algorithms for fundamental frequency-based aggregates in both the sliding window and the infinite window settings. In the sliding window setting, we give a parallel algorithm for maintaining a space-bounded block counter (SBBC). Using SBBC, we derive algorithms for basic counting, frequency estimation, and heavy hitters that perform no more work than their best sequential counterparts. In the infinite window setting, we present algorithms for frequency estimation, heavy hitters, and count-min sketch. For both the infinite window and sliding window settings, our parallel algorithms process a "minibatch" of items using linear work and polylog parallel depth. We also prove a lower bound showing that the work of the parallel algorithm is optimal in the case of heavy hitters and frequency estimation. To our knowledge, these are the first parallel algorithms for these problems that are provably work efficient and have low depth.
Kanat Tangwongsan, Srikanta Tirthapura, Kun-Lung Wu
SPAA3
2014 Elastic Scaling for Data Stream Processing
abstract
This article addresses the profitability problem associated with auto-parallelization of general-purpose distributed data stream processing applications. Auto-parallelization involves locating regions in the application's data flow graph that can be replicated at run-time to apply data partitioning, in order to achieve scale. In order to make auto-parallelization effective in practice, the profitability question needs to be answered: How many parallel channels provide the best throughput? The answer to this question changes depending on the workload dynamics and resource availability at run-time. In this article, we propose an elastic auto-parallelization solution that can dynamically adjust the number of channels used to achieve high throughput without unnecessarily wasting resources. Most importantly, our solution can handle partitioned stateful operators via run-time state migration, which is fully transparent to the application developers. We provide an implementation and evaluation of the system on an industrial-strength data stream processing platform to validate our solution.
Bugra Gedik, Scott Schneider 0001, Martin Hirzel, Kun-Lung Wu
IEEE Trans. Parallel Distributed Syst.4
2013 Efficient processing of streaming graphs for evolution-aware clustering
abstract
The clustering of vertices often evolves with time in a streaming graph, where graph update events are given as a stream of edge (vertex) insertions and deletions. Although a sliding window in stream processing naturally captures some cluster evolution, it alone may not be adequate, especially if the window size is large and the clustering within the windowed stream is unstable. Prior graph clustering approaches are mostly insensitive to clustering evolution. In this paper, we present an efficient approach to processing streaming graphs for evolution-aware clustering (EAC) of vertices. We incrementally manage individual connected components as clusters subject to a constraint on the maximal cluster size. For each cluster, we keep the relative recency of edges in a sorted order and favor more recent edges in clustering. We evaluate the effectiveness of EAC and compare it with a previous state-of-the-art evolution-insensitive clustering (EIC) approach. The results show that EAC is both effective and efficient in capturing evolution in a streaming graph. Moreover, we implement EAC as a streaming graph operator on IBM's InfoSphere Streams, a large-scale distributed middleware for stream processing, and show snapshots of the user cluster evolution in a streaming Twitter mention graph.
Mindi Yuan, Kun-Lung Wu, Gabriela Jacques-Silva, Yi Lu 0001
CIKM2
2013 SLIM: A Scalable Location-Sensitive Information Monitoring Service
abstract
Location-sensitive information monitoring services are a centerpiece of the technology for disseminating content-rich information from massive data streams to mobile users. The key challenges for such monitoring services are characterized by the combination of spatial and non-spatial attributes being monitored and the wide spectrum of update rates. A typical example of such services is "alert me when the gas price at a gas station within 5 miles of my current location drops to 4 per gallon". Such a service needs to monitor the gas price changes in conjunction with the highly dynamic nature of location information. Scalability of such location sensitive and content rich information monitoring services in the presence of different update rates and monitoring thresholds poses a big technical challenge. In this paper, we present SLIM, a scalable location sensitive information monitoring service framework with two unique features. First, we make intelligent use of the correlation between spatial and non-spatial attributes involved in the information monitoring service requests to devise a highly scalable distributed spatial trigger evaluation engine. Second, we introduce single and multi-dimensional safe value containment techniques to efficiently perform selective distributed processing of spatial triggers to reduce the amount of unnecessary trigger evaluations. Through extensive experiments, we show that SLIM offers high scalability for location-sensitive, content-rich information monitoring services in terms of the number of information sources being monitored, number of users and monitoring requests.
Bhuvan Bamba, Kun-Lung Wu, Bugra Gedik, Ling Liu 0001
ICWS2
2013 Testing properties of dataflow program operators
abstract
Dataflow programming languages, which represent programs as graphs of data streams and operators, are becoming increasingly popular and being used to create a wide array of commercial software applications. The dependability of programs written in these languages, as well as the systems used to compile and run these programs, hinges on the correctness of the semantic properties associated with operators. Unfortunately, these properties are often poorly defined, and frequently are not checked, and this can lead to a wide range of problems in the programs that use the operators. In this paper we present an approach for improving the dependability of dataflow programs by checking operators for necessary properties. Our approach is dynamic, and involves generating tests whose results are checked to determine whether specific properties hold or not. We present empirical data that shows that our approach is both effective and efficient at assessing the status of properties.
Martin Hirzel, Gregg Rothermel, Kun-Lung Wu
ASE4
2013 Counting and Sampling Triangles from a Graph Stream
abstract
This paper presents a new space-efficient algorithm for counting and sampling triangles--and more generally, constant-sized cliques--in a massive graph whose edges arrive as a stream. Compared to prior work, our algorithm yields significant improvements in the space and time complexity for these fundamental problems. Our algorithm is simple to implement and has very good practical performance on large graphs.
Aduri Pavan, Kanat Tangwongsan, Srikanta Tirthapura, Kun-Lung Wu
Proc. VLDB Endow.4
2013 Streaming Algorithms for k-core Decomposition
abstract
A k -core of a graph is a maximal connected subgraph in which every vertex is connected to at least k vertices in the subgraph. k -core decomposition is often used in large-scale network analysis, such as community detection, protein function prediction, visualization, and solving NP-Hard problems on real networks efficiently, like maximal clique finding. In many real-world applications, networks change over time. As a result, it is essential to develop efficient incremental algorithms for streaming graph data. In this paper, we propose the first incremental k -core decomposition algorithms for streaming graph data. These algorithms locate a small subgraph that is guaranteed to contain the list of vertices whose maximum k -core values have to be updated, and efficiently process this subgraph to update the k -core decomposition. Our results show a significant reduction in run-time compared to non-incremental alternatives. We show the efficiency of our algorithms on different types of real and synthetic graphs, at different scales. For a graph of 16 million vertices, we observe speedups reaching a million times, relative to the non-incremental algorithms.
Ahmet Erdem Sariyüce, Bugra Gedik, Gabriela Jacques-Silva, Kun-Lung Wu, Ümit V. Çatalyürek
Proc. VLDB Endow.4
2012 Auto-parallelizing stateful distributed streaming applications
abstract
Streaming applications transform possibly infinite streams of data and often have both high throughput and low latency requirements. They are comprised of operator graphs that produce and consume data tuples. The streaming programming model naturally exposes task and pipeline parallelism, enabling it to exploit parallel systems of all kinds, including large clusters. However, it does not naturally expose data parallelism, which must instead be extracted from streaming applications. This paper presents a compiler and runtime system that automatically extract data parallelism for distributed stream processing. Our approach guarantees safety, even in the presence of stateful, selective, and user-defined operators. When constructing parallel regions, the compiler ensures safety by considering an operator's selectivity, state, partitioning, and dependencies on other operators in the graph. The distributed runtime system ensures that tuples always exit parallel regions in the same order they would without data parallelism, using the most efficient strategy as identified by the compiler. Our experiments using 100 cores across 14 machines show linear scalability for standard parallel regions, and near linear scalability when tuples are shuffled across parallel regions.
Scott Schneider 0001, Martin Hirzel, Bugra Gedik, Kun-Lung Wu
PACT4
2012 Clustering Streaming Graphs
abstract
In this paper, we propose techniques for clustering large-scale "streaming" graphs where the updates to a graph are given in form of a stream of vertex or edge additions and deletions. Our algorithm handles such updates in an online and incremental manner and it can be easily parallel zed. Several previous graph clustering algorithms fall short of handling massive and streaming graphs because they are centralized, they need to know the entire graph beforehand and are not incremental, or they incur an excessive computational overhead. Our algorithm's fundamental building block is called graph reservoir sampling. We maintain a reservoir sample of the edges as the graph changes while satisfying certain desired properties like bounding number of clusters or cluster-sizes. We then declare connected components in the sampled sub graph as clusters of the original graph. Our experiments on real graphs show that our approach not only yields clusterings with very good quality, but also obtains orders of magnitude higher throughput, when compared to offline algorithms.
Ahmed Eldawy, Rohit Khandekar, Kun-Lung Wu
ICDCS3
2012 Examining the impact of data-access cost on XML twig pattern matching
SangKeun Lee 0001, Byung-Gul Ryu, Kun-Lung Wu
Inf. Sci.3
2012 Building User-defined Runtime Adaptation Routines for Stream Processing Applications
abstract
Stream processing applications are deployed as continuous queries that run from the time of their submission until their cancellation. This deployment mode limits developers who need their applications to perform runtime adaptation, such as algorithmic adjustments, incremental job deployment, and application-specific failure recovery. Currently, developers do runtime adaptation by using external scripts and/or by inserting operators into the stream processing graph that are unrelated to the data processing logic. In this paper, we describe a component called orchestrator that allows users to write routines for automatically adapting the application to runtime conditions. Developers build an orchestrator by registering and handling events as well as specifying actuations. Events can be generated due to changes in the system state (e.g., application component failures), built-in system metrics (e.g., throughput of a connection), or custom application metrics (e.g., quality score). Once the orchestrator receives an event, users can take adaptation actions by using the orchestrator actuation APIs. We demonstrate the use of the orchestrator in IBM's System S in the context of three different applications, illustrating application adaptation to changes on the incoming data distribution, to application failures, and on-demand dynamic composition.
Gabriela Jacques-Silva, Bugra Gedik, Rohit Wagle, Kun-Lung Wu, Vibhore Kumar
Proc. VLDB Endow.4
2012 Handling Selfishness in Replica Allocation over a Mobile Ad Hoc Network
abstract
In a mobile ad hoc network, the mobility and resource constraints of mobile nodes may lead to network partitioning or performance degradation. Several data replication techniques have been proposed to minimize performance degradation. Most of them assume that all mobile nodes collaborate fully in terms of sharing their memory space. In reality, however, some nodes may selfishly decide only to cooperate partially, or not at all, with other nodes. These selfish nodes could then reduce the overall data accessibility in the network. In this paper, we examine the impact of selfish nodes in a mobile ad hoc network from the perspective of replica allocation. We term this selfish replica allocation. In particular, we develop a selfish node detection algorithm that considers partial selfishness and novel replica allocation techniques to properly cope with selfish replica allocation. The conducted simulations demonstrate the proposed approach outperforms traditional cooperative replica allocation techniques in terms of data accessibility, communication cost, and average query delay.
Jae-Ho Choi 0001, Kyu-Sun Shim, SangKeun Lee 0001, Kun-Lung Wu
IEEE Trans. Mob. Comput.4
2012 On the optimization of schedules for MapReduce workloads in the presence of shared scans
Joel L. Wolf, Andrey Balmin, Deepak Rajan, Kirsten Hildrum, Rohit Khandekar, Sujay S. Parekh, Kun-Lung Wu, Rares Vernica
VLDB J.7
2011 Modeling stream processing applications for dependability evaluation
abstract
This paper describes a modeling framework for evaluating the impact of faults on the output of streaming applications. Our model is based on three abstractions: stream operators, stream connections, and tuples. By composing these abstractions within a Stochastic Activity Network, we allow the modeling of complete applications. We consider faults that lead to data loss and to silent data corruption (SDC). Our framework captures how faults originating in one operator propagate to other operators down the stream processing graph. We demonstrate the extensibility of our framework by evaluating three different fault tolerance techniques: checkpointing, partial graph replication, and full graph replication. We show that under crashes that lead to data loss, partial graph replication has a great advantage in maintaining the accuracy of the application output when compared to checkpointing. We also show that SDC can break the no data duplication guarantees of a full graph replication-based fault tolerance technique.
Gabriela Jacques-Silva, Zbigniew T. Kalbarczyk, Bugra Gedik, Henrique Andrade, Kun-Lung Wu, Ravishankar K. Iyer
DSN5
2011 Processing high data rate streams in System S
Henrique Andrade, Bugra Gedik, Kun-Lung Wu, Philip S. Yu
J. Parallel Distributed Comput.3
2011 Load Shedding in Mobile Systems with MobiQual
abstract
In location-based, mobile continual query (CQ) systems, two key measures of quality-of-service (QoS) are: freshness and accuracy. To achieve freshness, the CQ server must perform frequent query reevaluations. To attain accuracy, the CQ server must receive and process frequent position updates from the mobile nodes. However, it is often difficult to obtain fresh and accurate CQ results simultaneously, due to 1) limited resources in computing and communication and 2) fast-changing load conditions caused by continuous mobile node movement. Hence, a key challenge for a mobile CQ system is: How do we achieve the highest possible quality of the CQ results, in both freshness and accuracy, with currently available resources? In this paper, we formulate this problem as a load shedding one, and develop MobiQual—a QoS-aware approach to performing both update load shedding and query load shedding. The design of MobiQual highlights three important features. 1) Differentiated load shedding: We apply different amounts of query load shedding and update load shedding to different groups of queries and mobile nodes, respectively. 2) Per-query QoS specification: Individualized QoS specifications are used to maximize the overall freshness and accuracy of the query results. 3) Low-cost adaptation: MobiQual dynamically adapts, with a minimal overhead, to changing load conditions and available resources. We conduct a set of comprehensive experiments to evaluate the effectiveness of MobiQual. The results show that, through a careful combination of update and query load shedding, the MobiQual approach leads to much higher freshness and accuracy in the query results in all cases, compared to existing approaches that lack the QoS-awareness properties of MobiQual, as well as the solutions that perform query-only or update-only load shedding.
Bugra Gedik, Kun-Lung Wu, Ling Liu 0001, Philip S. Yu
IEEE Trans. Knowl. Data Eng.2
2010 DEDUCE: at the intersection of MapReduce and stream processing
abstract
MapReduce and stream processing are two emerging, but different, paradigms for analyzing, processing and making sense of large volumes of modern day data. While MapReduce offers the capability to analyze several terabytes of stored data, stream processing solutions offer the ability to process, possibly, a few million updates every second. However, there is an increasing number of data processing applications which need a solution that effectively and efficiently combines the benefits of MapReduce and stream processing to address their data processing needs. For example, in the automated stock trading domain, applications usually require periodic analysis of large amounts of stored data to generate a model using MapReduce, which is then used to process a stream of incident updates using a stream processing system. This paper presents Deduce, which extends IBM's System S stream processing middleware with support for MapReduce by providing (1) language and runtime support for easily specifying and embedding MapReduce jobs as elements of a larger data-flow, (2) capability to describe reusable modules that can be used as map and reduce tasks, and (3) configuration parameters that can be tweaked to control and manage the usage of shared resources by the MapReduce and stream processing components. We describe the motivation for Deduce and the design and implementation of the MapReduce extensions for System S, and then present experimental results.
Vibhore Kumar, Henrique Andrade, Bugra Gedik, Kun-Lung Wu
EDBT4
2010 A Universal Calculus for Stream Processing Languages
Robert Soulé, Martin Hirzel, Robert Grimm 0001, Bugra Gedik, Henrique Andrade, Vibhore Kumar, Kun-Lung Wu
ESOP7
2010 FLEX: A Slot Allocation Scheduling Optimizer for MapReduce Workloads
Joel L. Wolf, Deepak Rajan, Kirsten Hildrum, Rohit Khandekar, Vibhore Kumar, Sujay S. Parekh, Kun-Lung Wu, Andrey Balmin
Middleware7
2010 Towards proximity pattern mining in large graphs
abstract
Mining graph patterns in large networks is critical to a variety of applications such as malware detection and biological module discovery. However, frequent subgraphs are often ineffective to capture association existing in these applications, due to the complexity of isomorphism testing and the inelastic pattern definition.
Arijit Khan 0001, Xifeng Yan, Kun-Lung Wu
SIGMOD Conference3
2010 Efficient B-tree Based Indexing for Cloud Data Processing
abstract
A Cloud may be seen as a type of flexible computing infrastructure consisting of many compute nodes, where resizable computing capacities can be provided to different customers. To fully harness the power of the Cloud, efficient data management is needed to handle huge volumes of data and support a large number of concurrent end users. To achieve that, a scalable and high-throughput indexing scheme is generally required. Such an indexing scheme must not only incur a low maintenance cost but also support parallel search to improve scalability. In this paper, we present a novel, scalable B + -tree based indexing scheme for efficient data processing in the Cloud. Our approach can be summarized as follows. First, we build a local B + -tree index for each compute node which only indexes data residing on the node. Second, we organize the compute nodes as a structured overlay and publish a portion of the local B + -tree nodes to the overlay for efficient query processing. Finally, we propose an adaptive algorithm to select the published B + -tree nodes according to query patterns. We conduct extensive experiments on Amazon's EC2, and the results demonstrate that our indexing scheme is dynamic, efficient and scalable.
Sai Wu, Dawei Jiang, Beng Chin Ooi, Kun-Lung Wu
Proc. VLDB Endow.4
2010 From a Stream of Relational Queries to Distributed Stream Processing
abstract
Applications from several domains are now being written to process live data originating from hardware and software-based streaming sources. Many of these applications have been written relying solely on database and data warehouse technologies, despite their lack of need for transactional support and ACID properties. In several extreme high-load cases, this approach does not scale to the processing speeds that these applications demand. In this paper we demonstrate an application acceleration approach whereby a regular ODBC-based application is converted into a true streaming application with minimal disruption from a software engineering standpoint. We showcase our approach on three real-world applications. We experimentally demonstrate the substantial performance improvements that can be observed when contrasting the accelerated implementation with the original database-oriented implementation.
Qiong Zou, Huayong Wang, Robert Soulé, Martin Hirzel, Henrique Andrade, Bugra Gedik, Kun-Lung Wu
Proc. VLDB Endow.7
2009 A code generation approach to optimizing high-performance distributed data stream processing
abstract
We present a code-generation-based optimization approach to bringing performance and scalability to distributed stream processing applications. We express stream processing applications using an operator-based, stream-centric language called SPADE, which supports composing distributed data flow graphs out of toolkits of type-generic operators. A major challenge in building such applications is to find an effective and flexible way of mapping the logical graph of operators into a physical one that can be deployed on a set of distributed nodes. This involves finding how best operators map to processes and how best processes map to computing nodes. In this paper, we take a two-stage optimization approach, where an instrumented version of the application is first generated by the SPADE compiler to profile and collect statistics about the processing and communication characteristics of the operators within the application. In the second stage, the profiling information is fed to an optimizer to come up with a physical data flow graph that is deployable across nodes in a computing cluster. This approach not only creates highly optimized applications that are tailored to the underlying computing and networking infrastructure, but also makes it possible to re-target the application to a different hardware setup by simply repeating the optimization step and re-compiling the application to match the physical flow graph produced by the optimizer. Using real-world applications, from diverse domains such as finance and radio-astronomy, we demonstrate the effectiveness of our approach on System S -- a large-scale, distributed stream processing platform.
Bugra Gedik, Henrique Andrade, Kun-Lung Wu
CIKM3
2009 Characterizing, constructing and managing resource usage profiles of system S applications: challenges and experience
abstract
We describe the challenges of characterizing, constructing and managing the usage profiles of System S applications. A running System S application is a directed graph with software processing elements(PEs) as vertices and data streams as edges connecting the PEs. The resource usage of each PE is a critical input to the runtime scheduler for proper resource allocation. We represent the resource usage of PEs in terms of resource functions (RFs) that are used by the System S scheduler, with one RF per resource per PE. The first challenge is that it is difficult to build good RFs that can accurately predict the resource usage of a PE because the PEs perform arbitrary computations. A second set of challenges arises in managing the RFs and performance data so that we can apply them for PEs that are re-run or reused by the same or different applications or users. We report our experience in overcoming these challenges. Specifically, we present an empirical characterization of PE RFs from several real streaming applications running in a System S testbed. This indicates that our simple models of resource usage that build on the data-flow nature of the underlying application can be effective, even for complex PEs. To illustrate our methodology, we evaluate and analyze the performance of these applications as a function of the quality of our resource profile models. The system automatically learns the models from the raw metrics data collected from running PEs. We describe our approach to managing the metrics and RF models, which allows us to construct generalizable RFs and eliminates the learning time for new PEs by intelligently storing and reusing the metrics data.
Sujay S. Parekh, Kirsten Hildrum, Deepak Rajan, Joel L. Wolf, Kun-Lung Wu
CIKM5
2009 Language level checkpointing support for stream processing applications
abstract
Many streaming applications demand continuous processing of live data with little or no downtime, therefore, making high-availability a crucial operational requirement. Fault tolerance techniques are generally expensive and when directly applied to streaming systems with stringent throughput and latency requirements, they might incur a prohibitive performance overhead. This paper describes a flexible, light-weight fault tolerance solution in the context of the SPADE language and the System S distributed stream processing engine. We devised language extensions so users can define and parameterize check-point policies easily. This configurable fault tolerance solution is implemented through code generation in SPADE, which reduces the overall application fault tolerance costs by incurring them only for the parts of the application that require it. In this paper we focus on the overall design of our checkpoint mechanism and we also describe an incremental checkpointing algorithm that is suitable for on-the-fly processing of high-rate data streams.
Gabriela Jacques-Silva, Bugra Gedik, Henrique Andrade, Kun-Lung Wu
DSN4
2009 PROUD: a probabilistic approach to processing similarity queries over uncertain data streams
abstract
We present PROUD -- A PRObabilistic approach to processing similarity queries over Uncertain Data streams, where the data streams here are mainly time series streams. In contrast to data with certainty, an uncertain series is an ordered sequence of random variables. The distance between two uncertain series is also a random variable. We use a general uncertain data model, where only the mean and the deviation of each random variable at each timestamp are available. We derive mathematical conditions for progressively pruning candidates to reduce the computation cost. We then apply PROUD to a streaming environment where only sketches of streams, like wavelet synopses, are available. Extensive experiments are conducted to evaluate the effectiveness of PROUD and compare it with Det, a deterministic approach that directly processes data without considering uncertainty. The results show that, compared with Det, PROUD offers a flexible trade-off between false positives and false negatives by controlling a threshold, while maintaining a similar computation cost. In contrast, Det does not provide such flexibility. This trade-off is important as in some applications false negatives are more costly, while in others, it is more critical to keep the false positives low.
Mi-Yen Yeh, Kun-Lung Wu, Philip S. Yu, Ming-Syan Chen
EDBT2
2009 Scale-Up Strategies for Processing High-Rate Data Streams in System S
abstract
High performance stream processing is critical in sense-and-respond application domains – from environmental monitoring to algorithmic trading. In this paper, we focus on language and runtime support for improving the performance of sense-and-respond applications in processing data from high rate streams. The central tenet of this work is the definition of a streaming architectural pattern for these application domains and the programming model and the code generation framework to support it. Using IBM Research's System S middleware and the SPADE language, we demonstrate how to scale up a financial trading application.
Henrique Andrade, Bugra Gedik, Kun-Lung Wu, Philip S. Yu
ICDE3
2009 Auto-vectorization through code generation for stream processing applications
abstract
We describe language- and code generation-based approaches to providing access to architecture-specific vectorization support for high-performance data stream processing applications. We provide an experimental performance evaluation of several stream operators, contrasting our code generation approach with the native auto-vectorization support available in the GNU gcc and Intel icc compilers.
Huayong Wang, Henrique Andrade, Bugra Gedik, Kun-Lung Wu
ICS4
2009 Elastic scaling of data parallel operators in stream processing
abstract
We describe an approach to elastically scale the performance of a data analytics operator that is part of a streaming application. Our techniques focus on dynamically adjusting the amount of computation an operator can carry out in response to changes in incoming workload and the availability of processing cycles. We show that our elastic approach is beneficial in light of the dynamic aspects of streaming workloads and stream processing environments. Addressing another recent trend, we show the importance of our approach as a means to providing computational elasticity in multicore processor-based environments such that operators can automatically find their best operating point. Finally, we present experiments driven by synthetic workloads, showing the space where the optimizing efforts are most beneficial and a radioastronomy imaging application, where we observe substantial improvements in its performance-critical section.
Scott Schneider 0001, Henrique Andrade, Bugra Gedik, Alain Biem, Kun-Lung Wu
IPDPS5
2009 Job Admission and Resource Allocation in Distributed Streaming Systems
Joel L. Wolf, Nikhil Bansal 0001, Kirsten Hildrum, Sujay S. Parekh, Deepak Rajan, Rohit Wagle, Kun-Lung Wu
JSSPP7
2009 COLA: Optimizing Stream Processing Applications via Graph Partitioning
Rohit Khandekar, Kirsten Hildrum, Sujay S. Parekh, Deepak Rajan, Joel L. Wolf, Kun-Lung Wu, Henrique Andrade, Bugra Gedik
Middleware6
2009 Answering linear optimization queries with an approximate stream index
Gang Luo 0001, Kun-Lung Wu, Philip S. Yu
Knowl. Inf. Syst.2
2009 Tools and strategies for debugging distributed stream processing applications
abstract
Abstract Distributed data stream processing applications are often characterized by data flow graphs consisting of a large number of built‐in and user‐defined operators connected via streams. These flow graphs are typically deployed on a large set of nodes. The data processing is carried out on‐the‐fly, as tuples arrive at possibly very high rates, with minimum latency. It is well known that developing and debugging distributed, multi‐threaded, and asynchronous applications, such as stream processing applications, can be challenging. Thus, without domain‐specific debugging support, developers struggle when debugging distributed applications. In this paper, we describe tools and language support to support debugging distributed stream processing applications. Our key insight is to view debugging of stream processing applications from four different, but related, perspectives. First, debugging the semantics of the application involves verifying the operator‐level composition and inspecting the flows at the logical level. Second, debugging the user‐defined operators involves traditional source‐code debugging, but strongly tied to the stream‐level interactions. Third, debugging the deployment details of the application require understanding the runtime physical layout and configuration of the application. Fourth, debugging the performance of the application requires inspecting various performance metrics (such as communication rates, CPU utilization, etc.) associated with streams, operators, and nodes in the system. In light of this characterization, we developed several tools such as adebugger‐awarecompiler and an associatedstream debugger,composition and deployment visualizers, andperformance visualizers, as well as language support, such as configuration knobs for logging and tracing, deployment configurations such as operator‐to‐process and process‐to‐node mappings, monitoring directives to inspect streams, and special sink adapters to intercept and dump streaming data to files and sockets, to name a few. We describe these tools in the context ofSpade—a language for creating distributed stream processing applications, andSystem S—a distributed stream processing middleware under development at the IBM Watson Research Center. Published in 2009 by John Wiley & Sons, Ltd.
Bugra Gedik, Henrique Andrade, Andy Frenkiel, Wim De Pauw, Michael Pfeifer, Paul Allen, Norman Cohen, Kun-Lung Wu
Softw. Pract. Exp.8
2008 Efficient Construction of Compact Shedding Filters for Data Stream Processing
abstract
High-volume source streams, coupled with fluctuating rates, necessitate adaptive load shedding in data stream processing. When ignored, a continual query (CQ) server may randomly drop items, when its capacity is inadequate to handle the arriving data, and degrade the quality of the query results. To alleviate this problem, filters can be used at the source nodes. However, regular source filtering in itself is not sufficient to prevent random dropping, because the amount of data passing through the filters can still surpass the server's capacity. In this case, intelligent load shedding can be applied by the source filters to minimize the degradation in result quality. In this paper, we introduce a novel type of load-shedding source filters, called non- uniformly regulated (NR) sifters. An NR sifter judiciously applies varying amounts of load shedding to different regions of the data space within the sifter. We formulate the problem of constructing NR sifters as an optimization one. NR sifters are compact and quickly configurable, allowing frequent adaptations, and provide fast lookup f.or deciding if a data item should be dropped. We structure NR sifters as a set of (sifter region, drop threshold) pairs to achieve compactness, develop query consolidation techniques to enable quick construction, and introduce flexible space partitioning mechanisms to realize fast lookup.
Bugra Gedik, Kun-Lung Wu, Philip S. Yu
ICDE2
2008 MobiQual: QoS-aware Load Shedding in Mobile CQ Systems
abstract
Freshness and accuracy are two key measures of quality of service (QoS) in location-based, mobile continual queries (CQs). However, it is often difficult to provide both fresh and accurate CQ results due to (a) limited resources in computing and communication and (b) fast-changing load conditions caused by continuous mobile node movement. Thus a key challenge for a mobile CQ system is: How do we achieve the highest possible quality of the query results, in both freshness and accuracy, with currently available resources under changing load conditions? In this paper, we formulate this problem as a load shedding one, and develop MobiQual - a QoS-aware framework for performing both update load shedding and query load shedding. The design of MobiQual highlights three important features. (1)Differentiatedloadshedding: Different amounts of query and update load shedding are applied to different groups of queries and mobile nodes, respectively. (2)Per-queryQoSspecifications: The overall freshness and accuracy of the query results are maximized with individualized QoS specifications. (3)Low-costadaptation: MobiQual dynamically adapts, with a minimal overhead, to changing load conditions and available resources. We show that, through a careful combination of update and query load shedding, the MobiQual approach leads to much higher freshness and accuracy in the query results in all cases, compared to existing approaches.
Bugra Gedik, Kun-Lung Wu, Philip S. Yu, Ling Liu 0001
ICDE2
2008 SODA: An Optimizing Scheduler for Large-Scale Stream-Based Distributed Computer Systems
Joel L. Wolf, Nikhil Bansal 0001, Kirsten Hildrum, Sujay S. Parekh, Deepak Rajan, Rohit Wagle, Kun-Lung Wu, Lisa Fleischer
Middleware7
2008 SPADE: the system s declarative stream processing engine
abstract
In this paper, we present Spade - the System S declarative stream processing engine. System S is a large-scale, distributed data stream processing middleware under development at IBM T. J. Watson Research Center. As a front-end for rapid application development for System S, Spade provides (1) an intermediate language for flexible composition of parallel and distributed data-flow graphs, (2) a toolkit of type-generic, built-in stream processing operators, that support scalar as well as vectorized processing and can seamlessly inter-operate with user-defined operators, and (3) a rich set of stream adapters to ingest/publish data from/to outside sources. More importantly, Spade automatically brings performance optimization and scalability to System S applications. To that end, Spade employs a code generation framework to create highly-optimized applications that run natively on the Stream Processing Core (SPC), the execution and communication substrate of System S, and take full advantage of other System S services. Spade allows developers to construct their applications with fine granular stream operators without worrying about the performance implications that might exist, even in a distributed system. Spade's optimizing compiler automatically maps applications into appropriately sized execution units in order to minimize communication overhead, while at the same time exploiting available parallelism. By virtue of the scalability of the System S runtime and Spade's effective code generation and optimization, we can scale applications to a large number of nodes. Currently, we can run Spade jobs on ≈ 500 processors within more than 100 physical nodes in a tightly connected cluster environment. Spade has been in use at IBM Research to create real-world streaming applications, ranging from monitoring financial market feeds to radio telescopes to semiconductor fabrication lines.
Bugra Gedik, Henrique Andrade, Kun-Lung Wu, Philip S. Yu, Myungcheol Doo
SIGMOD Conference3
2008 Correlating burst events on streaming stock market data
Michail Vlachos, Kun-Lung Wu, Shyh-Kwei Chen, Philip S. Yu
Data Min. Knowl. Discov.2
2008 LeeWave: level-wise distribution of wavelet coefficients for processing kNN queries over distributed streams
abstract
We present LeeWave --- a bandwidth-efficient approach to searching range-specified k -nearest neighbors among distributed streams by LEvEl-wise distribution of WAVElet coefficients. To find the k most similar streams to a range-specified reference one, the relevant wavelet coefficients of the reference stream can be sent to the peer sites to compute the similarities. However, bandwidth can be unnecessarily wasted if the entire relevant coefficients are sent simultaneously. Instead, we present a level-wise approach by leveraging the multi-resolution property of the wavelet coefficients. Starting from the top and moving down one level at a time, the query initiator sends only the single-level coefficients to a progressively shrinking set of candidates. However, there is one difficult challenge in LeeWave: how does the query initiator prune the candidates without knowing all the relevant coefficients? To overcome this challenge, we derive and maintain a similarity range for each candidate and gradually tighten the bounds of this range as we move from one level to the next. The increasingly tightened similarity ranges enable the query initiator to effectively prune the candidates without causing any false dismissal. Extensive experiments with real and synthetic data show that, when compared with prior approaches, LeeWave uses significantly less bandwidth under a wide range of conditions.
Mi-Yen Yeh, Kun-Lung Wu, Philip S. Yu, Ming-Syan Chen
Proc. VLDB Endow.2
2007 Lira: Lightweight, Region-aware Load Shedding in Mobile CQ Systems
abstract
To provide high-quality results for location-based, continual queries (CQs) in a mobile system, the query processor usually demands receiving frequent position updates from the mobile nodes. However, processing frequent updates often causes the query processor to become overloaded, under which updates must be dropped randomly, bringing down the query-result accuracy and negating the benefits of frequent updates. In this paper, we develop LIRA - a lightweight, region-aware load-shedding technique for preventively reducing the position-update load of a query processor, while maintaining high-quality query results. Instead of receiving too many updates and then randomly dropping some of them, LIRA uses a region-aware partitioning mechanism to identify the most beneficial shedding regions to cut down the position updates sent by the mobile nodes within those regions. Based on the densities of mobile nodes and queries in a region, LIRA judiciously applies different amounts of update reduction for different regions, aiming to minimize the negative impacts of load shedding on query-result accuracy. Experimental results show that LIRA is vastly superior to random update dropping and clearly outperforms other alternatives that do not possess region-aware load-shedding capabilities. Moreover, due to its lightweight nature, LIRA introduces very little overhead.
Bugra Gedik, Ling Liu 0001, Kun-Lung Wu, Philip S. Yu
ICDE3
2007 A Load Shedding Framework and Optimizations for M-way Windowed Stream Joins
abstract
Tuple dropping, though commonly used for load shedding in most stream operations, is inadequate for m-way, windowed stream joins. The join output rate can be overly reduced because it fails to exploit the time correlations likely to exist among interrelated streams. In this paper, we introduce GrubJoin; an adaptive, m-way, windowed stream join that effectively performs time correlation-aware CPU load shedding. GrubJoin maximizes the output rate by achieving near-optimal window harvesting, which picks only the most profitable window segments for the join. Due to combinatorial explosion of possible m-way join sequences involving window segments, m-way, windowed stream joins pose several unique challenges. We focus on addressing two of them: (1) How can we quickly determine the optimal window harvesting configuration for any m-way, windowed stream join? (2) How can we monitor and learn the time correlations among the streams with high accuracy and minimal overhead? To tackle these challenges, we formalize window harvesting as an optimization problem, develop greedy heuristics to determine near-optimal window harvesting configurations and use approximation techniques to capture the time correlations. Our experimental results show that GrubJoin is vastly superior to tuple dropping when time correlations exist and is equally effective when time correlations are nonexistent.
Bugra Gedik, Kun-Lung Wu, Philip S. Yu, Ling Liu 0001
ICDE2
2007 SAO: A Stream Index for Answering Linear Optimization Queries
abstract
Linear optimization queries retrieve the top-K tuples in a sliding window of a data stream that maximize/minimize the linearly weighted sums of certain attribute values. To efficiently answer such queries against a large relation, an onion index was previously proposed to properly organize all the tuples in the relation. However, such an onion index does not work in a streaming environment due to fast tuple arrival rate and limited memory. In this paper, we propose a SAO index to approximately answer arbitrary linear optimization queries against a data stream. It uses a small amount of memory to efficiently keep track of the most "important" tuples in a sliding window of a data stream. The index maintenance cost is small because the great majority of the incoming tuples do not cause any changes to the index and are quickly discarded. At any time, for any linear optimization query, we can retrieve from the SAO index the approximate top-K tuples in the sliding window almost instantly. The larger the amount of available memory, the better the quality of the answers is. More importantly, for a given amount of memory, the quality of the answers can be further improved by dynamically allocating a larger portion of the memory to the outer layers of the SAO index. We evaluate the effectiveness of this SAO index through a prototype implementation.
Gang Luo 0001, Kun-Lung Wu, Philip S. Yu
ICDE2
2007 Challenges and Experience in Prototyping a Multi-Modal Stream Analytic and Monitoring Application on System S
Kun-Lung Wu, Philip S. Yu, Bugra Gedik, Kirsten Hildrum, Charu C. Aggarwal, Eric Bouillet, Wei Fan 0001, Xiaohui Gu, Gang Luo 0001, Haixun Wang
VLDB1
2007 CPU load shedding for binary stream joins
Bugra Gedik, Kun-Lung Wu, Philip S. Yu, Ling Liu 0001
Knowl. Inf. Syst.2
2007 GrubJoin: An Adaptive, Multi-Way, Windowed Stream Join with Time Correlation-Aware CPU Load Shedding
abstract
Tuple dropping, though commonly used for load shedding in most data stream operations, is generally inadequate for multiway windowed stream joins. The join output rate can be unnecessarily reduced because tuple dropping fails to exploit the time correlations that are likely to exist among interrelated streams. In this paper, we introduce GrubJoin-an adaptive multiway windowed stream join that effectively performs time correlation-aware CPU load shedding. GrubJoin maximizes the output rate by achieving near-optimal window harvesting, which picks only the most profitable segments of individual windows for the join. Due mainly to the combinatorial explosion of possible multiway join sequences involving different window segments, GrubJoin faces unique challenges that do not exist for binary joins, such as determining the optimal window harvesting configuration in a time-efficient manner and learning the time correlations among the streams without introducing overhead. To tackle these challenges, we formalize window harvesting as an optimization problem, develop greedy heuristics to determine near-optimal window harvesting configurations, and use approximation techniques to capture the time correlations. Our experimental results show that GrubJoin is vastly superior to tuple dropping when time correlations exist and is equally effective when time correlations are nonexistent.
Bugra Gedik, Kun-Lung Wu, Philip S. Yu, Ling Liu 0001
IEEE Trans. Knowl. Data Eng.2
2006 On-Demand Index for Efficient Structural Joins
Kun-Lung Wu, Shyh-Kwei Chen, Philip S. Yu
WAIM1
2006 A practical approach to extracting DTD-conforming XML documents from heterogeneous data sources
Shyh-Kwei Chen, Ming-Ling Lo, Kun-Lung Wu, Jih-Shyr Yih, Colleen Viehrig
Inf. Sci.3
2006 Query indexing with containment-encoded intervals for efficient stream processing
Kun-Lung Wu, Shyh-Kwei Chen, Philip S. Yu
Knowl. Inf. Syst.1
2006 Processing Moving Queries over Moving Objects Using Motion-Adaptive Indexes
abstract
This paper describes a motion-adaptive indexing scheme for efficient evaluation of moving continual queries (MCQs) over moving objects. It uses the concept of motion-sensitive bounding boxes (MSBs) to model moving objects and moving queries. These bounding boxes automatically adapt their sizes to the dynamic motion behaviors of individual objects. Instead of indexing frequently changing object positions, we index less frequently changing object and query MSBs, where updates to the bounding boxes are needed only when objects and queries move across the boundaries of their boxes. This helps decrease the number of updates to the indexes. More importantly, we use predictive query results to optimistically precalculate query results, decreasing the number of searches on the indexes. Motion-sensitive bounding boxes are used to incrementally update the predictive query results. Furthermore, we introduce the concepts of guaranteed safe radius and optimistic safe radius to extend our motion-adaptive indexing scheme to evaluating moving continual k-nearest neighbor (kNN) queries. Our experiments show that the proposed motion-adaptive indexing scheme is efficient for the evaluation of both moving continual range queries and moving continual kNN queries.
Bugra Gedik, Kun-Lung Wu, Philip S. Yu, Ling Liu 0001
IEEE Trans. Knowl. Data Eng.2
2006 Incremental Processing of Continual Range Queries over Moving Objects
abstract
Efficient processing of continual range queries over moving objects is critically important in providing location-aware services and applications. A set of continual range queries, each defining the geographical region of interest, can be periodically (re)evaluated to locate moving objects that are currently within individual query boundaries. We study a new query indexing method, called CES-based indexing, for incremental processing of continual range queries over moving objects. A set of containment-encoded squares (CES) are predefined, each with a unique ID. CESs are virtual constructs (VC) used to decompose query regions and to store indirectly precomputed search results. Compared with a prior VC-based approach, the number of VCs visited in a search operation is reduced from (4L2-1)/3 to log(L)+1, where L is the maximal side length of a VC. Search time is hence significantly lowered. Moreover, containment encoding among the CESs makes it easy to identify all those VCs that need not be visited during an incremental query (re)evaluation. We study the performance of CES-based indexing and compare it with a prior VC-based approach
Kun-Lung Wu, Shyh-Kwei Chen, Philip S. Yu
IEEE Trans. Knowl. Data Eng.1
2005 Adaptive load shedding for windowed stream joins
abstract
We present an adaptive load shedding approach for windowed stream joins. In contrast to the conventional approach of dropping tuples from the input streams, we explore the concept of selective processing for load shedding. We allow stream tuples to be stored in the windows and shed excessive CPU load by performing the join operations, not on the entire set of tuples within the windows, but on a dynamically changing subset of tuples that are learned to be highly beneficial. We support such dynamic selective processing through three forms of runtime adaptations: adaptation to input stream rates, adaptation to time correlation between the streams and adaptation to join directions. Indexes are used to further speed up the execution of stream joins. Experiments are conducted to evaluate our adaptive load shedding in terms of output rate. The results show that our selective processing approach to load shedding is very effective and significantly outperforms the approach that drops tuples from the input streams.
Bugra Gedik, Kun-Lung Wu, Philip S. Yu, Ling Liu 0001
CIKM2
2005 On Incremental Processing of Continual Range Queries for Location-Aware Services and Applications
abstract
A set of continual range queries, each defining the geographical region of interest, can be periodically re-evaluated to locate moving objects. Processing these continual queries efficiently and incrementally hence becomes important for location-aware services and applications. In this paper, we study a new query indexing method, called CES-based indexing, for incremental processing of continual range queries over moving objects. A set of containment-encoded squares (CES) are predefined, each with a unique ID. CES's are virtual constructs (VC) used to decompose query regions and to store indirectly pre-computed search results. Compared with a prior VC-based approach, the number of VC's visited in an index search in CES-based indexing is reduced from (4L/sup 2/-1)/3 to log(L)+1, where L is the maximal side length of a VC. Search time is hence significantly lowered. Moreover, containment encoding among the CES's makes it easy to identify all those VC's that need not be visited during an incremental query reevaluation. We study the performance of CES-based indexing and compare it with a prior VC-based approach.
Kun-Lung Wu, Shyh-Kwei Chen, Philip S. Yu
MobiQuitous1
2005 Fast Burst Correlation of Financial Data
Michail Vlachos, Kun-Lung Wu, Shyh-Kwei Chen, Philip S. Yu
PKDD2
2004 Monitoring Continual Range Queries
Philip S. Yu, Kun-Lung Wu, Shyh-Kwei Chen
APWeb2
2004 Motion adaptive indexing for moving continual queries over moving objects
abstract
This paper describes a motion adaptive indexing scheme for efficient evaluation of moving continual queries (MCQs) over moving objects. It uses the concept of motion-sensitive bounding boxes (MSBs) to model moving objects and moving queries. These bounding boxes automatically adapt their sizes to the dynamic motion behaviors of individual objects. Instead of indexing frequently changing object positions, we index less frequently changing object and query MSBs, where updates to the bounding boxes are needed only when objects and queries move across the boundaries of their boxes. This helps decrease the number of updates to the indexes. More importantly, we use predictive query results to optimistically precalculate query results, decreasing the number of searches on the indexes. Motion-sensitive bounding boxes are used to incrementally update the predictive query results. Our experiments show that the proposed motion adaptive indexing scheme is efficient for the evaluation of moving continual range queries.
Bugra Gedik, Kun-Lung Wu, Philip S. Yu, Ling Liu 0001
CIKM2
2004 Interval query indexing for efficient stream processing
abstract
A large number of continual range queries can be issued against a data stream. Usually, a main memory-based query index with a small storage cost and a fast search time is needed, especially if the stream is rapid. In this paper, we present a CEI-based query index that meets both criteria for efficient processing of continual interval queries in a streaming environment. This new query index is centered around a set of predefined virtual containment-encoded intervals, or CEIs. The CEIs are used to first decompose query intervals and then perform efficient search operations. The CEIs are defined and labeled such that containment relationships among them are encoded in their IDs. The containment encoding makes decomposition and search operations efficient because integer additions and logical shifts can be used to carry out most of the operations. Simulations are conducted to evaluate the effectiveness of the CEI-based query index and to compare it with alternative approaches. The results show that the CEI-based query index significantly outperforms existing approaches in terms of both storage cost and search time.
Kun-Lung Wu, Shyh-Kwei Chen, Philip S. Yu
CIKM1
2004 Processing Continual Range Queries over Moving Objects Using VCR-Based Query Indexes
abstract
This paper describes VCR-based query indexes for efficient processing of continual range queries over moving objects. A set of virtual construct rectangles (VCR) is predefined, each with a unique ID. One or more VCRs is used to strictly cover the entire region defined by a range query. The query index maintains a mapping from each VCR to the range queries that contain that VCR. The use of VCRs provides an indirect and cost-effective way of precomputing the search result for any object position, making possible efficient search operations. More importantly, it allows the processing of continual range queries to capitalize on incremental changes in object locations. Computation can be saved for objects that have not moved out of VCR boundaries. We study different strategies to cover a query region with VCRs and conduct simulations to compare them.
Kun-Lung Wu, Shyh-Kwei Chen, Philip S. Yu
MobiQuitous1
2004 The CHAMPS system: change management with planning and scheduling
abstract
Change management is a process by which IT systems are modified to accommodate considerations such as software fixes, hardware upgrades and performance enhancements. This paper discusses the CHAMPS system, a prototype under development at IBM Research for Change Management with Planning and Scheduling. The CHAMPS system is able to achieve a very high degree of parallelism for a set of tasks by exploiting detailed factual knowledge about the structure of a distributed system from dependency information at runtime. In contrast, today's systems expect an administrator to provide such insights, which is often not the case. Furthermore, the optimization techniques we employ allow the CHAMPS system to come up with a very high quality solution for a mathematically intractable problem in a time which scales nicely with the problem size. We have implemented the CHAMPS system and have applied it in a TPC-W environment that implements an on-line book store application.
Alexander Keller 0002, Joseph L. Hellerstein, Joel L. Wolf, Kun-Lung Wu, Vijaya Krishnan
NOMS (1)4
2004 Segmentation of multimedia streams for proxy caching
abstract
Proxy caching of large multimedia objects on the edge of the Internet has become increasingly important for reducing network latency. For a large media object, such as a two-hour video, treating the whole media as a single object for caching is not appropriate. In this paper, we study three media segmentation approaches to proxy caching: fixed, pyramid, and skyscraper. Blocks of a media stream are grouped into various segments for cache management. The cache admission and replacement policies attach different caching priorities to individual segments, taking into account the access frequency of the media object and the segment distance from the start of the media. These caching policies give preferential treatment to the beginning segments. As such, most user requests can be quickly played back from the proxy servers without delay. Event-driven simulations are conducted to evaluate the segmentation approaches and compare them with whole media caching. The results show that: 1) compared with whole media caching, segmentation-based caching is more effective not only in increased byte-hit ratio but also in lowered fraction of requests that requires delayed start; 2) pyramid segmentation, where segment size increases exponentially, is the best segmentation approach; and 3) segmentation-based caching is especially advantageous when the cache size is limited, when the set of hot media objects changes over time, when the media file size is large, and when there are a large number of distinct media objects.
Kun-Lung Wu, Philip S. Yu, Joel L. Wolf
IEEE Trans. Multim.1
2003 Epi-SPIRE: a system for environmental and public health activity monitoring
abstract
Health activity monitoring (HAM) has received increasing attention due to the rapid advances of both hardware and software technologies and strong environmental and public health needs. In this paper, we describe the architecture and implementation of the Epi-SPIRE prototype, which is a novel health activity monitoring system that generates alerts from environmental, behavioral, and public health data sources. A model-based approach is used to develop disease and behavior models from multi-modal heterogeneous data sources. Furthermore, a model-based indexing technique has been developed to speed up the data access and retrieval. This system has been successfully applied to various genuine and simulated diseases outbreaks scenarios'.
Chung-Sheng Li, Charu C. Aggarwal, Murray Campbell, Yuan-Chi Chang, Gregory Glass, Vijay S. Iyengar, Mahesh Joshi, Ching-Yung Lin, Milind R. Naphade, John R. Smith, Belle L. Tseng, Min Wang 0001, Kun-Lung Wu, Philip S. Yu
ICME13
2003 Replication for Load Balancing and Hot-Spot Relief on Proxy Web Caches with Hash Routing
Kun-Lung Wu, Philip S. Yu
Distributed Parallel Databases1
2003 Optimizing Index Allocation for Sequential Data Broadcasting in Wireless Mobile Computing
abstract
Energy saving is one of the most important issues in wireless mobile computing. Among others, one viable approach to achieving energy saving is to use an indexed data organization to broadcast data over wireless channels to mobile units. Using indexed broadcasting, mobile units can be guided to the data of interest efficiently and only need to be actively listening to the broadcasting channel when the relevant information is present. We explore the issue of indexing data with skewed access for sequential broadcasting in wireless mobile computing. We first propose methods to build index trees based on access frequencies of data records. To minimize the average cost of index probes, we consider two cases: one for fixed index fanouts and the other for variant index fanouts, and devise algorithms to construct index trees for both cases. We show that the cost of index probes can be minimized not only by employing an imbalanced index tree that is designed in accordance with data access skew, but also by exploiting variant fanouts for index nodes. Note that, even for the same index tree, different broadcasting orders of data records will lead to different average data access times. To address this issue, we develop an algorithm to determine the optimal order for sequential data broadcasting to minimize the average data access time. Performance evaluation on the algorithms proposed is conducted. Examples and remarks are given to illustrate our results.
Ming-Syan Chen, Kun-Lung Wu, Philip S. Yu
IEEE Trans. Knowl. Data Eng.2
2002 Efficient query monitoring using adaptive multiple key hashing
abstract
Monitoring continual queries or subscriptions is to determine the subset of all queries or subscriptions whose predicates match a given event. Predicates contain not only equality but also non-equality clauses. Event matching is usually accomplished by first identifying a "small" candidate set of subscriptions for an event and then determining the matched subscriptions from the candidate set. Prior work has focused on using equality clauses to identify the candidate set. However, we found that completely ignoring non-equality clauses can result in a much larger candidate set. In this paper, we present and evaluate an adaptive multiple key hashing (AMKH) method to judiciously include an effective subset of non-equality clauses in candidate set identification. Each subscription is mapped to a data point in a multidimensional space based on its predicate clauses. AMKH is then used to maintain subscriptions and perform event matching. AMKH further provides a controlling mechanism to limit the hash range of a non-equality clause, hence reducing the size of the candidate set. Simulations are conducted to study the performance of AMKH. The results show that (1) a small number of non-equality clauses can be effectively included by AMKH and (2) the attributes whose overall non-equality predicate clauses are most selective should be chosen for inclusion by AMKH.
Kun-Lung Wu, Philip S. Yu
CIKM1
2001 Segment-based proxy caching of multimedia streams
abstract
As streaming video and audio over the Internet becomes popular, proper proxy caching of large multimedia objects has become increasingly important. For a large media object, such as a 2-hour video, treating the whole video as a single web object for caching is not appropriate. In this paper, we present and evaluate a segment-based bu er management approach to proxy caching of large media streams. Blocks of a media stream received by a proxy server are grouped into variable-sized segments. The cache admission and replacement policies then attach di erent caching values to di erent segments, taking into account the segment distance from the start of the media. These caching policies give preferential treatments to the beginning segments. As such, users can quickly play back the media objects without much delay. Event-driven simulations are conducted to evaluate this segment-based proxy caching approach. The results show that (1) segment-based caching is e ective not only in increasing byte-hit ratio (or reducing total traAEc) but also in lowering the number of requests that require delayed starts; (2) segment-based caching is especially advantageous when the cache size is limited, when the set of hot media objects changes over time, when the media le size is large, and when many users may stop playing the media after only a few initial blocks.
Kun-Lung Wu, Philip S. Yu, Joel L. Wolf
WWW1
2000 TabSum: A Flexible and Dynamic Table Summarization Approach
abstract
Many diverse small devices, such as smart phones and personal digital assistants, are being deployed to access the Internet. Small devices typically have limited display and processing capabilities. It is thus difficult to effectively present a large table of information on these different devices such that the table can be easily browsed by the users. We present the design and prototype implementation of TabSum, a flexible and dynamic table summarization approach to reducing tables into smaller but still meaningful representations for various devices. TabSum supports easy browsing of a large table by various devices and provides personalization by taking into account both device capabilities and user preferences. A multi-source multi-layered specification methodology is proposed to describe summarization rules preferred by various sources. On request, a set of table reduction rules are dynamically applied to the rows and columns of a table to reduce its size. The approach is flexible and dynamic. It summarizes both numeric and non-numeric data types.
Ming-Ling Lo, Kun-Lung Wu, Philip S. Yu
ICDCS2
2000 Latency-sensitive hashing for collaborative Web caching
Kun-Lung Wu, Philip S. Yu
Comput. Networks1
2000 Workfile Disk Management for Concurrent Mergesorts in a Multiprocessor Database System
Kun-Lung Wu, Philip S. Yu, Jen-Yao Chung, James Z. Teng
Distributed Parallel Databases1
1999 Local Replication for Proxy Web Caches with Hash Routing
abstract
This paper studies controlled local replication for hash routing, such as CARP, among a collection of loosely-coupled proxy web cache servers. Hash routing partitions the entire URL space among the shared web caches, creating a single logical cache. Each partition is assigned to a cache server. Duplication of cache contents is eliminated and total incoming traffic to the shared web caches is minimized. Client requests for non-assigned-partition objects are forwarded to sibling caches. However, request forwarding increases not only inter-cache traffic but also cpu utilization, thus slows the client response time. We propose a controlled local replication of non-assigned-partition objects in each cache server to effectively reduce the inter-cache traffic. We use a multiple-exit LRU to implement controlled local replication. Trace-driven simulations are conducted to study the performance impact of local replication. The results show that (1) regardless of cache sizes, with a controlled local replication, the average response time, inter-cache traffic and CPU overhead can be effectively reduced without noticeable increases in incoming traffic; (2) for very large cache sizes, a larger amount of local replication can be allowed to reduce inter-cache traffic without increasing incoming traffic; and (3) local replication is effective even if clients are dynamically assigned to different cache servers.
Kun-Lung Wu, Philip S. Yu
CIKM1
1999 Load Balancing and Hot Spot Relief for Hash Routing among a Collection of Proxy Caches
abstract
Hash routing partitions the entire URL space among a collection of cooperating proxy caches. Each partition is assigned to a cache server. Duplication of cache contents is eliminated. Client requests to a cache server for non-assigned partition objects are forwarded to proper sibling caches. As a result, the load level of the cache servers can be quite unbalanced. We examine an adaptable controlled replication (ACR) of non-assigned partition objects in each cache server to reduce the load imbalance and relieve the problem of hot-spot references. Trace-driven simulations are conducted to study the effectiveness of ACR. The results show that: (1) access skew exists, and the load of the cache servers tends to be unbalanced in hash routing; (2) with a relatively small amount of ACR, say 10% of the cache size, significant improvements in load balance can be achieved; and (3) ACR provides a very effective remedy for load imbalance due to hot-spot references.
Kun-Lung Wu, Philip S. Yu
ICDCS1
1999 Horting Hatches an Egg: A New Graph-Theoretic Approach to Collaborative Filtering
abstract
Article Free Access Share on Horting hatches an egg: a new graph-theoretic approach to collaborative filtering Authors: Charu C. Aggarwal IBM T. J. Watson Research Center, Yorktown Heights, NY IBM T. J. Watson Research Center, Yorktown Heights, NYView Profile , Joel L. Wolf IBM T. J. Watson Research Center, Yorktown Heights, NY IBM T. J. Watson Research Center, Yorktown Heights, NYView Profile , Kun-Lung Wu IBM T. J. Watson Research Center, Yorktown Heights, NY IBM T. J. Watson Research Center, Yorktown Heights, NYView Profile , Philip S. Yu IBM T. J. Watson Research Center, Yorktown Heights, NY IBM T. J. Watson Research Center, Yorktown Heights, NYView Profile Authors Info & Claims KDD '99: Proceedings of the fifth ACM SIGKDD international conference on Knowledge discovery and data miningAugust 1999 Pages 201–212https://doi.org/10.1145/312129.312230Published:01 August 1999Publication History 218citation2,375DownloadsMetricsTotal Citations218Total Downloads2,375Last 12 Months144Last 6 weeks44 Get Citation AlertsNew Citation Alert added!This alert has been successfully added and will be sent to:You will be notified whenever a record that you have chosen has been cited.To manage your alert preferences, click on the button below.Manage my AlertsNew Citation Alert!Please log in to your account Save to BinderSave to BinderCreate a New BinderNameCancelCreateExport CitationPublisher SiteeReaderPDF
Charu C. Aggarwal, Joel L. Wolf, Kun-Lung Wu, Philip S. Yu
KDD3
1999 Run Placement Policies for Concurrent Mergesorts Using Parallel Prefetching
Kun-Lung Wu, Philip S. Yu, James Z. Teng
Knowl. Inf. Syst.1
1998 Range-Based Bitmap Indexing for High Cardinality Attributes with Skew
abstract
Bitmap indexing, though effective for low cardinality attributes, can be rather costly in storage overhead for high cardinality attributes. Range-based bitmap (RBM) indexing can be used to reduce this storage overhead. The attribute values are partitioned into ranges and a bitmap vector is used to represent a range. With RBM, however, the number of records assigned to different ranges can be highly uneven, resulting in non-uniform search times for different queries. We present and evaluate a dynamic bucket expansion and contraction (DBEC) approach to simultaneously constructing range-based bitmap indexes for multiple high-cardinality attributes. Simulations are conducted to evaluate this DBEC approach. Both synthetic and real data are used in the simulations. The results show that (1) with highly skewed data, DBEC performs quite well compared with a simple approach and (2) DBEC compares favorably with the optimal approach.
Kun-Lung Wu, Philip S. Yu
COMPSAC1
1998 Energy-Efficient Mobile Cache Invalidation
Kun-Lung Wu, Philip S. Yu, Ming-Syan Chen
Distributed Parallel Databases1
1998 Increasing Multimedia System Throughput with Consumption-based Buffer Management
Kun-Lung Wu, Philip S. Yu
Multim. Syst.1
1997 Indexed Sequential Data Broadcasting in Wireless Mobile Computing
abstract
Energy saving is one of the most important issues in wireless mobile computing. Among others, one viable approach to achieving energy saving is to use an indexed data organization to broadcast data over wireless channels to mobile units. We explore the issue of indexing data with skewed access for sequential broadcasting in wireless mobile computing. We propose methods to build index trees based on access frequencies of data records. To minimize the average cost of index probes, we consider two cases: one for fixed index fanouts and the other for variant index fanouts, and devise algorithms to construct index trees for both cases. We show that the cost of index probes can be minimized not only by employing an imbalanced index tree that is designed in accordance with data access skew, but also by exploiting variant fanouts for index nodes.
Ming-Syan Chen, Philip S. Yu, Kun-Lung Wu
ICDCS3
1997 Divergence Control Algorithms for Epsilon Serializability
abstract
The paper presents divergence control methods for epsilon serializability (ESR) in centralized databases. ESR alleviates the strictness of serializability (SR) in transaction processing by allowing for limited inconsistency. The bounded inconsistency is automatically maintained by divergence control (DC) methods in a way similar to SR is maintained by concurrency control (CC) mechanisms. However, DC for ESR allows more concurrency than CC for SR. The authors first demonstrate the feasibility of ESR by showing the design of three representative DC methods: two-phase locking, timestamp ordering and optimistic approaches. DC methods are designed by systematically enhancing CC algorithms in two stages: extension and relaxation. In the extension stage, a CC algorithm is analyzed to locate the places where it identifies non-SR conflicts of database operations. In the relaxation stage, the non-SR conflicts are relaxed to allow for controlled inconsistency. They then demonstrate the applicability Of ESR by presenting the design of DC methods using other most known inconsistency specifications, such as absolute value, age and total number of nonserializably read data items. In addition, they present a performance study using an optimistic divergence control algorithm as an example to show that a substantial improvement in concurrency can be achieved in ESR by allowing for a small amount of inconsistency.
Kun-Lung Wu, Philip S. Yu, Calton Pu
IEEE Trans. Knowl. Data Eng.1
1996 Energy-Efficient Caching for Wireless Mobile Computing
abstract
Caching can reduce the bandwidth requirement in a mobile computing environment. However, due to battery power limitations, a wireless mobile computer may often be forced to operate in a doze (or even totally disconnected) mode. As a result, the mobile computer may miss some cache invalidation reports broadcast by a server, forcing it to discard the entire cache contents after waking up. In this paper, we present an energy-efficient cache invalidation method, called GCORE (Grouping with COld update-set REtention), that allows a mobile computer to operate in a disconnected mode to save the battery while still retaining most of the caching benefits after a reconnection. We present an efficient implementation of GCORE and conduct simulations to evaluate its caching effectiveness. The results show that GCORE can substantially improve mobile caching by reducing the communication bandwidth (or energy consumption) for query processing.
Kun-Lung Wu, Philip S. Yu, Ming-Syan Chen
ICDE1
1996 Optimization of Parallel Execution for Multi-Join Queries
abstract
We study the subject of exploiting interoperator parallelism to optimize the execution of multi-join queries. Specifically, we focus on two major issues: (1) scheduling the execution sequence of multiple joins within a query, and (2) determining the number of processors to be allocated for the execution of each join operation obtained in (1). For the first issue, we propose and evaluate by simulation several methods to determine the general join sequences, or bushy trees. Despite their simplicity, the heuristics proposed can lead to the general join sequences that significantly outperform the optimal sequential join sequence. The quality of the join sequences obtained by the proposed heuristics is shown to be fairly close to that of the optimal one. For the second issue, it is shown that the processor allocation for exploiting interoperator parallelism is subject to more constraints-such as execution dependency and system fragmentation-than those in the study of intraoperator parallelism for a single join. The concept of synchronous execution time is proposed to alleviate these constraints. Several heuristics to deal with the processor allocation, categorized by bottom-up and top-down approaches, are derived and are evaluated by simulation. The relationship between issues (1) and (2) is explored. Among all the schemes evaluated, the two-step approach proposed, which first applies the join sequence heuristic to build a bushy tree as if under a single processor system, and then, in light of the concept of synchronous execution time, allocates processors to execute each join in the bushy tree in a top-down manner, emerges as the best solution to minimize the query execution time.
Ming-Syan Chen, Philip S. Yu, Kun-Lung Wu
IEEE Trans. Knowl. Data Eng.3
1996 Performance Analysis of Dynamic Finite Versioning Schemes: Storage Cost vs. Obsolescence
abstract
Dynamic finite versioning (DFV) schemes are an effective approach to concurrent transaction and query processing, where a finite number of consistent, but maybe slightly out-of-date, logical snapshots of the database can be dynamically derived for query access. In DFV, the storage overhead for keeping additional versions of changed data to support the logical snapshots and the amount of obsolescence faced by queries are two major performance issues. We analyze the performance of DFV, with emphasis on the trade-offs between the storage cost and obsolescence. We develop analytical models based on a renewal process approximation to evaluate the performance of DFV using M/spl ges/2 snapshots. Asymptotic closed form results for high query arrival rates are given for the case of two snapshots. Simulation is used to validate the analytical models and to evaluate the tradeoffs between various strategies for advancing snapshots when M>2. The results show that (1) the analytical models match closely with simulation; (2) storage cost and obsolescence are sensitive to the snapshot advancing strategies, and (3) usually, increasing the number of snapshots demonstrates a trade-off between storage overhead and query obsolescence. For cases with skewed accessor low update rates, a small increase in the number of snapshots beyond two can substantially reduce the obsolescence. Such a reduction in obsolescence is more significant as the coefficient of variation of the query length distribution becomes larger. Moreover, for very low update rates, a large number of snapshots can be used to reduce the obsolescence to almost zero without increasing the storage overhead.
Arif Merchant, Kun-Lung Wu, Philip S. Yu, Ming-Syan Chen
IEEE Trans. Knowl. Data Eng.2
1995 A Performance Study of Workfile Disk Management for Concurrent Mergesorts in a Multiprocessor Database System
Kun-Lung Wu, Philip S. Yu, Jen-Yao Chung, James Z. Teng
VLDB1
1995 Divergence Control for Distributed Database Systems
Wenwey Hseush, Gail E. Kaiser, Calton Pu, Kun-Lung Wu, Philip S. Yu
Distributed Parallel Databases4
1994 Multiversion Divergence Control of Time Fuzziness
abstract
Epsilon Serializability (ESR) has been proposed to manage and control inconsistency in extending the classic transaction processing. ESR increases system concurrency by tolerating a bounded amount of inconsistency. In this paper, we present multiversion divergence control (mvDC) algorithms that support ESR with not only value but also time fuzziness in multiversion databases. Unlike value fuzziness, accumulating time fuzziness is semantically different. A simple summation of the length of two time intervals may either underestimate the total time fuzziness, resulting in incorrect execution, or overestimate the total time fuzziness, unnecessarily degrading the effectiveness of mvESR. We present a new operation, called TimeUnion, to accurately accumulate the total time fuzziness. Because of the accurate control of time and value fuzziness by the mvDC algorithm, mvESR is very suitable for the use of multiversion databases for real-time applications that may tolerate a limited degree of data inconsistency but prefer more data recency.
Calton Pu, Miu K. Tsang, Kun-Lung Wu, Philip S. Yu
CIKM3
1994 Data Placement and Buffer Management for Concurrent Mergesorts with Parallel Prefetching
abstract
Various data placement policies are studied for the merge phase of concurrent mergesorts using parallel prefetching, where initial sorted runs (input) of a merge and its final sorted run (output) are stored on multiple disks but each run resides only on a single disk. Since the merge phase involves only sequential references, parallel prefetching can be attractive an reducing the average response time for concurrent merges. However, without careful buffer control, severe thrashing may develop under certain run placement policies, reducing the benefits of prefetching. The authors examine through detailed simulations three different run placement policies. The results show that even though buffer thrashing can be almost avoided by placing the output run of a job on the same disk with at least one of its input runs, this thrashing-avoiding run placement policy can be substantially outperformed by other policies that use buffer thrashing control. With buffer thrashing avoidance, the best performance as achieved by a run placement policy that uses a proper subset of disks dedicated for writing the output runs while the rest of the disks are used for prefetching the input runs in parallel.>
Kun-Lung Wu, Philip S. Yu, James Z. Teng
ICDE1
1994 Dynamic Parity Grouping for Improving Write Performance of RAID-5 Disk Arrays
abstract
One major drawback of a RAIDS disk array system is that an update to a data block may involve four disk accesses. Such a high overhead is especially undesirable for workloads with a high update rate. In this paper, we present a dynamic parity grouping (DPG) scheme for efficient parity buffering to reduce the write overhead of a RAID-5 system. In DPG, special parity groups are dynamically created for data blocks with high write activity, referred to as the hot data blocks, in addition to default parity groups for the remaining cold data blocks. The parity blocks of the special parity groups are then buffered in the disk controller cache. As a result, the number of disk accesses on a write to a hot data block is reduced to two.
Philip S. Yu, Kun-Lung Wu, Asit Dan
ICPP (2)2
1994 On real-time databases: concurrency control and scheduling
abstract
In addition to maintaining database consistency as in conventional databases, real-time database systems must also handle transactions with timing constraints. While transaction response time and throughput are usually used to measure a conventional database system, the percentage of transactions satisfying the deadlines or a time-critical value function is often used to evaluate a real-time database system. Scheduling real-time transactions is far more complex than traditional real-time scheduling in the sense that (1) worst case execution times are typically hard to estimate, since not only CPU but also I/O requirement is involved; and (2) certain aspects of concurrency control may not integrate well with real-time scheduling. In this paper, we first develop a taxonomy of the underlying design space of concurrency control including the various techniques for achieving serializability and improving performance. This taxonomy provides us with a foundation for addressing the real-time issues. We then consider the integration of concurrency control with real-time requirements. The implications of using run policies to better utilize real-time scheduling in a database environment are examined. Finally, as timing constraints may be more important than data consistency in certain hard realtime database applications, we also discuss several approaches that explore the nonserializable semantics of real-time transactions to meet the hard deadlines.>
Philip S. Yu, Kun-Lung Wu, Kwei-Jay Lin, Sang Hyuk Son
Proc. IEEE2
1994 Optimal NODUP All-to-All Broadcast Schemes in Distributed Computing Systems
abstract
Broadcast, referring to a process of information dissemination in a distributed system whereby a message originating from a certain node is sent to all other nodes in the system, is a very important issue in distributed computing. All-to-all broadcast means the process by which every node broadcasts its certain piece of information to all other nodes. In this paper, we first develop the optimal all-to-all broadcast scheme for the case of one-port communication, which means that each node can only send out one message in one communication step, and then, extend our results to the case of multi-port communication, i.e., k-port communication, meaning that each node can send out k messages in one communication step. We prove that the proposed schemes are optimal for the model considered in the sense that they not only require the minimal number of communication steps, but also incur the minimal number of messages.>
Ming-Syan Chen, Philip S. Yu, Kun-Lung Wu
IEEE Trans. Parallel Distributed Syst.3
1993 Decentralized Consensus Protocols with Multi-Port Communication
abstract
The authors develop efficient decentralized consensus protocols for a distributed system with multi-port communication. Two classes of decentralized consensus protocols are considered: the one without an initiator and the one with an initiator. The case of one-port communication is first presented, i.e., each node can send out one message in one step, and then results are derived for the case of multi-port communication, i.e., each node can send out more than one message in one step. Given an arbitrary number of nodes in a system, the proposed protocols can reach the consensus in the minimal numbers of message steps. The number of messages incurred by each algorithm is also derived.>
Ming-Syan Chen, Philip S. Yu, Kun-Lung Wu
ICDCS3
1993 Distributed Divergence Control for Epsilon Serializability
abstract
Epsilon serializability (ESR) allows for more concurrency by permitting nonserializable interleavings of database operations among epsilon transactions (ETs). The authors present the design of distributed divergence control (DDC) algorithms for ESR in homogeneous and heterogeneous distributed databases. They first present a strict two-phase locking DDC algorithm (S2PLDDC) and an optimistic DDC algorithm (ODDC) for homogeneous distributed databases, where the local orderings of all the sub-ETs of a distributed ET are the same, and the total inconsistency of a distributed ET is simply the sum of that of all its sub-ETs. A superdatabase DDC algorithm is described for heterogeneous distributed databases, where the local orderings of all the sub-ETs of a distributed ET may not be the same, and the total inconsistency of a distributed ET may be greater than the sum of that of all its sub-ETs. As a result, in addition to local divergence control in each site, a global mechanism is needed to guarantee ESR.>
Calton Pu, Wenwey Hseush, Gail E. Kaiser, Kun-Lung Wu, Philip S. Yu
ICDCS4
1993 Dynamic Finite Versioning: An Effective Versioning Approach to Concurrent Transaction and Query Processing
abstract
Dynamic finite versioning (DFV) schemes that effectively support concurrent processing of transaction and queries are presented. Without acquiring locks, queries read from a small, fixed number of dynamically derived, transaction-consistent, possibly slightly obsolete, logical snapshots of the database. On the other hand, transactions access the most up-to-date data in the database without data contention from queries. Intermediate versions created between snapshots are automatically discarded. Dirty pages updated by active transactions are allowed to be written back into the database before commitment and, at the same time, consistent logical snapshots can be advanced automatically without quiescing the ongoing transactions or queries.>
Kun-Lung Wu, Philip S. Yu, Ming-Syan Chen
ICDE1
1993 Performance Comparison of Thrashing Control Policies for Concurrent Mergesorts with Parallel Prefetching
abstract
We study the performance of various run-time thrashing control policies for the merge phase of concurrent mergesorts using parallel prefetching, where initial sorted runs are stored on multiple disks and the final sorted run is written back to another dedicated disk. Parallel prefetching via multiple disks can be attractive in reducing the response times for concurrent mergesorts. However, severe thrashing may develop due to imbalances between input and output rates, thus a large number of prefetched pages in the buffer can be replaced before referenced. We evaluate through detailed simulations three run-time thrashing control policies: (a) disabling prefetching, (b) forcing synchronous writes and (c) lowering the prefetch quantity in addition to forcing synchronous writes. The results show that (1) thrashing resulted from parallel prefetching can severely degrade the system response time; (2) though effective in reducing the degree of thrashing, disabling prefetching may worsen the response time since more synchronous reads are needed; (3) forcing synchronous writes can both reduce thrashing and improve the response time; (4) lowering the prefetch quantity in addition to forcing synchronous writes is most effective in reducing thrashing and improving the response time.
Kun-Lung Wu, Philip S. Yu, James Z. Teng
SIGMETRICS1
1993 Performance comparison of dynamic policies for remote caching
abstract
Abstract In a distributed system, data servers (file systems and databases) can easily become bottlenecks. We propose an approach to offloading data access requests from overloaded data servers to nodes that are Idle or less busy. This approach is referred to asremote caching, and the idle or less busy nodes are calledmutual serversas they help out the busy server nodes on data accesses. In addition to server and client local caches, frequently accessed data are cached in the main memory of mutual servers, thus improving the data access time in the system. We evaluate several data propagation strategics among data servers and mutual servers. These include policies in which senders are active/passive and receivers are active/passive in initiating data propagation. For example, an active sender takes the initiative to offload data onto a passive receiver. Simulation results show that the active‐sender/passive‐receiver policy is the method of choice In most cases. Active‐Sender policies are best able to exploit the main memory of other Idle nodes in the expected normal condition where some nodes are overloaded and others are less loaded. AH active policies perform far better than the policy without remote caching even in the degenerated case where each node is equally loaded.
Calton Pu, Danilo Florissi, Patricia Soares, Philip S. Yu, Kun-Lung Wu
Concurr. Pract. Exp.5
1993 Rapid Transaction-Undo Recovery Using Twin-Page Storage Management
abstract
A twin-page storage method, which is an alternative to the TWIST (twin slot) approach by A. Reuter
Kun-Lung Wu, W. Kent Fuchs
IEEE Trans. Software Eng.1
1992 Performance Comparison of Active-Sender and Active-Receiver Policies for Distributed Caching
abstract
The authors propose a distributed caching approach to off-loading data access requests from overloaded data servers in a distributed system to nodes that are idle or less busy. Helping out the busy servers on data accesses, the idle or less busy nodes are called mutual servers. Frequently accessed data are cached in the main memory of mutual servers in addition to server and client local caches. The authors evaluate several data propagation strategies among data servers and mutual servers. Simulation results show that the active-sender passive-receiver policy is the method of choice in most cases. Active-sender policies are best able to exploit the main memory of other idle nodes in the expected normal condition where some nodes are overloaded and others are less loaded. All active policies perform far better than the policy without distributed caching.>
Calton Pu, Danilo Florissi, Patricia Soares, Kun-Lung Wu, Philip S. Yu
HPDC4
1992 Efficient Decentralized Consensus Protocols in a Distributed Computing System
abstract
Two classes of efficient decentralized consensus protocols for a distributed computing system consisting of an arbitrary number of nodes, one without an initiator and the other with an initiator, are described. It is shown that the protocol without an initiator can be systematically executed and completed in the minimal number of steps. The protocol with an initiator is divided into three phases: broadcasting phase, shuffling phase, and confirming phase. It is proved that under the protocol with initiator, a distributed system of p nodes reaches consensus with an initiator in the minimal number of steps required. The total number of messages required by the protocol with initiator is derived.>
Ming-Syan Chen, Kun-Lung Wu, Philip S. Yu
ICDCS2
1992 Scheduling and Processor Allocation for Parallel Execution of Multi-Join Queries
abstract
The authors deal with two major issues to exploit inter-operator parallelism within a multijoin query: join sequence scheduling and processor allocation. For the first issue, they propose and evaluate by simulation several methods to determine the general join sequences, or the bush execution trees. Despite their simplicity, the proposed heuristics can lead to the general join sequences which significantly outperform the optimal sequential join sequence. In addition, several heuristics to determine the processor allocation, categorized by bottom-up and top-down approaches, were derived and evaluated by simulation. As confirmed by the simulation, by first using the join sequence heuristics to build a busy tree and then applying the concept of synchronous execution time to the busy tree for processor allocation, an efficient two-step approach to schedule and execute multijoin queries in a multiprocessor system can be obtained.>
Ming-Syan Chen, Philip S. Yu, Kun-Lung Wu
ICDE3
1992 Divergence Control for Epsilon-Serializability
abstract
The authors present divergence control methods for epsilon-serializability (ESR) in centralized databases. ESR alleviates the strictness of serializability (SR) in transaction processing by allowing for limited inconsistency. The bounded inconsistency is automatically maintained by divergence control (DC) methods in a way similar to the manner in which SR is maintained by concurrency control mechanisms, but DC for ESR allows more concurrency. Concrete representative instances of divergence-control methods are described based on two-phase locking, timestamp ordering, and optimistic approaches. The applicability of ESR is demonstrated by presenting the designs of DC methods using other most known inconsistency specifications, such as absolute value, age, and total number of nonserializably read data items.>
Kun-Lung Wu, Philip S. Yu, Calton Pu
ICDE1
1992 Performance Analysis of Dynamic Finite Versioning for Concurrency Transaction and Query Processing
abstract
In this paper, we analyze the performance of dynamic finite versioning (DFV) schemes for concurrent transaction and query processing, where a finite number of consistent snapshots can be derived for query access. We develop analytical models based on a renewal process approximation to evaluate the performance of DFV using M ≥ 2 snapshots. The storage overhead and obsolescence faced by queries are measured. Simulation is used to validate the analytical models and to evaluate the trade-offs between various starategies for advancing snapshots when M > 2.
Arif Merchant, Kun-Lung Wu, Philip S. Yu, Ming-Syan Chen
SIGMETRICS2
1990 Twin-page storage management for rapid transaction-undo recovery
abstract
A twin-page disk-storage management scheme for rapid database transaction-undo recovery is presented and evaluated. In contrast to previous twin-page schemes, this approach uses static page mapping and allows dirty pages in the main memory to be written, at any instant, onto disk without the requirement of undo logging. No explicit undo is required when a transaction is aborted. Transaction undo is implicitly performed by not subsequently fetching from disk the invalid pages updated by the aborted transaction. Performance in terms of disk I/O and CPU overhead for transaction-undo recovery is analyzed and compared with a previous approach called TWIST. It is shown that the present scheme achieves rapid transaction-undo recovery without degrading average system performance for various workloads, and that the scheme is well suited for applications with a large number of updates and frequent transaction aborts.>
Kun-Lung Wu, W. Kent Fuchs
COMPSAC1
1990 Recoverable Distributed Shared Virtual Memory
abstract
The problem of rollback recovery in distributed shared virtual environments, in which the shared memory is implemented in software in a loosely coupled distributed multicomputer system, is examined. A user-transparent checkpointing recovery scheme and a new twin-page disk storage management technique are presented for implementing recoverable distributed shared virtual memory. The checkpointing scheme can be integrated with the memory coherence protocol for managing the shared virtual memory. The twin-page disk design allows checkpointing to proceed in an incremental fashion without an explicit undo at the time of recovery. The recoverable distributed shared virtual memory allows the system to restart computation from a checkpoint without a global restart.>
Kun-Lung Wu, W. Kent Fuchs
IEEE Trans. Computers1
1990 Error Recovery in Shared Memory Multiprocessors Using Private Caches
abstract
The problem of recovering from processor transient faults in shared memory multiprocessor systems is examined. A user-transparent checkpointing and recovery scheme using private caches is presented. Processes can recover from errors due to faulty processors by restarting from the checkpointed computation state. Implementation techniques using checkpoint identifiers and recovery stacks are examined as a means of reducing performance degradation in processor utilization during normal execution. This cache-based checkpointing technique prevents rollback propagation, provides rapid recovery, and can be integrated into standard cache coherence protocols. An analytical model is used to estimate the relative performance of the scheme during normal execution. Extensions to take error latency into account are presented.>
Kun-Lung Wu, W. Kent Fuchs, Janak H. Patel
IEEE Trans. Parallel Distributed Syst.1
1989 Cache-Based Error Recovery for Shared Memory Multiprocessor Systems
Kun-Lung Wu, W. Kent Fuchs, Janak H. Patel
ICPP (1)1
1987 Comparison and Diagnosis of Large Replicated Files
abstract
This paper examines the problem of comparing large replicated files in a context in which communication dominates the cost of comparison. A low-cost checking matrix is proposed for comparison of these replicated files. The checking matrix is composed of check symbols generated by a divide-and-conquer encoding algorithm. The matrix allows for detection and diagnosis of disagreeing pages with very little communication overhead. In contrast to a previous O(N) proposal, the storage requirement for the checking matrix is O(log N), where N is the number of pages in the file. The matrix can be stored in main memory without the need for extra accesses to disk during normal updates of pages.
W. Kent Fuchs, Kun-Lung Wu, Jacob A. Abraham
IEEE Trans. Software Eng.2