> Partial import: ProQuest exposes only a 24-page preview without institutional access. This Seed document contains the available front matter, abstract, contents, acknowledgments, and the visible opening excerpt of Chapter 1. It does not represent the complete dissertation.
David Hillel Gelernter
Doctor of Philosophy in Computer Science
State University of New York at Stony Brook
May 1983
ProQuest publication number: 8313989
Copyright: 1982, David Hillel Gelernter
Abstract
A network computer is a computer network designed to function not as a collection of autonomous hosts but as one machine. Systems intended to support experiments with asynchronous distributed programs—including distributed systems programs and distributed applications—form a major subclass of network computers. The goal of the Stony Brook Network (SBN) project is to construct such a general-purpose network computer.
Unlike most other systems in its class, SBN is a language-centered design. Its starting point is the distributed programming language Linda. To run on SBN, a distributed program must either be written in or preprocessed into Linda; the language is intended to be a powerful, flexible, and expressive vehicle for distributed programming. The role of SBN's hardware and communication software is to support Linda efficiently—to implement a Linda machine.
Part I discusses the Linda design, including its interprocess communication primitives and program-structuring devices. Communication takes place through a logically shared data structure called structured memory, which may be implemented over many memory-disjoint nodes. Part II discusses communication-system algorithms supporting Linda, including runtime rendezvous, store-and-forward deadlock prevention, and staged circuit switching.
Contents
Chapter 1: Introduction
Part I: Linda
Chapter 2: Description of Linda
Chapter 3: Programming in Linda
Part II: The SBN Structured Memory Machine
Chapter 4: The Global Name Space
Chapter 5: A DAG-Based Algorithm for Prevention of Store-and-Forward Packet Deadlock
Chapter 6: Staged Circuit Switching
Chapter 7: Conclusions
Acknowledgments
Thanks to my advisor Professor Arthur Bernstein, and to Professors John Cherniavsky and Hussein Badr as well. Particular thanks are due my student co-workers on this project, Ken Bertapelle, Bob Ensor, Todd Morgan, Suresh Jagganathan and above all Mauricio Arango. Their contributions made the work described here possible.
Chapter 1: Introduction
1.1 Network Computers
A network computer is one computer built out of many, a computer network designed to function not as a collection of autonomous hosts but as one machine. As the term is used here, support for distributed programs—for programs, that is, that execute on many processors simultaneously—distinguishes network computers from conventional multiprocessors on the one hand and computer networks on the other. The term refers to an idea, not an architecture. Network computers will ordinarily consist of a collection of identical nodes, each with local memory, joined by a communication network. But in certain cases the term will be applied to a collection of processor nodes that share access via a switching network to a collection of memory modules, and even to local area networks that support distributed programs. We discuss each of these
cases in the following sections.
Deferring for the moment any discussion of particular systems, two questions arise immediately: why build network computers? How ought they to be built? At present the answers to both questions are preliminary at best (ill-defined at worst), Network computer
building is an infant field.
A general answer to the first question is easily provided: one computer built out of many is potentially more powerful and more reliable than any one of its constituent computers. Given a program that may be divided into N simultaneously-executing components, a single computer built out of N identical constituent computers might be, in the theoretical best case, N times faster than any of its constituents. As a practical matter, it is unlikely that built systems will approach linear speedup in executing most distributed programs. But even if the average realized speedup in an N-node system is as modest as log₂N, network computers that are routinely an order of magnitude faster than their constituents in the execution of distributable programs seem a reasonable goal. Whether such average speedups can in fact be obtained
is, however, an open question.
The reliability of a network computer might in the abstract increase exponentially with its size. An N-node machine is completely disabled only if all N of its constituents are simultaneously disabled, In fact one expects practical considerations, such as the reliability and the distribution of peripherals, once again to hold realized gain to drastically below the theoretical maximum, and of course a network computer that survives the loss of many or most of its constituents will suffer a corresponding degradation in performance. Nonetheless, it seems reasonable to predict that a
well-designed network computer will be substantially more reliable
than its constituents.
A further potential of network computers is frequently overlooked. While a network computer composed of N identical nodes is unlikely to be N times faster than any of its constituents, it will certainly be N times larger—that is, it will incorporate N times as much RAM as each constituent does singly. Disregard, temporarily, the conventional conception of an N-node network computer as N computers hooked together. Consider the N-node machine instead as a large RAM divided into N segments, each segment with a processor attached. Consider a program P that occupies all N memory segments simultaneously. Each instruction in P is executed by the processor attached to the segment in which it is loaded. Cross-segment memory references obviously will involve large overheads, because they must be channelled through the network communication system. On the other hand, locality of reference might keep the number of cross-segment references relatively low, and the swapping overhead necessarily incurred were all N segments multiplexed on one node is avoided. The fact that P may execute in N
places simultaneously is a side-effect.
The intent of this Gedankenexperiment is not to urge the adoption of an untried approach to network computer architecture, rather to emphasize the network computer's potential contributions to the management of very large programs. Network computers might in fact allow the development of programs that are simply too large
for any conventional machine to handle.
Finally, two other potential advantages of network computers are frequently mentioned. Network computers provide modular expandability to the extent that they may be expanded incrementally by the addition of small, relatively cheap microprocessor nodes. Some researchers anticipate that network computers will be cheaper to
build than conventional machines of comparable power.
Network computers, then, offer potential gains in speed, reliability and capacity over conventional systems. To what end are these potential gains directed? In what kind of problem domains are network computers expected to be effective? Because speed has dominated network-computer development so far, problem domains of interest have been those in which distributed algorithms promise to speed computation. Interest has generally been divided between synchronous numerical algorithms and asynchronous algorithms, which may or may not be numerical. Systems that address a subset of the second of these two domains are of central interest in this thesis; we are concerned specifically with network computers that provide a general-purpose environment for experiments in asynchronous distri—- buted computing. We approach a general discussion of this domain
by way of brief descriptions of the domains that are not of con-
cern.
The ILliac-IV is the progenitor of a class of network computers intended to execute parallel, synchronous numerical algorithms.
Its designers speak of handling "wellL-formulated but
computationally massive problems" such as manipulations of very large matrices and the solution of sets of partial differential equations[1]. Machines of this sort have prompted susbtantial theoretical work in two related sreas—the design of parallel algorithms, and the design of network architectures, particularly of interconnection topologies, to support them efficiently. Thus Stone[2], for example, presents his perfect shuffle topology in context of synchronous parallel algorithms such as parallel fast Fourier Transform, polynomical evaluation and sorting that it is
intended to support.
Systems designed for the execution of asynchronous distributed algorithms may be divided into three general classes. (1) Those designed around a model of computation different from the conventional control-flow model. (2) Those optimized to a particular subclass of distributed algorithm within the conventional controlflow model. (3) Those designed for general experiments in distri-
buted programming within the control—flow model.
Treleaven et. al.[3] discuss two major computational models that differ from control-flow and are associated with networkcomputer projects—the dataflow and the demand-driven model. In dataflow programs, a statement becomes executable as soon as the data values it depends upon become available. As many statements as fulfill this requirement are executable simultaneously. A demand-
driven program is in essence an expression that calls for the
evaluation of its constituent sub-expressions as they are needed. Any expression's constituents may be evaluated in parallel. Major dataflow projects are described by Dennis[5], Gostelow and Tho-
mas[6] and others; demand-driven projects are described by
Keller[7] and others.
A number of network computer projects are devoted not to specialized models of computation but to specialized problem domains. One example is the Stony Brook Hierarchical Multicomputer[8,9], designed to support problem-solving by decomposition, in particular to support concurrent execution of parallel subtasks in examples such as the 8-queens and travelling salesman problems that Lend themselves to divide-and-conquor solutions. Another special-purpose system is FEM[10] designed to solve structural analysis problems; graph-structured problems are mapped direct onto the densely-
conected FEM processor array, with Logical graph nodes mapped one-
to-one onto processor nodes.
We arrive, then, at the class of systems that is of interest in this thesis—the class intended to support general experimentation in distributed programming. An initial (and highly complex) problem within this broad domain is the construction of generalpurpose distributed operating systems. These are intended to controt network machines and to make their power, particularly their
support for distributed algorithms, availeble to users.
Of course the construction of distributed operating systems, as formidable and interesting a problem as it is in its own right, is not in the long term a self-sufficient justification for work in network computers. Distributed operating system projects beg the question of what users are supposed to do with network-computer Power once it is made available to them. An interesting collection of possible answers is discussed by ‘Deminet11] in his performance analysis of the Cm network computer. He discusses programs thet perform distributed Guicksort, solve sets of partial differential equations, compute fust Fourier transforms and execute a distributed railway-network simulator. Jones and Schwans[12] discuss another application for Cm, a distributed filtering algorithm that processes image data. Construction of distributed and pipelined compilers is e problem of obvious interest (El-Dessouki et. at.{13]). The builders of the Arachne network computer discuss an application that is ultimately of the greatest importance in context of the work to be discussed in this thesis: distributed Artificial Intelligence programming. Finkel and Solomon[14] discuss a
game—playing program implemented on Arachne.
At the outset, two fundamental network-computer questions were raised. We have concentrated so far on the first, that is on describing the broad impetus for work in this field. We move now to the second: how ought network computers to be constructed? We
are concerned now exclusively with systems that fall within our
class of interest. The remainder of this section discusses proposed answers to the second question by way of a survey of systems in this class. The two following sections introduce the research
project that is the topic of this thesis: the SBN system.
Systems within our class of interest fall into two categories: microprocessor networks, and local area networks that support distributed programs. We discuss six systems in the first category— Cm*, Arachne, Micronet, Munet, X-tree and the NYU ULtracomputer—
and two in the second, Ethernet and Demos.
Cm[15] is a 50-node network computer that uses a hierarchical bus to implement a global system-wide address space. Nodes are organized into clusters; five clusters make up Cm. Each network node has local memory attached, but any node Ni can reference any other node Nj's local memory via the hierachical bus. The time required to perform non-local memory references depends upon whether Ni and Nj fall in the same or different clusters, and is
always greater than the time required for a local memory reference.
Cm is unique within this category in the size and sophistication of its implementation and the availebility of data on its performance. Two separate operating sysiems have been implemented on Cn (Star0S[16] and Medusa[17]); Cm*'s performance is discussed by Jones and Schwartz(15] and by Deminet[11]. Deminet's results indi-
cate that speedup factors produced by parallel processing on Cm*
depend strongly on the program being executed and are Limited in some cases by the hardware design. Thus, the performance of distributed Quicksort is constrained by heavy contention at the memory module in which the distributed program's shared data is stored; performance of a program to solve PDE's depends strongly on the way in which data and processes are mapped to network nodes (although observed speedup is close to the theoretical maximum when the mapping is performed carefully); in some versions of distributed FFT, the number of nodes that can usefully execute in parallel is Limited by saturation of the hierarchical bus. Handling of Cm is evidently made more difficult by what Jones and Schwartz themselves refer to as the "management problem"[15,p.134] that is associated with systems that contain nonhomogenous or asymmetric components. On Cm, communication costs between two nodes vary as a function of their mutual position in the network, and the variation occurs not gradually as on a Linked network but in one sharp Leap as communication moves from intra~ to inter-cluster. In some cases the hierarchical nature of the Cm hardware has also, as noted, proved troublesome. (It is of course essential to note that any negative criticism of Cm must be qualified by the fact that no comparable Performance data is available for any other system in this
category.)
Arechne[14] is a network computer consisting of five nodes
joined by word-psrallel point-to-point Links in an arbitrary topol—
ogy. (Arachne's operating system, also called Arachne, makes no assumptions about network topology; it requires only that the network graph be connected.) Control in Arachne is distributed over the network; resource allocation is handled by negotiation between Per-node resource~manager processes, and each node maintains only local state information. Conclusions drawn by the builders based on their experience with Arachne center lergely around the suitebility of the inter-process communication mechanism incorporated in the system. Tasking and IPC in Arachne are discussed in Chapter
Three.
Micronet as described by Wittie and Tillborg[18] consists of 16 nodes that communicate via shared contention bus. Each node has two bus ports and each bus may connect 17 nodes; the system may be configured arbitrarily within these constraints. One possible topology (for example) is a 4x4 grid of processors with four row busses, four column busses, and each node connected to one row and one column bus. Micronet's operating system ("Micros") consists of a kernel replicated on every node and a set of system processes supported by the kernel. The system is intended to support paral-
lel user programs eventually.
Micronet's hardware is sophisticated and interesting. Consider a hypothetical 7x7 Micronet grid, a network with about as many nodes as Cm* has. Communication overhead between an arbitrary
pair of nodes varies in a fashion similar to the variation in Cm*:
messages between two arbitrary nodes travel one hop if they are connected along the same bus, otherwise two hops. But unlike Cm, Micronet has no hierarchical ordering among busses and thus no intrinsic topological bottlenecks. Micronet is furthermore more densely-connected than Cm: in the 7x7 grid, 12 other nodes share either one or the other of a given node's busses, as opposed to 9 other nodes that fall within a given node's cluster (and thus share its bus) in the Cm configuration described above. {In the absence of performance studies on Micronet comparable to the cited studies of Cm, direct quantitative comparison between the two architec—
tures is nonetheless impossible.)
Munet is described by Halstead[19] as a 10-node reconfigurable network. Nodes communicate over point-to-point Links; each node may have from one to four neighbors. Computation on Munet is based upon an idiosyncratic form of message passing among program modules: a message is referred to as an "event" insofar as it initiates execution of the target program module on the arguments it encapsulates. Dynamic task distribution for Load-levelling is a central goal of the Munet design. The system supports transparent mobility of processes through the net by providing a logical-addressing tech~ nique called "reference trees". All processes or data modules {both instances of "objects" on Munet] that may reference a name P must be configured in a Logical tree whose edges are mapped to physical
network Links; the tree may grow or shrink dynamically. Each node
> The ProQuest preview ends here, mid-discussion of Munet. Consult the full dissertation through ProQuest or a subscribing institution for the remaining chapters.
Do you like what you are reading? Subscribe to receive updates.
Unsubscribe anytime