Previous Table of Contents Next


More recent designs of NUMA systems have overcome the high ratio of remote-to-local cache miss times integral to the earlier designs, so that changes to applications and the redistribution of data are not necessary. CC-NUMA and COMA are two of these more recent designs that use some additional hardware in each node. You can physically picture each node having one or more processors, each with its own cache attached via the memory subsystem to the local node memory (this is just the SMP configuration we have previously examined). Now imagine a separate remote-access device (RAD) connected to the memory subsystem and to the node interconnect network. The RAD connections for an AS/400 would be very similar to the I/O subsystems shown in Figure 10.1; the 6xx buses would be connected to one side and the SAN ports would be connected to the other side of the RAD.1


1I am calling this additional piece of hardware a RAD, but there is no agreed-upon name. Some systems call it a coherence controller; others call it a hub. Whatever its name, its function is to implement the directory-based cache-coherence protocol between the node memory subsystem and the interconnection network to the remote nodes.

In a CC-NUMA machine, the RAD contains a separate cache that contains only remote data. When a processor accesses data that is not in its own cache, the requested data is fetched from the node memory if the address is local or from the RAD’s cache if the address is remote. References to remote data not satisfied by the RAD’s cache must be sent across the interconnect network to the referenced memory page’s home node to obtain the requested block of data from the remote node’s memory and to enforce any necessary coherence actions.

The RAD’s cache improves the performance of a CC-NUMA machine by reducing the number of remote cache misses that must go across the interconnection network to a remote node. Obviously, the first reference to a remote memory page will incur the long wait time while the data is retrieved from the remote node’s memory and placed into the RAD’s cache. Subsequent references to the same remote memory page by any processor in the node will retrieve the data quickly, without the need to go across the interconnection network. As a result, the ratio of remote-to-local cache miss times is reduced to somewhere between 2:1 and 3:1 for current systems. This penalty for remote accesses is low enough that most applications do not have to be changed when you move from an SMP to a cluster implementation. Because of this, CC-NUMA systems are often referred to as scalable-SMP designs.

An example of a CC-NUMA system is the SGI/Cray Origin 2000. The Origin 2000 consists of up to 64 nodes interconnected by a scalable CrayLink network. Each node contains one or two processors, up to 4 GB of memory, and connections to the I/O subsystem. The maximum configuration today has 128 processors and 256 GB of total memory. The first processors to be used for the Origin 2000 are MIPS R10000s running at 195 MHz, each with 4 MB L2 caches.

The two processors in a single node do not function as an SMP configuration because there is no snoopy protocol between their L2 caches. Instead, they operate as two separate processors that share hardware connections to the node’s memory and I/O. The processors in a node are connected to a hub chip that in turn connects to the node’s memory, the I/O subsystem, and the CrayLink interconnection network. The hub chip passes local accesses directly to the node’s memory. A separate memory in the hub caches remote data. If the request for remote data is not in the hub memory, the remote access is sent over the interconnection network to the remote node. It is also interesting to note that the hub chip uses a crossbar switch to quickly route the information through the hub.

A Cray Origin 2000 system (configurations larger than 64 processors are designated as Cray systems; smaller models are attributed to SGI) is also the basis for the ASCI Blue Mountain project I introduced in Chapter 2. The Origin system is being deployed in stages at the Los Alamos National Laboratory, with the goal of reaching a 3072 processor configuration that is capable of 4 teraflops late in 1998. You may remember from our earlier discussion that there are two halves of the ASCI Blue Project. IBM’s ASCI Blue Pacific system will be deployed at the Lawrence Livermore National Laboratory during the same timeframe to achieve the same level of performance with 512 8-way SMP nodes. The success of these two systems should provide a very good indication of how far distributed shared-memory machines can take us.

Recently, there have been several studies and research activities to improve on the performance of a CC-NUMA machine and to further reduce the ratio of remote-to-local cache miss times. One configuration that shows promise is COMA. A COMA system uses the same directory-based cache-coherence protocol that CC-NUMA uses, but the COMA system allocates part of the local node’s main memory to act as a large cache for remote data. The separate remote data cache in the RAD is eliminated with COMA; instead, the remote data is allowed to reside in the node’s processor cache hierarchy and main memory.

The original COMA work, started in the early 1990s, allowed data in cache-block sizes to be brought into the main memory of the node (similar to the way the separate remote cache of CC-NUMA stores cache blocks). The problem with this design is that cache blocks are smaller than pages in memory, so additional hardware is required in the node to manage a second page size in the main memory. This hardware is essentially a duplicate of the virtual memory hardware that I described in Chapter 8, and which is already implemented in the node. More recent implementations of COMA, called simple-COMA or S-COMA, force remote data to always be stored in page-sized blocks in the node’s main memory. With this approach, both the local and remote data accesses can use the existing virtual-memory hardware. Of course, within an SMP node, hardware still must be included to implement the directory-based cache-coherence protocol for the remote data instead of the snoopy protocol used for local data.

S-COMA can potentially outperform CC-NUMA because it can exploit the node’s large memory to hold remote data. S-COMA can dynamically tailor the fraction of the memory used for remote data in response to an application’s need. On the other hand, S-COMA requires larger blocks of data to be transferred across the interconnection network whenever there is a remote miss in a node. In the next couple of years, we will see whether S-COMA or some variant is a better choice than the CC-NUMA design various computer system vendors are currently using.


Previous Table of Contents Next

Copyright © NEWS/400 Books