A multi-level caching architecture for stateful stream computation

Muhammed Tawfiqul Islam, Renata Borovica‐Gajic, Shanika A. Karunasekera · 2022

Stream processing is used for real-time applications that deal with large volumes, velocities, and varieties of data. Stream processing frameworks discretize continuous data streams to apply computations on smaller batches. For real-time stream-based data analytics algorithms, the intermediate states of computations might need to be retained in memory until the query is complete. Thus, a massive surge in memory demand needs to be satisfied to run these algorithms successfully. However, a worker/server node in a computing cluster may have limited memory capacity. In addition, multiple parallel processes might be running concurrently, sharing the primary memory. As a result, a streaming application might fail to run or complete due to a memory shortage. Although spilling state information to the disk can alleviate the problem by allowing the query to finish, it will cause significant performance overhead. An in-memory-based object store as the state backend will also perform poorly due to the added communications with the external object store and serializing/deserializing the objects. This paper proposes a multi-level caching architecture to mitigate the surge of memory demand from the processes running complex stateful streaming applications. The multiple levels of the cache span across the process heap space, in-memory distributed object store, and secondary storage. The objects/states required in the computation are always served from the fastest level of the cache to boost the application performance. We also provide a multi-level caching library in Java which can be used to implement scalable streaming algorithms. The underlying cache management completely abstracts the multi-level cache implementation from the application and handles seamless migration of states/objects across different levels of the cache. In addition, the multi-layer caching architecture is configurable (i.e., an application can choose to leave out a cache level.) The experimental results demonstrate that our proposed multi-level caching approach for state management can manage large computational windows and improve the performance of an actual streaming application up to three times compared to the in-memory object store-based state backend.

Read the paper · More papers on PaperTik