Software Cache Coherence with Memory Scaling
Nikos Hardavellas, Leonidas I. Kontothanassis, R. Nikhil, Robert J. Stets · 1998
Software-coherent, distributed shared memory has received conciderable amount of attention as an attractive programming model that leverages inexpensive hardware. DSM systems running on clusters of symmetric multiprocessors provide good performance for a large number of applications and provide an easy path to incremental scalability. Most DSM systems [3, 1, 19] use the main memory of each node in the cluster as a third level cache and they migrate and replicate data in that memory. Since computer memories tend to be much larger than caches DSM systems have largely ignored memory capacity issues, assuming there is always enough space in main memory into which to replicate data. This design decision places a hard limit on the scalability of DSM applications running on clusters. While adding more nodes in the cluster increases the amount of computational power available to the application it has no effect on the amount of shared memory the application can use, thus limiting the sizes of problems that can be tackled with a DSM approach. In early studies we have seen significant performance degradation when an application attempts to access more memory than is available on a single node of the cluster. We have designed a new protocol based on the Cashmere DSM system that attempts to take advantage of all the memory available in a cluster. The key insight behind the protocol is that while a node may access a larger amount of memory than it has locally, at any point in time the working set is likely to be smaller than the amount of local memory. We have therefore augmented the protocol with the capability to evict coherence blocks when memory pressure gets high. Our protocol takes advantage of DSM knowledge in order to minimize the cost of evicting page from a node. In particular we catagorize pages in three categories: