MetaMP : a higher level abstraction for message-passing programming

Steve W. Otto · 2018

The potential performance of distributed-memory parallel computers is very high, but their programming has proven to be difficult. The only successful approach so far has been to program them directly in the message-passing system of the machine. To a extent this forms the of the computer. Higher level programming abstractions are available, such as versions of parallel Fortran, but it has proven difficult to compile these to efficient distributed-memory code. Here, we propose a slightly more modest approach, whereby useful abstractions are supported by a compiler and run-time system (MetaMP), but these constructs are within a message-passing framework. The user still writes a messagepassing program, but the MetaMP compiler understands distributed data structures and is therefore able to help in powerful ways. A preliminary version of MetaMP has been written which supports simple multi-dimensional arrays. Extensions to more complex data structures (e.g., unstructured meshes, dynamically changing arrays) are planned. MetaMP programs have proven to be succinct and more understandable than their counterparts. The performance of MetaMP programs is close to that obtainable by manual programming. Currently, MetaMP compiles down to Express, a commercial message-passing system developed at Caltech and available on many parallel computers. Objectives and Relation to Other Work Parallel computers such as the Intel Touchstone, the Ncube 11, and the Meiko Computing Surface form a class of MIMD machines which can be termed message-passing. The processors inside these systems are, to a first approximation, conventional microprocessors with a large amount of memory (.5 to 16 Mbytes in 1990) and an interface to a hardware message passing system. Though the programming of these machines remains problematical, they have been successfully used in many specific cases. The potential performance is very high, since this architecture can be easily scaled to numbers of processors. To a great extent, these machines have been manually programmed, using the message-passing calls provided by the system directly. A fundamental property of message-passing machines is that message passing times are one to three orders of magnitude slower than fundamental floating point operation times. This necessitates a style of programming in which communications are carefully scheduled so that: the correctness of the program is preserved; the communications occur infrequently; and many data items are transferred per message. Message-passing programming has often been compared to assembly language programming. Intricate details of distributed data structures must be managed by the programmer. The question naturally arises: Can message-passing programming be abstracted to a more understandable form without losing much of the performance of custom programming? Many parallel languages and compilers have been proposed and implemented on MIMD computers [I-161. These systems often allow the programmer an extremely clean and simple model of the parallel computation. Typically, all elements of an array are accessible by any processor (shared memory), and synchronization is provided automatically by the compiler (e.g., the programmer just writes doall). Unfortunately, it appears to be difficult to compile from a parallel language such as this to a message-passing computer, with the restriction that the resultant code be efficient. The research proposed here concerns a set of language extensions and a compiler called MetaMP. In contrast to the systems mentioned previously, MetaMP has a somewhat less ambitious goal. MetaMP does not attempt to completely hide the message passing nature of the underlying hardware. This makes the compiler implementable while preserving the efficiency of the resultant code. The programmer is still given a message-passing view of the hardware, but it is an abstract, minimalist one. The user still writes a message passing program, but the MetaMP compiler understands distributed data structures and is therefore able to help in powerful ways. Programs written in this language have proven to be more compact and understandable than those written directly in the underlying message passing system. Scientific computing focuses on programs which construct and manipulate large, multi-dimensional arrays. Our first version of MetaMP provides support for these types of programs. The MetaMP compiler introduces auxiliary data structures which describe the shapes, sizes, offsets, etc., of distributed multi-dimensional arrays. Different arrays can be distributed (or decomposed) in different ways and MetaMP keeps track of each decomposition. Abstract loop constructs similar to doall are available and release the programmer from having to remember the decomposition details of each array. Communications are more easily expressed since the compiler understands the shapes of arrays and spread and reduction operations can be succinctly written. Array sections of different from one processor to the next are completely supported by MetaMP. This means that the problem of odd sizes (array dimensions not exactly conforming to the machine size) can be removed. As we will demonstrate in our examples below, programs which handle any size problem on any size machine can be written, yet they are still succinct. Locality often plays a role in scientific computing. In solving a set of partial differential equations for example, arrays (or meshes) representing spatial locations are distributed across the parallel computer. The locality of the differential operator reflects itself in the fact that the required communications are of the nearest neighborn type. Such algorithms require array elements from a narrow boundary strip (or face in three dimensions) in the array sections of neighboring processors. MetaMP provides full support for this. Guard strips, that is, extra array elements which map to neighboring array sections, can be specified within MetaMP. It turns out that there is an elegant way to do this which makes the extra, guard elements transparent to the programmer. This will be discussed later in the context of a two dimensional elliptical PDE solver. The MetaMP language consists of two components: normal, sequential C (or Fortran) containing for loops that run over the multi-dimensional arrays, MetaMP directives which modify the meaning of the sequential f o r loops to their parallel, distributed-memory counterparts. The directives always appear between delimiters, that is, they look like this: % directive Y,. Compile time checking is done to ensure that the directives to distribute for loops make sense. Loop indices have associated ranges and these are compared with the allowable ranges of the distributed arrays. This checking catches most simple types of programming error, such as mixing up array subscripts or combining distributed arrays in an incompatible way. A dependency analysis can also be done to check if the semantics of the loops has been altered by the parallel directives. This is planned, but is not done in this first version of MetaMP. Currently, if the user inserts a directive to distribute a for loop, MetaMP does it, even if the meaning of the program is altered. A first version of MetaMP has been developed and non-trivial programs have been written in the language. The current version compiles down to a commercially available parallel message passing system, Express. The programs can be executed on actual parallel hardware. There is good reason to believe that efficiencies close to that obtained by manual programming are being achieved, though these have not yet been measured. The syntax and semantics of the MetaMP directives seem to be clean; the programs succinctly state what is happening in the parallel machine. We will discuss our plan of development for MetaMP in a later section. First we will describe a bit more thoroughly what MetaMP is through the use of a few examples.

Read the paper · More papers on PaperTik