Parallel processing systems and method
Claim Score by NHIP
Abstract
Methods and systems for parallel computation of an algorithm using a plurality of nodes configured as a Howard Cascade. A home node of a Howard Cascade receives a request from a host system to compute an algorithm identified in the request. The request is distributed to processing nodes of the Howard Cascade in a time sequence order in a manner to minimize the time to so expand the Howard Cascade. The participating nodes then perform the designated portion of the algorithm in parallel. Partial results from each node are agglomerated upstream to higher nodes of the structure and then returned to the host system. The nodes each include a library of stored algorithms accompanied by data template information defining partitioning of the data used in the algorithm among the number of participating nodes.

Term
Term ended
Projected expiry passed 23 May 2026, 0.3 years ago.
- Priority
- Filed
- Published
- Projected expiry
- Today
104 claims: 15 independent, 89 dependent
- 1A parallel processing system, comprising:a plurality of processing nodes arranged in a Howard Cascade;a first home node, responsive to an algorithm processing request, for (a) broadcasting the algorithm processing request to the plurality of processing nodes in a time and sequence order within the Howard Cascade and for (b) broadcasting a dataset of an algorithm to at least top level processing nodes of the Howard Cascade simultaneously;contiguous processing nodes within the plurality of processing nodes being operable to process contiguous parts of the dataset and to agglomerate results contiguously to the first home node in reverse to the time and sequence order.
- 39A method for processing a context-based algorithm for enhanced parallel processing within a parallel processing architecture, comprising the steps of:A. determining whether work to be performed by the algorithm is intrinsic to the algorithm or to another algorithm;B. determining whether the algorithm requires data movement;and C. parallelizing the algorithm based upon whether the work is intrinsic to the algorithm and whether the algorithm requires data movement.
- 44A method for parallel processing an algorithm for use within a Howard Cascade, comprising the steps of:extract input and output data descriptions for the algorithm;acquire data for the algorithm;process the algorithm on nodes of the Howard Cascade;agglomerate node results through the Howard Cascade;and return results to a remote host requesting parallel processing of the algorithm.
- 45A method for parallel computation comprising:transmitting an algorithm computation request and associated data from a requesting host to a home node of a computing system wherein the request includes a requested number (N) of processing nodes to be applied to computation of the request;distributing the computation request from the home node to a plurality of processing nodes wherein the plurality of processing nodes includes N processing nodes coupled to the home node and wherein the distribution is in a hierarchical ordering;broadcasting the associated data from the home node to all of the plurality of processing nodes;agglomerating a final computation result from partial computation results received from the plurality of processing nodes wherein the agglomeration is performed in the reverse order of the hierarchical ordering;and returning the final computation result from the home node to the requesting host.
- 52A method of distributing an algorithm computation request comprising:receiving within a home node of a distributed parallel computing system a computation request and associated data from a requesting host system;determining a number of processing nodes (N) of the parallel computing system to be applied to performing the computation request;partitioning the associated data to identify a portion of the data associated with each processing node;and recursively communicating the computing request and information regarding the partitioned data from the home node to each of the N processing nodes over a plurality of communication channels during a sequence of discrete time intervals, wherein each communication channel is used for communication by at most one processing node or home node during any one time interval, and wherein the number of discrete time intervals to recursively communicate to all N processing nodes is minimized.
- 59A method for distributing an algorithm computation request for a complex algorithm in a parallel processing system comprising:receiving from a requesting host a computation request for a complex algorithm wherein the complex algorithm includes a plurality of computation sections;expanding the computation request to a plurality of nodes configured as a Howard Cascade;computing within the Howard Cascade a first computation section to generate a partial result;returning the partial result to a control device;receiving further direction from the control device;computing within the Howard Cascade a next computation section to generate a partial result in response to receipt of further direction to compute the next computation section;repeating the steps of returning, receiving and computing the next computation section in response to receipt of further direction to compute the next computation section;and returning the partial result to the requesting host as a final result in response to further direction to complete processing of the complex algorithm.
- 62Broadest claimClaim Score 85, broad(NHIP)A method for parallelizing an algorithm comprising:receiving a new algorithm description;automatically annotating the new algorithm description with template information relating to data used by the new algorithm and relating to data generated by the new algorithm;and storing the annotated new algorithm in each processing node of a Howard Cascade parallel processing system.
- 65A computer readable storage medium tangibly embodying program instructions for a method for parallel computation, the method comprising:transmitting an algorithm computation request and associated data from a requesting host to a home node of a computing system wherein the request includes a requested number (N) of processing nodes to be applied to computation of the request;distributing the computation request from the home node to a plurality of processing nodes wherein the plurality of processing nodes includes N processing nodes coupled to the home node and wherein the distribution is in a hierarchical ordering;broadcasting the associated data from the home node to all of the plurality of processing nodes;agglomerating a final computation result from partial computation results received from the plurality of processing nodes wherein the agglomeration is performed in the reverse order of the hierarchical ordering;and returning the final computation result from the home node to the requesting host.
- 72A computer readable storage medium tangibly embodying program instructions for a method of distributing an algorithm computation request, the method comprising:receiving within a home node of a distributed parallel computing system a computation request and associated data from a requesting host system;determining a number of processing nodes (N) of the parallel computing system to be applied to performing the computation request;partitioning the associated data to identify a portion of the data associated with each processing node;and recursively communicating the computing request and information regarding the partitioned data from the home node to each of the N processing nodes over a plurality of communication channels during a sequence of discrete time intervals, wherein each communication channel is used for communication by at most one processing node or home node during any one time interval, and wherein the number of discrete time intervals to recursively communicate to all N processing nodes is minimized.
- 79A computer readable storage medium tangibly embodying program instructions for a method for distributing an algorithm computation request for a complex algorithm in a parallel processing system, the method comprising:receiving from a requesting host a computation request for a complex algorithm wherein the complex algorithm includes a plurality of computation sections;expanding the computation request to a plurality of nodes configured as a Howard Cascade;computing within the Howard Cascade a first computation section to generate a partial result;returning the partial result to a control device;receiving further direction from the control device;computing within the Howard Cascade a next computation section to generate a partial result in response to receipt of further direction to compute the next computation section;repeating the method steps of returning, receiving and computing the next computation section in response to receipt of further direction to compute the next computation section;and returning the partial result to the requesting host as a final result in response to further direction to complete processing of the complex algorithm.
- 82A computer readable storage medium tangibly embodying program instructions for a method for parallelizing an algorithm, the method comprising:receiving a new algorithm description;automatically annotating the new algorithm description with template information relating to data used by the new algorithm and relating to data generated by the new algorithm;and storing the annotated new algorithm in each processing node of a Howard Cascade parallel processing system.
- 85A system for parallel computation comprising:means for transmitting an algorithm computation request and associated data from a requesting host to a home node of a computing system wherein the request includes a requested number (N) of processing nodes to be applied to computation of the request;means for distributing the computation request from the home node to a plurality of processing nodes wherein the plurality of processing nodes includes N processing nodes coupled to the home node and wherein the distribution is in a hierarchical ordering;means for broadcasting the associated data from the home node to all of the plurality of processing nodes;means for agglomerating a final computation result from partial computation results received from the plurality of processing nodes wherein the agglomeration is performed in the reverse order of the hierarchical ordering;and means for returning the final computation result from the home node to the requesting host.
- 92A system of distributing an algorithm computation request comprising:means for receiving within a home node of a distributed parallel computing system a computation request and associated data from a requesting host system;means for determining a number of processing nodes (N) of the parallel computing system to be applied to performing the computation request;means for partitioning the associated data to identify a portion of the data associated with each processing node;and means for recursively communicating the computing request and information regarding the partitioned data from the home node to each of the N processing nodes over a plurality of communication channels during a sequence of discrete time intervals, wherein each communication channel is used for communication by at most one processing node or home node during any one time interval, and wherein the number of discrete time intervals to recursively communicate to all N processing nodes is minimized.
- 99A system for distributing an algorithm computation request for a complex algorithm in a parallel processing system comprising:means for receiving from a requesting host a computation request for a complex algorithm wherein the complex algorithm includes a plurality of computation sections;means for expanding the computation request to a plurality of nodes configured as a Howard Cascade;means for computing within the Howard Cascade a first computation section to generate a partial result;means for returning the partial result to a control device;means for receiving further direction from the control device;means for computing within the Howard Cascade a next computation section to generate a partial result in response to receipt of further direction to compute the next computation section;means for repeating the steps of returning, receiving and computing the next computation section in response to receipt of further direction to compute the next computation section;and means for returning the partial result to the requesting host as a final result in response to further direction to complete processing of the complex algorithm.
- 102A system for parallelizing an algorithm comprising:means for receiving a new algorithm description;means for automatically annotating the new algorithm description with template information relating to data used by the new algorithm and relating to data generated by the new algorithm;and means for storing the annotated new algorithm in each processing node of a Howard Cascade parallel processing system.
Independent claims15
464 paragraphs in 6 sections, as filed
RELATED APPLICATIONS
[0001] This application is a continuation-in-part of commonly-owned and co-pending U.S. patent application Ser. No. 09/603,020, filed on Jun. 26, 2000, entitled MASSIVELY PARALLEL INTERNET COMPUTING, and incorporated herein by reference. This application also claims priority to U.S. Patent Application No. 60/347,325, filed on Jan. 10, 2002, entitled PARALLEL PROCESSING SYSTEMS AND METHODS, and incorporated herein by reference.
BACKGROUND
[0002] Prior art programming methods are implemented with parallel processing architectures called “clusters.” Such clusters are generally one of two types: the shared memory cluster and the distributed memory cluster. A shared memory cluster consists of multiple computers connected via RAM memory through a back plane. Because of scaling issues, the number of processors that may be linked together via a shared back plane is limited. A distributed memory cluster utilizes a network rather than a back plane. Among other problems, one limitation of a distributed memory cluster is the bandwidth of the network switch array.
[0003] More particularly, each node of an N-node distributed memory cluster must obtain part of an initial dataset before starting computation. The conventional method for distributing data in such a cluster is to have each node obtain its own data from a central source. For problems that involve a large dataset, this can represent a significant fraction of the time it takes to solve the problem. Although the approach is simple, it has several deficiencies. First, the central data source is a bottleneck: only one node can access the data source at any given time, while others nodes must wait. Second, for large clusters, the number of collisions that occur when the nodes attempt to access the central data source leads to a significant inefficiency. Third, N separate messages are required to distribute the dataset over the cluster. The overhead imposed by N separate messages represents an inefficiency that grows directly with the size of a cluster; this is a distinct disadvantage for large clusters.
[0004] Shared memory clusters of the prior art operate to transfer information from one node to another as if the memory is shared. Because the data transfer cost of a shared memory model is very low, the data transfer technique is also used within clustered, non-shared memory machines. Unfortunately, using a shared memory model in non-shared memory architectures imposes a very low efficiency; the cluster inefficiency is approximately three to seven percent of the actual processor power of the cluster.
[0005] Although increasing the performance of the central data source can reduce the impact of these deficiencies, adding additional protocol layers on the communication channel to coordinate access to the data, or to increase the performance of the communication channel, adds cost and complexity. These costs scale directly as the number of nodes increase, which is another significant disadvantage for large clusters in the prior art.
[0006] Finally, certain high performance clusters of the prior art also utilize invasive “parallelization” methods. In such methods, a second party is privy to the algorithms used on a cluster. Such methods are, however, commercially unacceptable, as the users of such clusters desire confidentiality of the underlying algorithms.
[0007] The prior art is familiar with four primary parallel programming methods: nested-mixed model parallelism, POSIX Pthreads, compiler extensions and work sharing. Nested-mixed model parallelism is where one task spawns subtasks. This has the effect of assigning more processors to assist with tasks that could benefit with parallel programming. It is however difficult to predict, a priori, how job processing will occur as the amount of increase in computational speed remains unknown until after all of subtasks are created. Further, because only parts of a particular job benefit from the parallel processing, and because of the high computational cost of task spawning, the total parallel activity at the algorithm level is decreased. According to the so-called Amdahl's Law of the prior art, even a small percentage change in parallel activity generates large effective computational cost.
[0008] POSIX Pthreads are used in shared memory architectures. Each processor in the shared memory is treated as a separate processing thread that may or may not pass thread-safe messages in communicating with other threads. Although this may work well in a shared memory environment, it does not work well in a distributed processor environment. The inability to scale to large numbers of processors even in a shared memory environment has been well documented in the prior art. Because of bus speed limits, memory contention, cache issues, etc., most shared memory architectures are limited to fewer than sixty-four processors working on a single problem. Accordingly, efficient scaling beyond this number is problematic. The standard method of handling this problem is to have multiple, non-communicating algorithms operating simultaneously. This still limits the processing speedup achievable by a single algorithm.
[0009] Compiler extensions, such as distributed pointers, sequence points, and explicit synchronization, are tools that assist in efficiently programming hardware features. Accordingly, compiler extensions tools offer little to enhance parallel processing effects as compared to other prior art methods.
[0010] Work-Sharing models are characterized by how they divide the work of an application as a function of user-supplied compiler directives or library function calls. The most popular instance of work sharing is loop unrolling. Another work-sharing tool is the parallel region compiler directives in OpenMP, which again provides for limited parallel activity at the algorithm level.
[0011] It is interesting to note that prior art parallel processing techniques are injected into the algorithms and are not “intrinsic” to the algorithms. This may be a result of the historic separation between programming and algorithm development, in the prior art.
SUMMARY
[0012] The following U.S. patents provide useful background to the teachings hereinbelow and are incorporated herein by reference: <tables id="TABLE-US-00001" num="1"><table frame="none" colsep="0" rowsep="0"><tgroup align="left" colsep="0" rowsep="0" cols="3"><colspec colname="OFFSET" colwidth="28PT" align="left" /><colspec colname="1" colwidth="105PT" align="left" /><colspec colname="2" colwidth="84PT" align="left" /><thead><row><entry /><entry /></row><row><entry /><entry namest="OFFSET" nameend="2" align="center" rowsep="1" /></row></thead><tbody valign="top"><row><entry /><entry>U.S. Pat. No. 4,276,643</entry><entry>June 1981</entry></row><row><entry /><entry>U.S. Pat. No. 4,667,287</entry><entry>May 1987</entry></row><row><entry /><entry>U.S. Pat. No. 4,719,621</entry><entry>January 1988</entry></row><row><entry /><entry>U.S. Pat. No. 4,958,273</entry><entry>September 1990</entry></row><row><entry /><entry>U.S. Pat. No. 5,023,780</entry><entry>June 1991</entry></row><row><entry /><entry>U.S. Pat. No. 5,079,765</entry><entry>January 1992</entry></row><row><entry /><entry>U.S. Pat. No. 5,088,032</entry><entry>February 1992</entry></row><row><entry /><entry>U.S. Pat. No. 5,093,920</entry><entry>March 1992</entry></row><row><entry /><entry>U.S. Pat. No. 5,109,515</entry><entry>April 1992</entry></row><row><entry /><entry>U.S. Pat. No. 5,125,081</entry><entry>June 1992</entry></row><row><entry /><entry>U.S. Pat. No. 5,166,931</entry><entry>November 1992</entry></row><row><entry /><entry>U.S. Pat. No. 5,185,860</entry><entry>February 1993</entry></row><row><entry /><entry>U.S. Pat. No. 5,224,205</entry><entry>June 1993</entry></row><row><entry /><entry>U.S. Pat. No. 5,371,852</entry><entry>December 1994</entry></row><row><entry /><entry namest="OFFSET" nameend="2" align="center" rowsep="1" /></row></tbody></tgroup></table></tables>
[0013] The present inventions solves numerous problems in parallel computing by providing methods and systems for configuration of and operation of a Howard Cascade (HC). Algorithm computation requests are transmitted to a home node of a HC and distributed (cascaded) through nodes of the HC to expand the algorithm computation request to a desired number of processing nodes. The algorithm to be computed may be predefined by a user and is stored within each processing node of the HC. Associated data is then broadcast to all processing nodes participating in the requested computation. Partial results are agglomerated upstream in the opposite order to the expansion of the HC and then, eventually, forwarded to the requesting host.
[0014] A first aspect of the invention provides a parallel processing system having a plurality of processing nodes arranged in a Howard Cascade (HC). The HC couples to a first home node, responsive to an algorithm processing request, for (a) broadcasting the algorithm processing request to the plurality of processing nodes in a time and sequence order within the Howard Cascade and for (b) broadcasting a dataset of an algorithm to at least top level processing nodes of the Howard Cascade simultaneously. The HC then provides for contiguous processing nodes within the plurality of processing nodes being operable to process contiguous parts of the dataset and to agglomerate results contiguously to the first home node in reverse to the time and sequence order.
[0015] Another aspect of the invention provides systems and a method for processing a context-based algorithm for enhanced parallel processing within a parallel processing architecture. The method first determining whether work to be performed by the algorithm is intrinsic to the algorithm or to another algorithm. The method then determining whether the algorithm requires data movement; and parallelizing the algorithm based upon whether the work is intrinsic to the algorithm and whether the algorithm requires data movement.
[0016] Another aspect of the invention provides systems and a method for parallel processing an algorithm for use within a Howard Cascade. The method being operable to extract input and output data descriptions for the algorithm; acquire data for the algorithm; process the algorithm on nodes of the Howard Cascade; agglomerate node results through the Howard Cascade; and return results to a remote host requesting parallel processing of the algorithm.
[0017] Still another aspect of the invention provides systems and a method for parallel computation. The method first transmitting an algorithm computation request and associated data from a requesting host to a home node of a computing system wherein the request includes a requested number (N) of processing nodes to be applied to computation of the request. The method then distributes the computation request from the home node to a plurality of processing nodes wherein the plurality of processing nodes includes N processing nodes coupled to the home node and wherein the distribution is in a hierarchical ordering. Then the method broadcasts the associated data from the home node to all of the plurality of processing nodes, agglomerates a final computation result from partial computation results received from the plurality of processing nodes wherein the agglomeration is performed in the reverse order of the hierarchical ordering, and returns the final computation result from the home node to the requesting host.
[0018] Another aspect of the invention provides systems and a method of distributing an algorithm computation request. The method comprises the step of receiving within a home node of a distributed parallel computing system a computation request and associated data from a requesting host system. The method also includes the step of determining a number of processing nodes (N) of the parallel computing system to be applied to performing the computation request and the step of partitioning the associated data to identify a portion of the data associated with each processing node. Lastly the method comprises the step of recursively communicating the computing request and information regarding the partitioned data from the home node to each of the N processing nodes over a plurality of communication channels during a sequence of discrete time intervals, wherein each communication channel is used for communication by at most one processing node or home node during any one time interval, and wherein the number of discrete time intervals to recursively communicate to all N processing nodes is minimized.
[0019] Yet another aspect of the invention provides systems and a method for distributing an algorithm computation request for a complex algorithm in, a parallel processing system. The method comprising the steps of: receiving from a requesting host a computation request for a complex algorithm wherein the complex algorithm includes a plurality of computation sections; expanding the computation request to a plurality of nodes configured as a Howard Cascade; computing within the Howard Cascade a first computation section to generate a partial result; returning the partial result to a control device; receiving further direction from the control device; computing within the Howard Cascade a next computation section to generate a partial result in response to receipt of further direction to compute the next computation section; repeating the steps of returning, receiving and computing the next computation section in response to receipt of further direction to compute the next computation section; and returning the partial result to the requesting host as a final result in response to further direction to complete processing of the complex algorithm.
[0020] Still another aspect of the invention provides systems and a method for parallelizing an algorithm. The method first receive a new algorithm description from a host system. Next the new algorithm is automatically annotated with template information relating to data used by the new algorithm and relating to data generated by the new algorithm. Lastly the annotated new algorithm is stored in each processing node of a Howard Cascade parallel processing system.
BRIEF DESCRIPTION OF THE DRAWINGS
[0021]FIG. 1 is a block diagram of a prior art parallel processing cluster that processes an application running on a remote host.
[0022]FIG. 2 is a block diagram illustrating further detail of the prior art cluster, application and remote host of FIG. 1.
[0023]FIG. 3 schematically illustrates algorithm partitioning and processing through the prior art remote host and cluster of FIG. 1, FIG. 2.
[0024]FIG. 4 is a block schematic illustrating further complexity of one parallel data item within an algorithm processed through the prior art parallel application of FIG. 1, FIG. 2.
[0025]FIG. 5 is a block diagram of one Howard Cascade Architecture System connected to a remote host.
[0026]FIG. 6 is a block schematic of one Howard Cascade.
[0027]FIG. 7 is a block schematic illustrating broadcast messaging among the Howard Cascade of FIG. 6.
[0028]FIG. 8 illustrates agglomeration through the Howard Cascade of FIG. 6.
[0029]FIG. 9 is a graph illustrating processing times for nodes in an unbalanced algorithm for calculating Pi to 1000 digits for six different cluster sizes.
[0030]FIG. 10 is a graph illustrating processing times for nodes in a balanced algorithm for calculating Pi to 1000 digits for six different cluster sizes.
[0031]FIG. 11 is a block schematic illustrating expansion of a problem dataset in two time units among nodes of a Howard Cascade, each node including one processor and one communication channel.
[0032]FIG. 12 is a block schematic illustrating expansion of a problem dataset in two time units among nodes of a Howard Cascade, each node including two processors and two communication channels.
[0033]FIG. 14 illustrates one Howard Cascade segregated into three levels among three strips.
[0034]FIG. 15 illustrates the Howard Cascade of FIG. 14 during agglomeration.
[0035]FIG. 16 illustrates further detail of the Howard Cascade of FIG. 14, including switches and routers.
[0036]FIG. 17 illustrates reconfiguring of cascade strips within the Howard Cascade of FIG. 14, to accommodate boundary conditions.
[0037]FIG. 18 shows one Howard Cascade with multiple home nodes.
[0038]FIG. 19 shows one home node network communication configuration for a home node within the Howard Cascade of FIG. 18.
[0039]FIG. 20 shows one processing node network communication configuration for a processing node within the Howard Cascade of FIG. 18.
[0040]FIG. 21 shows one Howard Cascade with seven processing nodes and one unallocated processing node.
[0041]FIG. 22 shows the Howard Cascade of FIG. 21 reallocated upon failure of one of the seven processing nodes.
[0042]FIG. 23 shows further detail of the Howard Cascade of FIG. 21.
[0043]FIG. 24 shows one Howard Cascade during agglomeration.
[0044]FIG. 25 shows the Howard Cascade of FIG. 24 reallocated to accommodate a failed processing node during agglomeration.
[0045]FIG. 26 shows one Howard Cascade during distribution of an algorithm processing request.
[0046]FIG. 27 shows the Howard Cascade of FIG. 26 reallocated to accommodate a failed processing node during distribution of the algorithm processing request.
[0047]FIG. 28 shows one Howard Cascade configured to recast an algorithm processing request to acquire additional processing nodes at lower cascade levels.
[0048]FIG. 29 shows the Howard Cascade of FIG. 28 recasting the algorithm processing request.
[0049]FIG. 30 shows one Howard Cascade with a spare home node.
[0050]FIG. 31 shows the Howard Cascade of FIG. 30 reconfigured after a failed home node.
[0051]FIG. 32 illustrates two single-processor nodes of a Howard Cascade utilizing a single communication channel.
[0052]FIG. 33 illustrates two two-processor nodes of a Howard Cascade utilizing a double communication channel.
[0053]FIG. 34 shows one interface between proprietary algorithms and one Howard Cascade Architecture System.
[0054]FIG. 35 illustrates further definition associated with the interface of FIG. 34.
[0055]FIG. 36 shows one Howard Cascade utilizing the interfaces of FIG. 34 and FIG. 35.
[0056]FIG. 37-FIG. 62 illustrate processes for implementing algorithms for one Howard Cascade Architecture System.
[0057]FIG. 63 illustrates one example of a complex algorithm.
[0058]FIG. 64 one method of implementing complex algorithms in an HCAS.FIG. 66 illustrates a 2D dataset for an ECADM category algorithm;
[0059]FIG. 66 illustrates a 2D dataset for an ECADM category algorithm;
[0060]FIG. 67 illustrates row data distribution of the 2D dataset of FIG. 66 among nodes in a cluster.
[0061]FIG. 68 illustrates column data distribution of the 2D dataset of FIG. 66 among nodes in a cluster.
[0062]FIG. 69 is a block diagram illustrating one Howard Cascade processing a 2D dataset as in FIG. 66.
[0063]FIG. 70 illustrates one example of two processing nodes each utilizing two communication channels to connect to two network switches.
[0064]FIG. 71 illustrates one example of application threads running on a processing node utilizing a multiple channel software API for managing multiple communication channels.
[0065]FIG. 72 illustrates one advantage of using two communication channels per node in an HCAS.
DETAILED DESCRIPTION OF ILLUSTRATED EMBODIMENTS
[0066]FIG. 1 shows a block diagram illustrating a prior art parallel processing cluster <b>10</b> that processes a parallel application <b>28</b> operating on a remote host <b>26</b>. As those skilled in the art appreciate, cluster <b>10</b> is for example a Beowulf cluster that connects several processors <b>12</b> together, through a communication channel <b>14</b> and a switch <b>16</b>, to process parallel application <b>28</b>. The number N of processors <b>12</b> is carefully matched to cluster type, e.g., sixteen processors for a Beowulf cluster. Remote host <b>26</b> communicates with cluster <b>10</b> via a data path <b>18</b> to access the collective computing power of processors <b>12</b> within cluster <b>10</b> to run parallel application <b>28</b>, with a goal of reducing the execution time of parallel application <b>28</b>.
[0067]FIG. 2 shows parallel processing application <b>28</b> using cluster <b>10</b> in more detail. Nodes <b>15</b>(<b>1</b>), <b>15</b>(<b>2</b>), <b>15</b>(<b>3</b>) . . . <b>15</b>(N) represent processors <b>12</b>, FIG. 1, connected together via communication channel <b>14</b> that facilitates the movement of computer programs, data, and inter-process messages through nodes <b>15</b>. As those skilled in the art appreciate, communication channel <b>14</b> may consist of computer buses, network linkages, fiber optic channels, Ethernet connections, switch meshes (e.g., switch <b>16</b>, FIG. 1) and/or token rings, for example. Cluster <b>10</b> further has a gateway <b>24</b> that provides an interface to cluster <b>10</b>. Gateway <b>24</b> typically provides one or more interfaces—such as a network interface <b>30</b>, an Internet interface <b>34</b>, and a peer-to-peer interface <b>35</b>—to communicate with remote host <b>26</b> through a cluster boundary <b>36</b>.
[0068] Remote host <b>26</b> is for example a computer system that runs parallel application <b>28</b>. Illustratively, FIG. 2 shows application <b>28</b> and remote host <b>26</b> connected to gateway node <b>24</b> through network interface <b>30</b>. In operation, gateway <b>24</b> enables communication between parallel application <b>28</b>, running on host <b>26</b>, and individual processing nodes <b>15</b> within cluster <b>10</b>. Application <b>28</b> loads algorithm code and data onto individual processing nodes <b>15</b> prior to parallel processing of the algorithm code and data.
[0069] There are two main prior art methods for application <b>28</b> to load algorithm code and data onto nodes <b>15</b>. If cluster <b>10</b> is an encapsulated cluster, only the network address of gateway <b>24</b> is known by application <b>28</b>; application <b>28</b> communicates with gateway <b>24</b> which in turn routes the algorithm code and data to nodes <b>15</b>. Although internal nodes <b>15</b> are hidden from host <b>26</b> in this approach, application <b>28</b> must still have knowledge of nodes <b>15</b> in order to specify the number of nodes to utilize and the handling of node communications through gateway <b>24</b>.
[0070] In a non-encapsulated cluster <b>10</b>, application <b>28</b> has information concerning the organization and structure of cluster <b>10</b>, as well as knowledge to directly access individual nodes <b>15</b>. In one example, application <b>28</b> uses a remote procedure call (RPC) to address each node <b>15</b> individually, to load algorithm and data onto specific nodes <b>15</b>; in turn, nodes <b>15</b> individually respond back to application <b>28</b>.
[0071] Unfortunately, neither of these prior art techniques make cluster <b>10</b> behave like a single machine to application <b>28</b>. Application <b>28</b> must therefore be designed a priori for the particular internal architecture and topology of cluster <b>10</b> so that it can appropriately process algorithm code and data in parallel.
[0072]FIG. 3 schematically illustrates how algorithm code and data are partitioned and processed through remote host <b>26</b> and cluster <b>10</b>. Such algorithm code and data are logically illustrated in FIG. 3 as algorithm design <b>40</b>. An initialization section <b>42</b> first prepares input data prior to algorithm execution; this typically involves preparing a data structure and loading the data structure with the input data. A process results section <b>44</b> handles data produced during algorithm execution; this typically involves storing and outputting information. Designated parallel section <b>56</b> contains parallel-specific data items <b>48</b>, <b>50</b>, and <b>52</b>. Specifically, parallel data items <b>48</b>, <b>50</b> and <b>52</b> represent parts of algorithm code that are to be executed on different nodes <b>15</b> in cluster <b>10</b>. FIG. 4 illustrates further detail of one parallel data item <b>48</b>; parallel data items <b>50</b> and <b>52</b> of FIG. 3 are similarly constructed. Process results section <b>44</b> combines data output from execution of parallel data items <b>48</b>, <b>50</b> and <b>52</b> to form a complete data set.
[0073] Each parallel data item <b>48</b>, <b>50</b> and <b>52</b> is transferred as algorithmic code and data to cluster <b>10</b> for processing. Synchronization with algorithm execution is attained by inter-process communication between parallel application <b>28</b> and individual nodes <b>15</b> of cluster <b>10</b>. Parallel processing libraries facilitate this inter-process communication, but still require that algorithm design <b>40</b> be formatted, a priori, for the topology of nodes <b>15</b>. That is, the topology of cluster <b>10</b> must be known and fixed when application <b>28</b> is compiled.
[0074] Arrow <b>58</b> represents the step of compiling algorithm design <b>40</b> into parallel application <b>28</b>, for execution on remote host <b>26</b>. Elements <b>42</b>, <b>44</b>, <b>46</b>, <b>48</b>, <b>50</b>, <b>52</b> and <b>56</b> are illustratively shown as compiled elements <b>42</b>′, <b>44</b>′, <b>46</b>′, <b>48</b>′, <b>50</b>′, <b>52</b>′, and <b>56</b>′ in parallel application <b>28</b>.
[0075] Parallel application <b>28</b> is executed by remote host <b>26</b> in the following sequence: initialization section <b>42</b>′, parallel data initialization <b>46</b>′, compiled parallel data item <b>48</b>′, compiled parallel data item <b>50</b>′, compiled parallel data item <b>52</b>′, and process results section <b>44</b>′. In the illustrated example of FIG. 3, compiled parallel data item <b>48</b>′ transfers algorithm code and input data to node <b>15</b>(<b>2</b>) in cluster <b>10</b> via algorithm transfer section <b>70</b> and data input section <b>72</b>, respectively, of FIG. 4. Compiled parallel data item <b>48</b>′ also controls the execution of its algorithm on node <b>15</b>(<b>2</b>) using process synchronization <b>74</b> (FIG. 4) to determine when data processing has completed. When node <b>15</b>(<b>2</b>) indicates that the execution of parallel data item <b>48</b>′ has completed, data output transfer section <b>76</b> (FIG. 4) transfers results from node <b>15</b>(<b>2</b>) back into parallel application <b>28</b>.
[0076] In similar fashion, compiled parallel data item <b>50</b>′ and compiled parallel data item <b>52</b>′ utilize nodes <b>15</b>(<b>4</b>) and <b>15</b>(<b>5</b>), respectively, such that all algorithms of parallel application <b>28</b> are processed concurrently by nodes <b>15</b>(<b>2</b>), <b>15</b>(<b>4</b>), <b>15</b>(<b>5</b>) of cluster <b>10</b>.
[0077] Process results section <b>44</b>′ of parallel application <b>28</b> operates when compiled parallel data items <b>48</b>′, <b>50</b>′, and <b>52</b>′, are processed by cluster <b>10</b> and returned by respective data output transfer sections <b>76</b>. Data returned by nodes <b>15</b>(<b>2</b>), <b>15</b>(<b>4</b>), <b>15</b>(<b>5</b>) of cluster <b>10</b>, in this example, are combined by process results section <b>44</b>′ so as to continue processing of the completed data set.
[0078] It should be apparent from the foregoing description of the prior art that the focus of architectural designs in high performance computing has concentrated on various hardware-centric solutions. This has been true since the beginnings of parallel processing, particularly since the 1970s. Moreover, the current trend for new parallel processing systems is to use multiple processors working in tandem to achieve higher performance at lower cost. Historically, therefore, hardware advances such as improved processor performance, faster networks, improvements in memory and caching technologies, etc, have far outpaced software advances.
[0079] This has remained true for the prior art parallel processing designs despite the fact that Amdahl's law (which predicts parallel performance for a given application on multiple processors) does not expressly address hardware parameters. Amdahl's law is represented by the following equation: speedup =1/((1−f)+f/p), where f is the percent of parallel activity within the algorithm, and p equals the number of processors. “Speedup” determines the processor speed multiplier, e.g., 1×, 2×, etc., of the current processor speed of the individual linked processors in the parallel processing design.
[0080] Amdahl's law only takes into consideration the degree of parallel activity at the algorithm level and the number of processors used in the calculations. Finding the limit of Amdahl's law (with respect to the number of processors) is a standard operation that yields the disheartening understanding of how little serial activity must be present before the parallel processing effect becomes unusable. That is, the “maximum speedup” under Amdahl's law is given by the following relationship: lim<sub>p→∞</sub>[1/((1−f)+f/p)]=1/(1−f) where “maximum speedup” equals the processor speed multiplier, e.g., 1×, 2×, etc., of the current processor speed of the individual linked processors if there are an infinite number of processors p, and f is the percent of parallel activity within the algorithm. Table 1 below shows the maximum speedup for given values of f. <tables id="TABLE-US-00002" num="2"><table frame="none" colsep="0" rowsep="0"><tgroup align="left" colsep="0" rowsep="0" cols="1"><colspec colname="1" colwidth="217PT" align="center" /><thead><row><entry namest="1" nameend="1" align="center">TABLE 1</entry></row></thead><tbody valign="top"><row><entry /></row><row><entry namest="1" nameend="1" align="center" rowsep="1" /></row><row><entry>Maximum Speedup under Amdahl's Law</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="4"><colspec colname="OFFSET" colwidth="21PT" align="left" /><colspec colname="1" colwidth="28PT" align="center" /><colspec colname="2" colwidth="77PT" align="center" /><colspec colname="3" colwidth="91PT" align="left" /><tbody valign="top"><row><entry /><entry /><entry>Maximum</entry><entry /></row><row><entry /><entry>f</entry><entry>Speedup</entry></row><row><entry /><entry namest="OFFSET" nameend="3" align="center" rowsep="1" /></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="4"><colspec colname="OFFSET" colwidth="21PT" align="left" /><colspec colname="1" colwidth="28PT" align="center" /><colspec colname="2" colwidth="77PT" align="char" char="." /><colspec colname="3" colwidth="91PT" align="left" /><tbody valign="top"><row><entry /><entry>0.10000</entry><entry>1.11</entry><entry>Processor Equivalent</entry></row><row><entry /><entry>0.20000</entry><entry>1.25</entry><entry>Processor Equivalent</entry></row><row><entry /><entry>0.30000</entry><entry>1.43</entry><entry>Processor Equivalent</entry></row><row><entry /><entry>0.40000</entry><entry>1.67</entry><entry>Processor Equivalent</entry></row><row><entry /><entry>0.50000</entry><entry>2.00</entry><entry>Processor Equivalent</entry></row><row><entry /><entry>0.60000</entry><entry>2.50</entry><entry>Processor Equivalent</entry></row><row><entry /><entry>0.70000</entry><entry>3.33</entry><entry>Processor Equivalent</entry></row><row><entry /><entry>0.80000</entry><entry>5.00</entry><entry>Processor Equivalent</entry></row><row><entry /><entry>0.85000</entry><entry>6.67</entry><entry>Processor Equivalent</entry></row><row><entry /><entry>0.90000</entry><entry>10.00</entry><entry>Processor Equivalent</entry></row><row><entry /><entry>0.95000</entry><entry>20.00</entry><entry>Processor Equivalent</entry></row><row><entry /><entry>0.99000</entry><entry>100.00</entry><entry>Processor Equivalent</entry></row><row><entry /><entry>0.99900</entry><entry>1000.00</entry><entry>Processor Equivalent</entry></row><row><entry /><entry>0.99990</entry><entry>10000.00</entry><entry>Processor Equivalent</entry></row><row><entry /><entry>0.99999</entry><entry>100000.00</entry><entry>Processor Equivalent</entry></row><row><entry /><entry namest="OFFSET" nameend="3" align="center" rowsep="1" /></row></tbody></tgroup></table></tables>
[0081] As described in more detail below, since the key parameter in Amdahl's law is “f”, the problem of generating high performance computing with multiple processors is overcome, in part, by (a) approaching parallel processing from an algorithm-centric perspective, and/or by utilizing the Howard Cascade (HC) which increases the parallel activity of cross-communication, each described in more detail below.
[0082] The Howard Cascade Architecture System
[0083] The Howard Cascade Architecture System (“HCAS”) provides a topology for enabling a cluster of nodes (each containing one or more processors) to act as a single machine to one or more remote host computers outside of the inter-cluster network. Unlike the prior art, this enables each remote host to communicate and parallel process algorithm code and data within the cluster but without (a) direct communication with individual cluster nodes and/or (b) detailed knowledge of cluster topology. More particularly, in one embodiment, the HCAS uses the following features to create a single machine experience for a remote host.
[0084] Complete mathematical and logical algorithms are stored on each node of the cluster prior to being accessed by the remote host.
[0085] Mathematical and logical algorithms are made parallel by changing the data sets processed by the mathematical and logical algorithms, though no additions, changes or deletions occur to the algorithm itself.
[0086] The remote host sends algorithm processing requests and data only to a gateway processor attached to a network; mathematical and logical algorithms are not communicated by the remote host to the home node.
[0087] The remote host knows only the IP address and a single port identity of the gateway of the HCAS; multiple ports and/or multiple IP addresses are not required to make use of multiple processing nodes in the cluster.
[0088] The remote host requires no knowledge of the internal parameters of the cluster, including the number of processors used by the cluster and the connectivity model used by the cluster.
[0089] The gateway communicates with a home node within the cascade; the home node facilitates communication to the connected processing nodes.
[0090]FIG. 5 illustrates one HCAS <b>80</b> that provides algorithm-centric parallel processing. HCAS <b>80</b> has a gateway <b>107</b>, a home node <b>110</b> and, illustratively, three processing nodes <b>88</b>, <b>90</b> and <b>92</b>. Those skilled in the art should appreciate that three nodes <b>88</b>, <b>90</b>, <b>92</b> are shown for purposes of illustration, and that many more nodes may be included within HCAS <b>80</b>. Processing nodes <b>88</b>, <b>90</b>, <b>92</b> (and any other nodes of HCAS <b>80</b>) are formulated into a HC, described in more detail below. Gateway <b>107</b> communicates with a remote host <b>82</b> and with home node <b>110</b>; home node <b>110</b> facilitates communication to and among processing nodes <b>88</b>, <b>90</b>, <b>92</b>.
[0091] Each processing node <b>88</b>, <b>90</b> and <b>92</b> has an algorithm library <b>99</b> that contains computationally intensive algorithms; algorithm library <b>99</b> preferably does not contain graphic user interfaces, application software, and/or computationally non-intensive functions. A remote host <b>82</b> is shown with a remote application <b>84</b> that has been constructed using computationally intensive algorithm library API <b>86</b>. Computationally intensive algorithm library API <b>86</b> defines an interface for computationally intense functions in the algorithm library <b>99</b> of processing nodes <b>88</b>, <b>90</b> and <b>92</b>.
[0092] In operation, remote host <b>82</b> sends an algorithm processing request <b>83</b>, generated by computationally intensive algorithm library API <b>86</b>, to gateway <b>107</b>. Gateway <b>107</b> communicates request <b>83</b> to controller <b>108</b> of home node <b>110</b>, via data path <b>109</b>. Since the computationally intensive algorithms of libraries <b>99</b> are unchanged, and remain identical across processing nodes, “parallelization” within HCAS <b>80</b> occurs as a function of how an algorithm traverses its data set. Each of the algorithms, when placed on processing nodes <b>88</b>, <b>90</b> and <b>92</b>, is integrated with a data template <b>102</b>. Controller <b>108</b> adds additional information to algorithm processing request <b>83</b> and distributes the request and the additional information to processing nodes <b>88</b>, <b>90</b>, <b>92</b> via data paths <b>104</b>, <b>106</b>, <b>108</b>, respectively; the additional information details (a) the number of processing nodes (e.g., N=3 in this example) and (b) data distribution information. Each processing node <b>88</b>, <b>90</b>, <b>94</b> has identical control software <b>94</b> that routes algorithm processing request <b>83</b> to data template software <b>102</b>. Data template software <b>102</b> computes data indexes and input parameters to communicate with a particular algorithm identified by algorithm processing request <b>83</b> in algorithm library <b>99</b>.
[0093] Data template software <b>102</b> determines whether or not a particular computationally intensive algorithm requires data. If the algorithm requires data, data template <b>102</b> requests such data from home node <b>110</b>. The algorithm in library <b>99</b> is then invoked with the appropriate parameters, including where to find the data, how much data there is, and where to place results. There is no need for remote host <b>82</b> to have information concerning HCAS <b>80</b> since only the data set is being manipulated. Specifically, remote host <b>82</b> does not directly send information, data, or programs to any processing node <b>88</b>, <b>90</b>, <b>92</b>. HCAS <b>80</b> appears as a single machine to remote host <b>82</b>, via gateway <b>107</b>. Once HCAS <b>80</b> completes its processing, results from each node <b>88</b>, <b>90</b>, <b>92</b> are agglomerated (described in more detail below) and communicated to remote host <b>82</b> as results <b>85</b>.
[0094] In one embodiment, the HCAS maximizes the number of nodes that can communicate in a given number of time units. The HCAS avoids the inefficiencies (e.g., collisions in shared memory environments, the bottle-neck of a central data source, and the requirement of N messages for an N node cluster) in the prior art by, for example, broadcasting the full data set to all processing nodes at once. Even though the same amount of data is transferred over the communication channel, the broadcasting eliminates the overhead of using N separate messages. An important advantage of the broadcasting is that the overhead of sending data is independent of the number of nodes in the HCAS. This is especially important in maintaining efficiency of a large cluster.
[0095] Each HCAS has at least one HC, such as HC <b>100</b> of FIG. 6. HC <b>100</b> is illustratively shown within seven processing nodes <b>112</b>, <b>114</b>, <b>116</b>, <b>118</b>, <b>120</b>, <b>122</b> and <b>124</b>. Algorithm processing requests <b>140</b>, <b>142</b>, <b>144</b>, <b>146</b>, <b>148</b>, <b>150</b> and <b>152</b> are messages passed between nodes of HC <b>100</b>. As shown, three time units (time <b>1</b>, time <b>2</b>, time <b>3</b>) are used to expand HC <b>100</b> to all seven nodes. A time unit is, for example, the transmission time for one message in the HCAS, regardless of the transmission medium. For example, the message may be transmitted across any of the following media; LAN, WAN, shared bus, fiber-optic interface etc.
[0096] HC <b>100</b> transmits the algorithm processing requests to an arbitrary group of nodes. Full expansion occurs when the problem set of all algorithm processing requests <b>140</b>, <b>142</b>, <b>144</b>, <b>146</b>, <b>148</b>, <b>150</b> and <b>152</b> has been transmitted to all required processing nodes (nodes <b>112</b>, <b>114</b>, <b>116</b>, <b>118</b>, <b>120</b>, <b>122</b> and <b>124</b> in this example). Home node <b>110</b> does not necessarily participate in parallel processing.
[0097] Once HC <b>100</b> is fully expanded, all nodes have received the algorithm processing request and are ready to accept associated data. Each processing node knows how much data to expect based upon information it received in its algorithm processing request message. Each processing node then prepares to receive the data by joining a multicast message group, managed by home node <b>110</b>, and by listening on a broadcast communication channel. Through the multicast message group, home node <b>110</b> broadcasts a single message that is received by all nodes that are listening on the broadcast communication channel. Home node <b>110</b> waits for all nodes to join the multicast message group and then broadcasts the dataset to all nodes simultaneously. An example of the data broadcast on a broadcast communication channel <b>121</b> of HC <b>100</b> is shown in FIG. 7.
[0098] The data decomposition scheme used by HC <b>100</b> partitions the dataset based on node number. For example, processing node <b>112</b> receives the first piece of data, processing node <b>114</b> receives the second piece of data, and so on, until the last needed processing node of HC <b>100</b> receives the last piece of data. In this manner, the dataset is distributed such that the data for each processing node is adjacent to the data received by downstream processing nodes.
[0099] As each node of HC <b>100</b> receives data, it can choose to store the entire dataset or only the data it needs. If only the needed data is stored, the remaining data may be discarded when received. This inter-node choice reduces the amount of memory used to store data while maintaining flexibility that facilitates improved parallel-processing performance, as described below.
[0100] Once a processing node receives its data, it leaves the multicast message group. Home node <b>110</b> monitors which nodes have left the multicast message group, thereby providing positive acknowledgment that each node has received its data. Home node <b>110</b> can thereby detect failed processing nodes by noting whether a processing node remains in the multicast message group after a specified timeout; it may then open a discrete communication channel with the processing node to attempt recovery.
[0101] When the processing nodes have produced their results, an agglomeration process commences. Agglomeration refers to (a) the gathering of individual results from each of the processing nodes and (b) the formatting of these results into the complete solution. Each processing node sends its results to the processing node that is directly upstream. The flow of results thereby occurs in reverse sequence order of the initial expansion within HC <b>100</b>. An example of an agglomeration process by HC <b>100</b> is shown in FIG. 8.
[0102] A direct result of agglomeration is that the results from each node maintain the same ordered relationship as the decomposition of the initial dataset. Each processing node knows how many downstream processing nodes it has; and the subsequent downstream results, from the downstream nodes, form a contiguous block of data. Each of the processing nodes has its results data, and the location and size information that enables the upstream processing node to properly position the results, when received. As the results are sent upstream through HC <b>100</b>, the size of the result information expands contiguously until the entire result block is assembled at home node <b>110</b>.
[0103]FIG. 8, for example, illustrates how agglomeration occurs through the nodes of HC <b>100</b>. At a first time unit (“Time 1”), nodes <b>116</b>, <b>118</b>, <b>122</b>, <b>124</b> communicate respective results (identified as “processor 3 data”, “processor 4 data”, “processor 6 data”, and “processor 7 data”, respectively) to upstream nodes (i.e., to node <b>114</b>, node <b>112</b>, node <b>120</b>, and home node <b>110</b> (for node <b>124</b>)). Data agglomerates at these upstream nodes such that, at a second time unit (“Time 2”), nodes <b>114</b> and <b>120</b> communicate respective results (including the results of its downstream nodes) to upstream nodes <b>112</b> and <b>110</b>, respectively. At the third time unit (“Time 3”), agglomeration completes as the results from each downstream node has agglomerated at home node <b>110</b>, as shown.
[0104] The agglomeration process may therefore be considered the “reverse flow” of data through HC <b>100</b>. One important benefit of the reverse flow is that the communication network that connects the nodes is not saturated, as would be the case if all processing nodes attempt to simultaneously present results to home node <b>110</b>. Many cluster architectures of the prior art experience network bottlenecks when results are returned to a single node; this bottleneck can be exacerbated when the computational load is evenly distributed among the processing nodes, which is desired to obtain maximum computational efficiency.
[0105] Context-Based Algorithms for Parallel Processing Architectures
[0106] Certain context-based algorithms may be thought of as the work-defining elements of an algorithm. For example, a kernel is the work-defining element that is performed on a dataset involving a convolution algorithm. As described in more detail below, such context-based algorithms may be further distinguished with (a) internal or external contexts, and/or with (b) algorithmic data movement or no algorithmic data movement. Algorithmic data movement is data that moves as a natural process within the algorithm. An example of an algorithm with an internal requirement to move data is the matrix inversion algorithm, since data movement is intrinsic to the matrix inversion algorithm.
[0107] Accordingly, certain context-based algorithms distinguished by algorithm data movements defines five parallelization categories: Transactional, Internal Context No Algorithmic Data Movement (“ICNADM”), Internal Context with Algorithmic Data Movement (“ICADM”), External Context No Algorithmic Data Movement (“ECNADM”), and External Context with Algorithmic Data Movement (“ECADM”). These categories and their importance are discussed in more detail below, and are relevant to all types of parallel processing clusters, including HC <b>100</b>.
[0108] Transactional Category
[0109] Both Beowulf clusters and SETI@Home parallel processing architectures, known in the art, work best when the problems being processed are transactional in nature. Transactional problems are those in which the problem space consists of a large number of independent problems. Each problem is thus solved separately and there is no cross-communication or synchronization between the problems. Almost any grouping of computers may be used for this class of problem; the only challenge is moving the data onto and off of the group of computers.
[0110] By way of example, the following applications have, by and large, transactional problem sets: analyzing star spectra, certain banking applications, and genetic engineering. Such applications share several common traits: there is no cross-communication, there is a high degree of parallel processing, and there is high scalability.
[0111] HC <b>100</b> may serve to process problem sets in the transactional category.
[0112] Internal Context No Algorithmic Data Movement (ICNADM) Category
[0113] Many logical and mathematical problems fall into the ICNADM category. By way of definition, the ICNADM category is one in which the final work to be accomplished, for the particular logical/mathematical algorithm, is intrinsic to the algorithm and the data being transformed does not have to be moved as part of the algorithm. The following examples have ICNADM problem sets: series expansions, Sobel edge detection, convolutions, matrix multiplications, correlations and cross-correlations, one-dimensional fast-Fourier transforms (1D-FFTs), 1D-wavelets, etc.
[0114] With regard to parallelization, numerical computations involving series expansions have inherent imbalances in the computational complexity of low order terms at the beginning of, and high order terms at the end of, a series. When algorithms based on a series expansion are converted to a parallel implementation for execution by parallel processing, this imbalance presents an obstacle to achieving parallel efficiency. More particularly, when the terms of the series expansion are distributed among nodes of a parallel processing architecture in consecutive intervals, (e.g. 1-10, 11-20, 21-30, . . . ), the nodes that are assigned the first intervals have less computation than the nodes that are assigned the last intervals. Disparity in the computational load leads to inefficiency in parallel computations.
[0115] In accord with one embodiment hereof, the computational load is distributed equally across all nodes in order to achieve maximum efficiency in parallel computations. This distribution enhances parallel processing within HC <b>100</b>, or within existing parallel processing architectures of the prior art.
[0116] Given the diversity of series expansions, it is nonetheless difficult to predict the increase in computational complexity for each term. In accord with one embodiment hereof, every n<sup>th </sup>term is assigned to each node to equalize the computational loads. By making n equal to the number of nodes in the parallel processing architecture, each node's computational load has an equal number of low-and high-order terms. In cases where the total number of terms in a series is evenly divisible by the number of nodes, then each node will have an equal number of terms. An example of a forty-two term series expansion for a seven node array (e.g., HC <b>100</b>, FIG. 6) is shown in Table 2. <tables id="TABLE-US-00003" num="3"><table frame="none" colsep="0" rowsep="0" pgwide="1"><tgroup align="left" colsep="0" rowsep="0" cols="1"><colspec colname="1" colwidth="259PT" align="center" /><thead><row><entry namest="1" nameend="1" align="center">TABLE 2</entry></row></thead><tbody valign="top"><row><entry /></row><row><entry namest="1" nameend="1" align="center" rowsep="1" /></row><row><entry>Example of 42-term Series Expansion in 7-Node Architecture</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="2"><colspec colname="1" colwidth="35PT" align="center" /><colspec colname="2" colwidth="224PT" align="center" /><tbody valign="top"><row><entry>Node #</entry><entry>Series terms</entry></row><row><entry namest="1" nameend="2" align="center" rowsep="1" /></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="7"><colspec colname="1" colwidth="35PT" align="char" char="." /><colspec colname="2" colwidth="35PT" align="left" /><colspec colname="3" colwidth="42PT" align="left" /><colspec colname="4" colwidth="35PT" align="left" /><colspec colname="5" colwidth="42PT" align="left" /><colspec colname="6" colwidth="35PT" align="left" /><colspec colname="7" colwidth="35PT" align="left" /><tbody valign="top"><row><entry>1</entry><entry>1</entry><entry>8</entry><entry>15</entry><entry>22</entry><entry>29</entry><entry>36</entry></row><row><entry>2</entry><entry> 2</entry><entry> 9</entry><entry> 16</entry><entry> 23</entry><entry> 30</entry><entry> 37</entry></row><row><entry>3</entry><entry> 3</entry><entry> 10</entry><entry> 17</entry><entry> 24</entry><entry> 31</entry><entry> 38</entry></row><row><entry>4</entry><entry> 4</entry><entry> 11</entry><entry> 18</entry><entry> 25</entry><entry> 32</entry><entry> 39</entry></row><row><entry>5</entry><entry> 5</entry><entry> 12</entry><entry> 19</entry><entry> 26</entry><entry> 33</entry><entry> 40</entry></row><row><entry>6</entry><entry> 6</entry><entry> 13</entry><entry> 20</entry><entry> 27</entry><entry> 34</entry><entry> 41</entry></row><row><entry>7</entry><entry> 7</entry><entry> 14</entry><entry> 21</entry><entry> 28</entry><entry> 35</entry><entry> 42</entry></row><row><entry namest="1" nameend="7" align="center" rowsep="1" /></row></tbody></tgroup></table></tables>
[0117] In cases where the total number of terms in a series is not equally divisible by the number of nodes, then, in accord with one embodiment, the series terms are divided as equally as possible. An example of a thirty-nine term series expansion with a seven node array (e.g., HC <b>100</b>, FIG. 6) is shown in Table 3. <tables id="TABLE-US-00004" num="4"><table frame="none" colsep="0" rowsep="0" pgwide="1"><tgroup align="left" colsep="0" rowsep="0" cols="1"><colspec colname="1" colwidth="259PT" align="center" /><thead><row><entry namest="1" nameend="1" align="center">TABLE 3</entry></row></thead><tbody valign="top"><row><entry /></row><row><entry namest="1" nameend="1" align="center" rowsep="1" /></row><row><entry>Example of 39-term Series Expansion in 7-Node Architecture</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="2"><colspec colname="1" colwidth="35PT" align="center" /><colspec colname="2" colwidth="224PT" align="center" /><tbody valign="top"><row><entry>Node #</entry><entry>Series terms</entry></row><row><entry namest="1" nameend="2" align="center" rowsep="1" /></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="7"><colspec colname="1" colwidth="35PT" align="center" /><colspec colname="2" colwidth="35PT" align="left" /><colspec colname="3" colwidth="42PT" align="left" /><colspec colname="4" colwidth="35PT" align="left" /><colspec colname="5" colwidth="42PT" align="left" /><colspec colname="6" colwidth="35PT" align="left" /><colspec colname="7" colwidth="35PT" align="left" /><tbody valign="top"><row><entry>1</entry><entry>1</entry><entry>8</entry><entry>15</entry><entry>22</entry><entry>29</entry><entry>36</entry></row><row><entry>2</entry><entry> 2</entry><entry> 9</entry><entry> 16</entry><entry> 23</entry><entry> 30</entry><entry> 37</entry></row><row><entry>3</entry><entry> 3</entry><entry> 10</entry><entry> 17</entry><entry> 24</entry><entry> 31</entry><entry> 38</entry></row><row><entry>4</entry><entry> 4</entry><entry> 11</entry><entry> 18</entry><entry> 25</entry><entry> 32</entry><entry> 39</entry></row><row><entry>5</entry><entry> 5</entry><entry> 12</entry><entry> 19</entry><entry> 26</entry><entry> 33</entry></row><row><entry>6</entry><entry> 6</entry><entry> 13</entry><entry> 20</entry><entry> 27</entry><entry> 34</entry></row><row><entry>7</entry><entry> 7</entry><entry> 14</entry><entry> 21</entry><entry> 28</entry><entry> 35</entry></row><row><entry namest="1" nameend="7" align="center" rowsep="1" /></row></tbody></tgroup></table></tables>
[0118] As seen in Table 4, each node cuts off its terms when it exceeds the length of the series. This provides an efficient means of parallel distribution of any arbitrary length series on a multi-node array (e.g., within cluster <b>10</b>, FIG. 2, and within HC <b>100</b>, FIG. 6).
[0119] For the evenly divisible case (Table 2), there is still a slight imbalance between the first and last node. This arises from the imbalance that exists across each n<sup>th </sup>interval. The first node in the array computes the first term in every n<sup>th </sup>interval, while the last node computes the last term. For most series expansions, however, where the complexity between succeeding terms does not increase rapidly, and the series is not too long, this level of balancing is sufficient.
[0120] Nonetheless, an additional level of balancing may be achieved by advancing the starting term, by one, for each node on successive intervals, and by rotating the node that computed the last term on the last interval to the first term, eliminating the computational imbalance across each n<sup>th </sup>interval. An example of this rebalancing is shown in Table 4, illustrating a forty-two term series expansion distributed on a seven node array. <tables id="TABLE-US-00005" num="5"><table frame="none" colsep="0" rowsep="0" pgwide="1"><tgroup align="left" colsep="0" rowsep="0" cols="1"><colspec colname="1" colwidth="259PT" align="center" /><thead><row><entry namest="1" nameend="1" align="center">TABLE 4</entry></row></thead><tbody valign="top"><row><entry /></row><row><entry namest="1" nameend="1" align="center" rowsep="1" /></row><row><entry>Rebalanced Example of 42-term Series Expansion in 7-Node Architecture</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="2"><colspec colname="1" colwidth="35PT" align="center" /><colspec colname="2" colwidth="224PT" align="center" /><tbody valign="top"><row><entry>Node #</entry><entry>Series terms</entry></row><row><entry namest="1" nameend="2" align="center" rowsep="1" /></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="7"><colspec colname="1" colwidth="35PT" align="center" /><colspec colname="2" colwidth="35PT" align="left" /><colspec colname="3" colwidth="42PT" align="left" /><colspec colname="4" colwidth="35PT" align="left" /><colspec colname="5" colwidth="42PT" align="left" /><colspec colname="6" colwidth="35PT" align="left" /><colspec colname="7" colwidth="35PT" align="left" /><tbody valign="top"><row><entry>1</entry><entry>1</entry><entry> 9</entry><entry> 17</entry><entry> 25</entry><entry> 33</entry><entry> 41</entry></row><row><entry>2</entry><entry> 2</entry><entry> 10</entry><entry> 18</entry><entry> 26</entry><entry> 34</entry><entry> 42</entry></row><row><entry>3</entry><entry> 3</entry><entry> 11</entry><entry> 19</entry><entry> 27</entry><entry> 35</entry><entry>36</entry></row><row><entry>4</entry><entry> 4</entry><entry> 12</entry><entry> 20</entry><entry> 28</entry><entry>29</entry><entry> 37</entry></row><row><entry>5</entry><entry> 5</entry><entry> 13</entry><entry> 21</entry><entry>22</entry><entry> 30</entry><entry> 38</entry></row><row><entry>6</entry><entry> 6</entry><entry> 14</entry><entry>15</entry><entry> 23</entry><entry> 31</entry><entry> 39</entry></row><row><entry>7</entry><entry> 7</entry><entry>8</entry><entry> 16</entry><entry> 24</entry><entry> 32</entry><entry> 40</entry></row><row><entry namest="1" nameend="7" align="center" rowsep="1" /></row></tbody></tgroup></table></tables>
[0121] The distribution of Table 4 achieves near perfect balancing regardless of how fast the computational complexity of the series terms increases, or how many terms are computed.
[0122] Examples of the effectiveness of these techniques are shown in FIG. 9 and FIG. 10. In each example of FIG. 9 and FIG. 10, one-thousand digits of PI were computed using a parallel algorithm. The parallel algorithm was based on Machin's formula which uses arctan(x), chosen because arctan(x) is computed using a power series expansion and because the computation of successive terms in the series is highly unbalanced. The computation was performed using 1, 3, 7, 15, 31, and 63 nodes in an array. Specifically, FIG. 9 illustrates the unbalanced case where the terms of the series expansion are distributed among the nodes in consecutive intervals, (e.g. 1-10, 11-20, 21-30, . . . ). The computational imbalance is evident from the disparity in the individual node compute times. FIG. 10 on the other hand illustrates the balanced case where the terms of the series expansion are distributed among the nodes as shown in Table 2. The dramatic improvement in balancing is evident by the nearly equal computation times for each of the individual nodes.
[0123] The example of computing PI using Machin's formula fits the ICNADM category since the ultimate work to be performed is Machin's formula and since that algorithm does not intrinsically move data.
[0124] A connected group of computers (e.g., the nodes of HC <b>100</b>, FIG. 6) may be used to generate both computational scale-up and computational balance, such as provided for in Tables 2-4 and FIG. 9, FIG. 10.
[0125] In certain embodiments hereof, multiple functions within the ICNADM category can be chained together to build more complex functions. For example, consider Image<sub>output</sub>=2D-CONVOLUTION (Kernel, 2D-SOBEL(image<sub>input</sub>)), where image<sub>input </sub>is a base input image to a convolution process, 2-DSOBEL is a Sobel edge detection algorithm, Kernel is the convolution Kernel (e.g., an object to find in the image), 2D-CONVOLUTION is the convolution algorithm, and Image<sub>output </sub>is the output image. Since the output of the 2D-SOBEL edge detection algorithm produces an output image array that is compatible with the input requirements of the 2D-CONVOLUTION algorithm, no data translation needs to occur.
[0126] If however data translation were required, then the above equation may take the following form: Image<sub>output</sub>=2D-CONVOLUTION (Kernel, TRANSLATE(2D-SOBEL(image<sub>input</sub>))), where image<sub>input </sub>is the base input image, 2-DSOBEL is the Sobel edge detection algorithm, Kernel is the convolution Kernel (e.g., the object to find), 2D-CONVOLUTION is the convolution algorithm, Image<sub>output </sub>is the output image, and TRANSLATE is the hypothetical function needed to make the 2D-SOBEL output compatible with the input requirements of the 2D-CONVOLUTION function. In this form the following rules may thus apply:
[0127] Rule 1: If data movement is not required for the current TRANSLATE function then it can be treated as a ICNADM category function and no further effort is required.
[0128] Rule 2: If data movement is required for the current TRANSLATE function then the entire equation can be treated as either ICADM or ECADM category algorithms, discussed in more detail below.
[0129] In accord with one embodiment hereof, a number of functions are strung together in this manner to facilitate parallel processing. ICNADM class algorithms in particular may further scale as a function of the size of the largest dataset input to the algorithms. As such, the speed of a set of processors can be at or above 90% of the sum of the speeds of processors, provided there is sufficient data. To ensure overall end-to-end parallel processing performance, the I/O of the parallel processing architecture should be set so that the time it takes to move data onto and off of the nodes is a small percentage of the time it takes to process the data.
[0130] ICNADM category algorithms can be parallelized using a compiler or a run-time translation process. Preferably, a run-time translation process is used because it affords increased flexibility (i.e., compiled solutions do not dynamically allocate computational resources that maximally fit the work required for any given problem).
[0131] Internal Context with Algorithmic Data Movement (ICADM) Category
[0132] Certain logical and mathematical problems fall into this category. By way of definition, the ICADM category is one in which the final work to be accomplished, for a particular logical or mathematical problem, is intrinsic to the algorithm and the data being transformed requires data to be moved as part of the algorithm. Example algorithms in the ICADM category include the Matrix Transpose and Gaussian elimination algorithms.
[0133] ICADM category algorithms thus require the movement of data. This means that the faster the data movement occurs, the better the algorithm scales when applied to multiple nodes or processors. One approach to solving this dilemma is to focus on faster and faster point-to-point connection speeds, which is inherently expensive, but will work for the ICADM category of algorithm.
[0134] Nonetheless, in accord with one embodiment, data movement parallelism is used for ICADM category algorithms. The following model may be used to define a minimum time and maximum node expansion model for multiple processor communications within an array of connected nodes. This model is for example implemented as a logical overlay on top of a standard switch mesh or as a fixed connection model, the former (logical overlay) being preferred: <maths id="MATH-US-00001" num="1"><math overflow="scroll"><mtable><mtr><mtd><mrow><msub><mi>H</mi><mrow><mo>(</mo><mrow><mi>P</mi><mo>,</mo><mi>C</mi><mo>,</mo><mi>t</mi></mrow><mo>)</mo></mrow></msub><mo>=</mo><mrow><mrow><mi>P</mi><mo>*</mo><msup><mrow><mo>(</mo><mrow><mi>C</mi><mo>+</mo><mn>1</mn></mrow><mo>)</mo></mrow><mi>t</mi></msup></mrow><mo>-</mo><mi>P</mi></mrow></mrow></mtd></mtr><mtr><mtd><mrow><mo>=</mo><mrow><mi>P</mi><mo>+</mo><mrow><mo>(</mo><mrow><msub><mi>H</mi><mrow><mo>(</mo><mrow><mi>P</mi><mo>,</mo><mi>C</mi><mo>,</mo><mn>1</mn></mrow><mo>)</mo></mrow></msub><mo>+</mo><mi>P</mi></mrow><mo>)</mo></mrow><mo>+</mo><mrow><mo>(</mo><mrow><msub><mi>H</mi><mrow><mo>(</mo><mrow><mi>P</mi><mo>,</mo><mi>C</mi><mo>,</mo><mn>2</mn></mrow><mo>)</mo></mrow></msub><mo>+</mo><mi>P</mi></mrow><mo>)</mo></mrow><mo>+</mo><mstyle><mtext> </mtext></mstyle><mo></mo><mi>…</mi><mo></mo><mstyle><mtext> </mtext></mstyle><mo>+</mo><mrow><mo>(</mo><mrow><msub><mi>H</mi><mrow><mo>(</mo><mrow><mi>P</mi><mo>,</mo><mi>C</mi><mo>,</mo><mrow><mi>t</mi><mo>-</mo><mn>1</mn></mrow></mrow><mo>)</mo></mrow></msub><mo>+</mo><mi>P</mi></mrow><mo>)</mo></mrow></mrow></mrow></mtd></mtr><mtr><mtd><mrow><mo>=</mo><mrow><mi>P</mi><mo>+</mo><mrow><mo>(</mo><mrow><mrow><mi>P</mi><mo>*</mo><msup><mrow><mo>(</mo><mrow><mi>C</mi><mo>+</mo><mn>1</mn></mrow><mo>)</mo></mrow><mn>1</mn></msup></mrow><mo>-</mo><mi>P</mi><mo>+</mo><mi>P</mi></mrow><mo>)</mo></mrow><mo>+</mo><mrow><mo>(</mo><mrow><mrow><mi>P</mi><mo>*</mo><msup><mrow><mo>(</mo><mrow><mi>C</mi><mo>+</mo><mn>1</mn></mrow><mo>)</mo></mrow><mn>2</mn></msup></mrow><mo>-</mo><mi>P</mi><mo>+</mo><mi>P</mi></mrow><mo>)</mo></mrow><mo>+</mo><mstyle><mtext> </mtext></mstyle><mo></mo><mi>…</mi><mo></mo><mstyle><mtext> </mtext></mstyle><mo>+</mo><mrow><mo>(</mo><mrow><mrow><mi>P</mi><mo>*</mo><msup><mrow><mo>(</mo><mrow><mi>C</mi><mo>+</mo><mn>1</mn></mrow><mo>)</mo></mrow><mrow><mi>t</mi><mo>-</mo><mn>1</mn></mrow></msup></mrow><mo>-</mo><mi>P</mi><mo>+</mo><mi>P</mi></mrow><mo>)</mo></mrow></mrow></mrow></mtd></mtr><mtr><mtd><mrow><mo>=</mo><mrow><munderover><mo>∑</mo><mrow><mi>x</mi><mo>=</mo><mn>0</mn></mrow><mrow><mi>t</mi><mo>-</mo><mn>1</mn></mrow></munderover><mo></mo><mrow><mi>P</mi><mo>*</mo><msup><mrow><mo>(</mo><mrow><mi>C</mi><mo>+</mo><mn>1</mn></mrow><mo>)</mo></mrow><mi>x</mi></msup></mrow></mrow></mrow></mtd></mtr></mtable></math><img file="US20030195938A1-20031016-M00001.TIF" id="EMI-M00001" he="61.9164" wi="274.11615" img-format="tif" img-content="mf" /><attachments><attachment idref="MATHEMATICA-00001" attachment-type="nb" file="US20030195938A1-20031016-M00001.NB" /></attachments></maths>
[0135] where, H( ) represents a HC, P is the number of processors per motherboard, C is the number of communication channels per motherboard, t is the number of expansion time units, and x is the strip number (defined below) of the HC. Equally the expansion time, t, for a given number of nodes with P processors per motherboard and C channels per motherboard can be expressed as:
<i>t=φ</i> log((<i>N+P</i>)/<i>P</i>)/log(<i>C+</i>1)κ.
[0136] Pictorially, if P=1, C=1, and t=2, then the expansion is shown as in FIG. 11.
[0137] Data communication as a logical overlay may occur in one of three ways: (1) direct cascade communication for problem-set distribution and agglomeration; (2) one-to-many of initial data to all nodes (using the cascade position to allow for the selection of the requisite data by each node); and: (3) secondary one-to-many of intermediate data to all relevant nodes required sharing data. If P=2, C=2, and t=2, then the expansion is shown in FIG. 12. Systems <b>160</b>, <b>162</b>, <b>164</b>, <b>166</b>, <b>168</b>, <b>170</b>, <b>172</b>, <b>174</b>, and <b>176</b> represent a unit containing two processors and two communication channels (P=2, C=2), whereby each system represents two nodes in a HCAS. System <b>160</b>, in this example, represents two home nodes, H<b>1</b> and H<b>2</b>. System <b>162</b> represents two processing nodes, P<b>1</b> and P<b>2</b>. System <b>164</b> represents two processing nodes, P<b>3</b> and P<b>4</b>, and so on. Table 5 illustrates a comparison of a binary tree expansion rate to an expansion rate for a HC (with select numbers of processors P and communication channels C, per node), illustrating that even moderately constructed nodes (e.g., a motherboard with one processor and one network interface card (NIC), i.e., P=1, C=1) generates almost twice the expansion rate as a binary expansion. <tables id="TABLE-US-00006" num="6"><table frame="none" colsep="0" rowsep="0"><tgroup align="left" colsep="0" rowsep="0" cols="1"><colspec colname="1" colwidth="217PT" align="center" /><thead><row><entry namest="1" nameend="1" align="center">TABLE 5</entry></row></thead><tbody valign="top"><row><entry /></row><row><entry namest="1" nameend="1" align="center" rowsep="1" /></row><row><entry>HCSA</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="7"><colspec colname="OFFSET" colwidth="14PT" align="left" /><colspec colname="1" colwidth="21PT" align="center" /><colspec colname="2" colwidth="42PT" align="center" /><colspec colname="3" colwidth="28PT" align="center" /><colspec colname="4" colwidth="42PT" align="center" /><colspec colname="5" colwidth="28PT" align="center" /><colspec colname="6" colwidth="42PT" align="center" /><tbody valign="top"><row><entry /><entry>Time</entry><entry /><entry>P = 1,</entry><entry>P = 2,</entry><entry>P = 3,</entry><entry>P = 4,</entry></row><row><entry /><entry>Units</entry><entry>Binary</entry><entry>C = 1</entry><entry>C = 2</entry><entry>C = 3</entry><entry>C = 4</entry></row><row><entry /><entry namest="OFFSET" nameend="6" align="center" rowsep="1" /></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="7"><colspec colname="OFFSET" colwidth="14PT" align="left" /><colspec colname="1" colwidth="21PT" align="center" /><colspec colname="2" colwidth="42PT" align="center" /><colspec colname="3" colwidth="28PT" align="char" char="." /><colspec colname="4" colwidth="42PT" align="char" char="." /><colspec colname="5" colwidth="28PT" align="char" char="." /><colspec colname="6" colwidth="42PT" align="char" char="." /><tbody valign="top"><row><entry /><entry>1</entry><entry>1</entry><entry>1</entry><entry>4</entry><entry>9</entry><entry>16</entry></row><row><entry /><entry>2</entry><entry>2</entry><entry>3</entry><entry>16</entry><entry>45</entry><entry>96</entry></row><row><entry /><entry>3</entry><entry>4</entry><entry>7</entry><entry>52</entry><entry>189</entry><entry>496</entry></row><row><entry /><entry>4</entry><entry>8</entry><entry>15</entry><entry>160</entry><entry>765</entry><entry>2496</entry></row><row><entry /><entry namest="OFFSET" nameend="6" align="center" rowsep="1" /></row></tbody></tgroup></table></tables>
[0138] As above, ICADM category algorithms may be parallelized using a compiler or a run-time translation process. Preferably, the run-time translation process is used because of its flexibility (compiled solutions cannot dynamically allocate computational resources that maximally fit the work required for any given problem). Scaling, on the other hand, is a function of the effective bandwidth connecting the group of nodes (e.g., via point-to-point communications or via parallel data movement speeds).
[0139] External Context No Algorithmic Data Movement (ECNADA) Category
[0140] Certain logical and mathematical problems fall into this category. By way of definition, the ECNADA category is one in which the final work to be accomplished, for a given logical or mathematical problem, is intrinsic to another algorithm and the data being transformed does not require data to be moved as part of the collective algorithms. An example of an ECNADA category algorithm is the Hadamard Transform. In accord with one embodiment herein, an ECNADA category algorithm is treated like an ICNADM category algorithm.
[0141] External Context with Algorithmic Data Movement (ECADM) Category
[0142] Many logical and mathematical problems fall into this category. By way of definition, the ECADM category is one in which the final work to be accomplished, for a given logical or mathematical problem, is intrinsic to another algorithm and the base algorithm requires data movement. An arbitrary 2D dataset for an ECADM category algorithm is illustratively shown in FIG. 66 as an array of m rows and n columns. The 2D dataset of FIG. 66 is for example a bitmap image consisting of m rows and n columns of pixels, such as generated by a digital camera. When computing a two-dimensional FFT on such an array, it is done in two stages. First, a one-dimensional (1D) FFT is computed on the first dimension (either the m rows or n columns) of the input data array to produce an intermediate array. Second, a 1D FFT is computed on the second dimension of the intermediate array to produce the final result. The row and column operations are separable, so it does not matter whether rows or columns are processed first.
[0143] Then, the m rows are distributed over the P parallel processing nodes. The rows assigned to a node are defined by the starting row, referred to as a row index, and a row count. Rows are not split across nodes, so the row count is constrained to whole numbers. The distribution is done such that equal quantities of rows are assigned to nodes 1 through P-<b>1</b> , and the remainder rows are assigned to node P. A remainder is provided to handle cases where the number of rows does not divide equally into P nodes. The remainder may be computed as the largest integer, less than the row count, such that the sum of the row count times the number of nodes and the remainder equals the total number of rows. This adjustment achieves an equalized row distribution.
[0144] This step of distributing m rows over P nodes is illustrated in FIG. 67. Processor nodes 1 through P-<b>1</b> are assigned M<sub>r </sub>rows, and processor node P is assigned the remaining R<sub>r </sub>rows. Each processor node, i, is assigned a row index, IR<sub>i</sub>, equal to (i−1) times the row count M<sub>r</sub>. This mapping of node numbers to row indices makes the data partitioning simple and efficient.
[0145] The next step evenly distributes the n columns over the P parallel processing nodes. Again, columns are not split across nodes, so the distribution is constrained to whole numbers. The column assignments are done in the same manner as for the rows. This step of distributing n columns over P nodes is illustrated in FIG. 68. Processor nodes 1 through P-<b>1</b> are assigned M<sub>c </sub>columns, and processor node P is assigned the remaining R<sub>c </sub>columns. Each processor node, i, is assigned a row index, IR<sub>i</sub>, equal to (i−1) times the column count M<sub>c</sub>.
[0146] The relationship between row and column indices, row and column counts, and node numbers is further illustrated in the context of a 7-node HC <b>100</b>A of FIG. 69, which distributes an image bitmap array consisting of 1024 rows and 700 columns on HC <b>100</b>A. In this case, the row count, M<sub>r</sub>, is 147 and the row remainder, R<sub>r</sub>, is 142. Processing node 1 is assigned rows <b>1</b> through <b>147</b>, processing node 2 is assigned rows <b>148</b> through <b>294</b>, and so on, up to processing node 7, which is assigned the remainder rows <b>883</b> through <b>1024</b>. As a check, 147 times 6 plus 142 equals the total row count of 1024. The column count, Mc, is 100 and the column remainder, R<sub>c</sub>, is 0, since the number of columns evenly divides into the number of processing nodes. Processing node 1 is assigned columns <b>1</b> through <b>100</b>, processing node 2 is assigned columns <b>101</b> through <b>200</b>, and so on, up to processing node 7 which is assigned columns <b>601</b> through <b>700</b>.
[0147] Home node <b>110</b>A performs data partitioning as described above in connection with FIG. 6. Messages describing the command and data partitioning (in this case a 2D FFT) are sent out to processing nodes <b>112</b>, <b>114</b>, <b>116</b>, <b>118</b>, <b>120</b>, <b>122</b>, <b>124</b> (i.e., processing nodes 1-7, respectively). These messages are sent as illustrated by the arrows in FIG. 69. Once processing nodes 1-7 have received their command messages, each processing node waits for home node <b>110</b>A to send data. The data is broadcast such that all processing nodes in HC <b>100</b>A receives the data at the same time. Each node also receives the entire dataset, which is much more efficient than sending each individual node a separate message with only a portion of the dataset, particularly with large numbers of parallel processing nodes in the array.
[0148] Once a node receives its data, it proceeds with computing 1D FFTs on the columns independent of other nodes. When the column FFT results are ready, they are accumulated upstream to home node <b>110</b>A, in agglomeration. Home node <b>110</b>A then broadcasts the intermediate results to the processing nodes to compute the 1D FFTs on the rows. The individual column results are then accumulated upstream to home node <b>110</b>A, in agglomeration, and into the final result to complete the process. Accordingly, data distribution through HC <b>100</b>A can accommodate arbitrary sized 2D datasets in a simple and efficient manner. In an analogous fashion, 2D FFT are parallelized for any group of processors.
[0149] These techniques thus mitigate the difficulties of cross-communication normally found with this category of algorithm. Consider for example the following algorithm: Image<sub>output</sub>=2D-CONVOLUTION (Kernel, (1D-FFT<sub>colunms </sub>(TRANSLATE<sub>fft-Transpose</sub>(1D-FFT<sub>row </sub>(image<sub>input</sub>))))), where image<sub>input </sub>is the base input image, 1D-FFT<sub>column </sub>is the column form 1D FFT, 1D-FFT<sub>row </sub>is the row form 1D FFT, Kernel is the convolution Kernel (e.g., an object to find in the input image), 2D-CONVOLUTION is the convolution algorithm, Image<sub>output </sub>is the output image, and TRANSLATE<sub>fftTranspose </sub>is the Matrix Transpose for FFT. Since the work to be accomplished is not part of the 2D-FFT (i.e., the algorithm operates to translate data to the frequency domain), the work unit of the associated function may be used to limit cross-communication. The work unit in this example is defined by the Kernel parameter of the 2D-CONVOLUTION. As long as lowest frequency per node of the 2D-FFT is at least twice the lowest frequency of the Kernel parameter, then that defines the minimum size that the image<sub>input </sub>can be split into without cross-communicate between processors. This occurs because the TRANSLATE<sub>fft-Transpose </sub>function need not retain lower frequency data (the purpose of cross-communication) when it performs its matrix Transpose. In this rendition of a 2D-FFT, the FFT is broken into three parts: 1D-FFT<sub>row</sub>, TRANSLATE<sub>fft-Transpose</sub>, and 1D-FFT<sub>column</sub>. This functional breakup allows for the insertion of the required Transpose, but also allows that function to determine whether or not to move data between processors. Accordingly, the translation between an algorithm such as ALG(A) to ALG(B), below, may occur automatically.
[0150] ALG(A)=Image<sub>output</sub>=2D-CONVOLUTION (Kernel, (2D-FFT<sub>columns </sub>(image<sub>input</sub>))), where image<sub>input </sub>is the base input image, 1D-FFT<sub>columns </sub>is the column form 1D FFT, 1D-FFTrow is the row form 1D FFT, Kernel is the convolution Kernel (e.g., the object to locate), 2D-CONVOLUTION is the convolution algorithm, and Image<sub>output </sub>(ALG(A)) is the output image. ALG(B)=Image<sub>output</sub>=2D-CONVOLUTION (Kernel, (1D-FFTcolumns (TRANSLATE<sub>fft-Transpose </sub>(1D-FFT<sub>row </sub>(image<sub>input</sub>))))), where, image<sub>input </sub>is the base input image, 1D-FFT<sub>columns </sub>is the column form 1D FFT, 1D-FFT<sub>ro </sub>is the row form 1D FFT, Kernel is the convolution Kernel (e.g., the object to locate), 2D-CONVOLUTION is the convolution algorithm, Image<sub>output </sub>(ALG(B)) is the output image, and TRANSLATE<sub>fft-Transpose </sub>is the matrix Transpose for the FFT.
[0151] ECADM category algorithms can be parallelized using a compiler or a run-time translation process. Scaling is primarily a function of work that binds the cross-communication. Preferably, the run-time translation process is used because of its flexibility (compiled solutions cannot dynamically allocate computational resources that maximally fit the work required for any given problem).
[0152] The above discussion of ICNADM, ICADM, ECNADM and ECADM category processing may be used to enhance parallel processing of algorithms through parallel processing architectures, including the HC.
[0153] Large Scale Cluster Computer Network Switch Using HC
[0154] The HC provides certain advantages over the prior art. By way of example, the HC decreases the constraints on cluster size imposed by the back plane of a switching network of a shared memory cluster. The logical form of a two dimensional, level three HC <b>100</b>B is illustrated in FIG. 14. HC <b>100</b>B is illustratively divided into sections labeled “strips.” Each strip is connectively independent of other strips. This independency remains throughout processing, until agglomeration. During agglomeration, communication may occur between the strips, as illustrated in FIG. 15. More particularly, while distributing the problem dataset, home node <b>110</b>B only communicates with top-level nodes <b>112</b>, <b>120</b>, <b>124</b>, which in turn communicate with other processing nodes of each respective strip. However, during agglomeration, inter-node communication may occur in a single direction between strips and nodes <b>112</b>, <b>120</b>, <b>124</b> at the top level, as shown in FIG. 15. Accordingly, the strips of HC <b>100</b>B are separable in terms of switch topology, as illustrated further in FIG. 16.
[0155] In FIG. 16, a router <b>200</b> is used to communicate between home node <b>110</b>B and the top nodes <b>112</b>, <b>120</b>, <b>124</b>; this inter-strip communication may occur at problem distribution and/or agglomeration. Node connectivity within a strip is achieved by using switches that accommodate the number of nodes in the strip. In FIG. 16, switch <b>202</b> provides physical connectivity (i.e., inter-strip communication) between nodes <b>112</b>, <b>114</b>, <b>116</b>, <b>118</b> in strip 1, and switch <b>204</b> provides physical connectivity (i.e., inter-strip communication) between nodes <b>120</b>, <b>122</b> in strip 2. As strip 3 has only one node <b>124</b>, connected to router <b>200</b>, no additional switch is necessary. Since very little interaction occurs at the level of router <b>200</b>, it has a negligible affect on performance. The topology of HC <b>100</b>B thus allows for extremely large clusters with little cost.
[0156] Switches such as switch <b>202</b> or <b>204</b> have a limitation on the number of nodes that can be chained together. For example, a Hewlett-Packard HP2524 switch has twenty-four ports, and can have up to sixteen switches chained together. The total number of switching ports available in this switching array is therefore 384. A strip utilizing the HP2524 as its switch within a switching array may therefore connect up to 384 nodes, for example. An example of a router suitable for router <b>200</b> is a Hewlett-Packard HD <b>9308</b>m, which has a total of 168 ports. Using the HD<b>9308</b>m as router <b>200</b> may then connect (168×384) 64,512 nodes together.
[0157] The topology of HC <b>100</b>B may be modified to break up a very large strip so that it can perform with a boundary condition such as the afore-mentioned <b>384</b>-node boundary. Specifically, in one embodiment, the number of nodes (processors) within a strip may be limited according to the following: (a) a strip has at least one HC and (b) a HC consists of the sum of HCs plus a remainder function. That is, H=H(x)+H(y)+ . . . +R(z), where H( ) is a HC with a size no larger than H, and where R( ) is 0, 1, or 2 nodes. Thus, larger cascade strips may be broken into smaller groups with a switch associated with each group, as shown in FIG. 17. By decomposing the Howard cascade into its component cascades, and by properly associating the switching network, we maintain high switching speeds while minimizing the number of switches.
[0158] The Use of Multiple Expansion Channels In Cascading Computers
[0159] Computers connected as a HC, e.g., HC <b>100</b>, FIG. 6, may utilize multiple parallel interfaces to increase the rate of problem set expansion. The general expansion rate for the HC can be expressed by the following:
<i>N</i><sub>1</sub>=2<sup>t</sup>−1 (Equation 1)
[0160] where N is the number of nodes used in expansion, and t is the number of time units used in expansion. As can be seen, the expansion rate provides the following number of nodes:{0, 1, 3, 7, 15, 31, 63, 127, 255, 511, 1023, . . . }, which is a geometric expansion of base-<b>2</b>. Equation 1 can be generalized to:
<i>N</i><sub>2</sub><i>=p*</i>(<i>p</i>+1)<sup>t</sup><i>−p</i> (Equation 2)
[0161] where N is the number of nodes used in expansion, t is the number of time units used in expansion, and p is the number of parallel channels of expansion.
[0162] A parallel channel of expansion equates to the number of bus connections (network cards, processors, etc.), which can operate simultaneously, per node. In the case of networked computers, this implies both a network interface card (NIC) and a parallel bus. For a single processor computer (which has a single internal bus), the number of channels of expansion is one and Equation 1 applies. For the case where there are two processors, per node, the following applies (per Equation 2).
<i>N</i><sub>3</sub>=2(2+1)<sup>t</sup>−2 (Equation 3)
[0163] where N is the number of nodes used in expansion, and t is the number of time units used in expansion. Table 6 below sets forth the expansion rate, in nodes, for two parallel expansion channels. <tables id="TABLE-US-00007" num="7"><table frame="none" colsep="0" rowsep="0"><tgroup align="left" colsep="0" rowsep="0" cols="1"><colspec colname="1" colwidth="217PT" align="center" /><thead><row><entry namest="1" nameend="1" align="center">TABLE 6</entry></row></thead><tbody valign="top"><row><entry /></row><row><entry namest="1" nameend="1" align="center" rowsep="1" /></row><row><entry>Expansion Rate in a HC for Two Parallel Expansion Channels</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="3"><colspec colname="OFFSET" colwidth="49PT" align="left" /><colspec colname="1" colwidth="42PT" align="center" /><colspec colname="2" colwidth="126PT" align="center" /><tbody valign="top"><row><entry /><entry>Time Units</entry><entry>Nodes</entry></row><row><entry /><entry namest="OFFSET" nameend="2" align="center" rowsep="1" /></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="3"><colspec colname="OFFSET" colwidth="49PT" align="left" /><colspec colname="1" colwidth="42PT" align="char" char="." /><colspec colname="2" colwidth="126PT" align="char" char="." /><tbody valign="top"><row><entry /><entry>0</entry><entry>0</entry></row><row><entry /><entry>1</entry><entry>4</entry></row><row><entry /><entry>2</entry><entry>16</entry></row><row><entry /><entry>3</entry><entry>52</entry></row><row><entry /><entry>4</entry><entry>160</entry></row><row><entry /><entry>5</entry><entry>484</entry></row><row><entry /><entry>6</entry><entry>1456</entry></row><row><entry /><entry>7</entry><entry>4372</entry></row><row><entry /><entry>8</entry><entry>13120</entry></row><row><entry /><entry>9</entry><entry>39364</entry></row><row><entry /><entry>10</entry><entry>118096</entry></row><row><entry /><entry>11</entry><entry>354292</entry></row><row><entry /><entry namest="OFFSET" nameend="2" align="center" rowsep="1" /></row></tbody></tgroup></table></tables>
[0164] Comparing the computations of Equation 2 with Equation 3 generates Equation 4:
<i>R=N</i><sub>1</sub>/N<sub>3</sub> (Equation 4)
[0165] where R is the ratio of expansion, N<sub>1 </sub>is the number of nodes of expansion using Equation 1, and N<sub>3 </sub>is the number of nodes of expansion using Equation 3. Table 7 sets forth the ratios of such expansion. <tables id="TABLE-US-00008" num="8"><table frame="none" colsep="0" rowsep="0"><tgroup align="left" colsep="0" rowsep="0" cols="1"><colspec colname="1" colwidth="217PT" align="center" /><thead><row><entry namest="1" nameend="1" align="center">TABLE 7</entry></row></thead><tbody valign="top"><row><entry /></row><row><entry namest="1" nameend="1" align="center" rowsep="1" /></row><row><entry>Expansion Rates and Ratios</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="5"><colspec colname="OFFSET" colwidth="14PT" align="left" /><colspec colname="1" colwidth="42PT" align="center" /><colspec colname="2" colwidth="63PT" align="center" /><colspec colname="3" colwidth="49PT" align="center" /><colspec colname="4" colwidth="49PT" align="center" /><tbody valign="top"><row><entry /><entry>Time Units</entry><entry>Nodes p = 1</entry><entry>Nodes p = 2</entry><entry>Ratio</entry></row><row><entry /><entry namest="OFFSET" nameend="4" align="center" rowsep="1" /></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="5"><colspec colname="OFFSET" colwidth="14PT" align="left" /><colspec colname="1" colwidth="42PT" align="char" char="." /><colspec colname="2" colwidth="63PT" align="char" char="." /><colspec colname="3" colwidth="49PT" align="char" char="." /><colspec colname="4" colwidth="49PT" align="center" /><tbody valign="top"><row><entry /><entry>0</entry><entry>0</entry><entry>0</entry><entry>unknown</entry></row><row><entry /><entry>1</entry><entry>1</entry><entry>4</entry><entry>0.250000</entry></row><row><entry /><entry>2</entry><entry>3</entry><entry>16</entry><entry>0.187500</entry></row><row><entry /><entry>3</entry><entry>7</entry><entry>52</entry><entry>0.134615</entry></row><row><entry /><entry>4</entry><entry>15</entry><entry>160</entry><entry>0.093750</entry></row><row><entry /><entry>5</entry><entry>31</entry><entry>484</entry><entry>0.064049</entry></row><row><entry /><entry>6</entry><entry>63</entry><entry>1456</entry><entry>0.043269</entry></row><row><entry /><entry>7</entry><entry>127</entry><entry>4372</entry><entry>0.029048</entry></row><row><entry /><entry>8</entry><entry>255</entry><entry>13120</entry><entry>0.019435</entry></row><row><entry /><entry>9</entry><entry>511</entry><entry>39364</entry><entry>0.012981</entry></row><row><entry /><entry>10</entry><entry>1023</entry><entry>118096</entry><entry>0.008662</entry></row><row><entry /><entry>11</entry><entry>2047</entry><entry>354292</entry><entry>0.005777</entry></row><row><entry /><entry namest="OFFSET" nameend="4" align="center" rowsep="1" /></row></tbody></tgroup></table></tables>
[0166] Table 7 illustrates that the first expansion provides one fourth of the expansion efficiency, as compared to the second expansion. The foregoing also illustrates that by increasing the number of dimensions of expansion using the HC, the cluster efficiency is further enhanced.
[0167]FIG. 32 and FIG. 33 illustrate representative hardware configurations for parallel channel communication; each parallel channel of expansion consist of a computer processor and a communication channel. FIG. 32 shows a first processor <b>420</b> connected to a second processor <b>422</b> via a single communication channel <b>424</b>. Multiple parallel channels of expansion, on the other hand, consist of multiple computer processors and multiple communication channels, as shown in FIG. 33. Processor <b>430</b> and processor <b>432</b> of processing node <b>112</b> are connected to processor <b>34</b> and processor <b>436</b> of processing node <b>114</b> by two communication channels <b>438</b> and <b>440</b>, respectively. Other parallel channel configurations follow and may be implemented such that a HC moves data as efficiently as desired.
[0168] The HC and Node Number
[0169] An n-dimensional HC may consume a fixed number of processing nodes for each expansion level. For example, a 2-Dimensional HC may consume a number of nodes set forth by Table 8. <tables id="TABLE-US-00009" num="9"><table frame="none" colsep="0" rowsep="0"><tgroup align="left" colsep="0" rowsep="0" cols="1"><colspec colname="1" colwidth="217PT" align="center" /><thead><row><entry namest="1" nameend="1" align="center">TABLE 8</entry></row></thead><tbody valign="top"><row><entry /></row><row><entry namest="1" nameend="1" align="center" rowsep="1" /></row><row><entry>Node Expansion for N-dimensional HC</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="3"><colspec colname="OFFSET" colwidth="28PT" align="left" /><colspec colname="1" colwidth="63PT" align="center" /><colspec colname="2" colwidth="126PT" align="center" /><tbody valign="top"><row><entry /><entry>Expansion Level</entry><entry>Number of nodes</entry></row><row><entry /><entry namest="OFFSET" nameend="2" align="center" rowsep="1" /></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="3"><colspec colname="OFFSET" colwidth="28PT" align="left" /><colspec colname="1" colwidth="63PT" align="center" /><colspec colname="2" colwidth="126PT" align="char" char="." /><tbody valign="top"><row><entry /><entry>1</entry><entry>1</entry></row><row><entry /><entry>2</entry><entry>3</entry></row><row><entry /><entry>3</entry><entry>7</entry></row><row><entry /><entry>4</entry><entry>15</entry></row><row><entry /><entry>5</entry><entry>31</entry></row><row><entry /><entry>6</entry><entry>63</entry></row><row><entry /><entry>. . .</entry><entry>. . .</entry></row><row><entry /><entry>N</entry><entry>(p + 1)<sup>expansion level</sup> − 1</entry></row><row><entry /><entry namest="OFFSET" nameend="2" align="center" rowsep="1" /></row></tbody></tgroup></table></tables>
[0170] To be able to consume, for example, 4, 5, or 6 nodes in a two-dimensional HC, and extension to the above node expansion algorithm for the HC may be used, as described below.
[0171] More particularly, the expansion of a HC in time may be shown in the sets defined below. The natural cascade frequency of the HC is the expansion rate of Set 1:{1, 3, 7, 15, 31, . . . , 2<sup>t</sup>−1}, where t is the number of time units. In the more general case, Set 1 takes the form of Set <b>2: {d</b><sup>1</sup>−1, d<sup>2</sup>−1, d<sup>3</sup>−1, d<sup>t</sup>−1}, where d is the number of dimensions of expansion (=(p+1)), p is the number of parallel channels of expansion, and t is the number of time units. It is advantageous to obtain an expansion number that lies between the elements of Set 2. For example consider Set 3: {d<sup>1</sup>−1+1, d<sup>1</sup>−1+2, d<sup>1</sup>−1+3, . . . , d<sup>2</sup>−2}, where d is the number o dimensions of expansion (=p+1), and p is the number of parallel channels of expansion. Set 3 more specifically shows the set of values that lie between the first and second natural expansion terms. In the case of Set 1, this translates into Set 4: {2}. The general set for an in-between number of nodes is given in Set 5: {d<sup>t</sup>−1+1, d<sup>t</sup>−1+2, d<sup>t</sup>−1+3, . . . , d<sup>t</sup>−1+t, d<sup>(t+1)</sup>−2}, where d is the number of dimensions of expansion (=p+1), p is the number of parallel channels of expansion, and t is the number of time units. The general HC construction series is then given below: <maths id="MATH-US-00002" num="2"><math overflow="scroll"><mtable><mtr><mtd><mrow><mi>N</mi><mo>=</mo><mrow><msub><mi>H</mi><mrow><mo>(</mo><mrow><mi>P</mi><mo>,</mo><mo>,</mo><mi>t</mi></mrow><mo>)</mo></mrow></msub><mo>=</mo><mrow><mrow><mi>P</mi><mo>*</mo><msup><mrow><mo>(</mo><mrow><mi>P</mi><mo>+</mo><mn>1</mn></mrow><mo>)</mo></mrow><mi>t</mi></msup></mrow><mo>-</mo><mi>P</mi></mrow></mrow></mrow></mtd></mtr><mtr><mtd><mrow><mo>=</mo><mrow><mi>P</mi><mo>+</mo><mrow><mo>(</mo><mrow><msub><mi>H</mi><mrow><mo>(</mo><mrow><mi>P</mi><mo>,</mo><mo>,</mo><mn>1</mn></mrow><mo>)</mo></mrow></msub><mo>+</mo><mi>P</mi></mrow><mo>)</mo></mrow><mo>+</mo><mrow><mo>(</mo><mrow><msub><mi>H</mi><mrow><mo>(</mo><mrow><mi>P</mi><mo>,</mo><mo>,</mo><mn>2</mn></mrow><mo>)</mo></mrow></msub><mo>+</mo><mi>P</mi></mrow><mo>)</mo></mrow><mo>+</mo><mstyle><mtext> </mtext></mstyle><mo></mo><mi>…</mi><mo></mo><mstyle><mtext> </mtext></mstyle><mo>+</mo><mrow><mo>(</mo><mrow><msub><mi>H</mi><mrow><mo>(</mo><mrow><mi>P</mi><mo>,</mo><mrow><mi>t</mi><mo>-</mo><mn>1</mn></mrow></mrow><mo>)</mo></mrow></msub><mo>+</mo><mi>P</mi></mrow><mo>)</mo></mrow></mrow></mrow></mtd></mtr><mtr><mtd><mrow><mo>=</mo><mrow><mi>P</mi><mo>+</mo><mrow><mo>(</mo><mrow><mrow><mi>P</mi><mo>*</mo><msup><mrow><mo>(</mo><mrow><mi>P</mi><mo>+</mo><mn>1</mn></mrow><mo>)</mo></mrow><mn>1</mn></msup></mrow><mo>-</mo><mi>P</mi><mo>+</mo><mi>P</mi></mrow><mo>)</mo></mrow><mo>+</mo><mrow><mo>(</mo><mrow><mrow><mi>P</mi><mo>*</mo><msup><mrow><mo>(</mo><mrow><mi>P</mi><mo>+</mo><mn>1</mn></mrow><mo>)</mo></mrow><mn>2</mn></msup></mrow><mo>-</mo><mi>P</mi><mo>+</mo><mi>P</mi></mrow><mo>)</mo></mrow><mo>+</mo><mstyle><mtext> </mtext></mstyle><mo></mo><mi>…</mi><mo></mo><mstyle><mtext> </mtext></mstyle><mo>+</mo><mrow><mo>(</mo><mrow><mrow><mi>P</mi><mo>*</mo><msup><mrow><mo>(</mo><mrow><mi>P</mi><mo>+</mo><mn>1</mn></mrow><mo>)</mo></mrow><mrow><mi>t</mi><mo>-</mo><mn>1</mn></mrow></msup></mrow><mo>-</mo><mi>P</mi><mo>+</mo><mi>P</mi></mrow><mo>)</mo></mrow></mrow></mrow></mtd></mtr><mtr><mtd><mrow><mo>=</mo><mrow><munderover><mo>∑</mo><mrow><mi>x</mi><mo>=</mo><mn>0</mn></mrow><mrow><mi>t</mi><mo>-</mo><mn>1</mn></mrow></munderover><mo></mo><mrow><mi>P</mi><mo>*</mo><msup><mrow><mo>(</mo><mrow><mi>P</mi><mo>+</mo><mn>1</mn></mrow><mo>)</mo></mrow><mi>x</mi></msup></mrow></mrow></mrow></mtd></mtr></mtable></math><img file="US20030195938A1-20031016-M00002.TIF" id="EMI-M00002" he="61.9164" wi="274.11615" img-format="tif" img-content="mf" /><attachments><attachment idref="MATHEMATICA-00002" attachment-type="nb" file="US20030195938A1-20031016-M00002.NB" /></attachments></maths>
[0172] As can be seen, each term of the expansion is equal to the sum of all proceeding terms. Since each term corresponds to a cascade strip, the HC is balanced by adding additional nodes starting with the next highest term and continuing until all empty potential slots are filled, or by evenly spreading extra-nodes among the strips. Keeping the HC balanced by evenly spreading nodes among the cascade strips is accomplished by level-spreading among the nodes, as illustrated in Table 9 and Table 10: <tables id="TABLE-US-00010" num="10"><table frame="none" colsep="0" rowsep="0"><tgroup align="left" colsep="0" rowsep="0" cols="1"><colspec colname="1" colwidth="217PT" align="center" /><thead><row><entry namest="1" nameend="1" align="center">TABLE 9</entry></row></thead><tbody valign="top"><row><entry /></row><row><entry namest="1" nameend="1" align="center" rowsep="1" /></row><row><entry>A Three-level, Seven-Node HC with Strip Boundaries</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="4"><colspec colname="OFFSET" colwidth="14PT" align="left" /><colspec colname="1" colwidth="70PT" align="left" /><colspec colname="2" colwidth="42PT" align="left" /><colspec colname="3" colwidth="91PT" align="left" /><tbody valign="top"><row><entry /><entry>Strip 1</entry><entry>Strip 2</entry><entry>Strip 3</entry></row><row><entry /><entry namest="OFFSET" nameend="3" align="center" rowsep="1" /></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="5"><colspec colname="OFFSET" colwidth="14PT" align="left" /><colspec colname="1" colwidth="70PT" align="left" /><colspec colname="2" colwidth="42PT" align="left" /><colspec colname="3" colwidth="35PT" align="left" /><colspec colname="4" colwidth="56PT" align="center" /><tbody valign="top"><row><entry /><entry>Node 1</entry><entry>Node 5</entry><entry>Node 7</entry><entry>Level 1</entry></row><row><entry /><entry>Node 2 Node 4</entry><entry>Node 6</entry><entry /><entry>Level 2</entry></row><row><entry /><entry>Node 3</entry><entry /><entry /><entry>Level 3</entry></row><row><entry /><entry namest="OFFSET" nameend="4" align="center" rowsep="1" /></row></tbody></tgroup></table></tables>
[0173]<tables id="TABLE-US-00011" num="11"><table frame="none" colsep="0" rowsep="0"><tgroup align="left" colsep="0" rowsep="0" cols="1"><colspec colname="1" colwidth="217PT" align="center" /><thead><row><entry namest="1" nameend="1" align="center">TABLE 10</entry></row></thead><tbody valign="top"><row><entry /></row><row><entry namest="1" nameend="1" align="center" rowsep="1" /></row><row><entry>A Three-level, Eight-Node HC with Strip Boundaries</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="5"><colspec colname="OFFSET" colwidth="14PT" align="left" /><colspec colname="1" colwidth="56PT" align="left" /><colspec colname="2" colwidth="35PT" align="left" /><colspec colname="3" colwidth="42PT" align="left" /><colspec colname="4" colwidth="70PT" align="left" /><tbody valign="top"><row><entry /><entry>Strip 1</entry><entry>Strip 2</entry><entry>Strip 3</entry><entry>Strip 4</entry></row><row><entry /><entry namest="OFFSET" nameend="4" align="center" rowsep="1" /></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="6"><colspec colname="OFFSET" colwidth="14PT" align="left" /><colspec colname="1" colwidth="56PT" align="left" /><colspec colname="2" colwidth="35PT" align="left" /><colspec colname="3" colwidth="42PT" align="left" /><colspec colname="4" colwidth="28PT" align="left" /><colspec colname="5" colwidth="42PT" align="center" /><tbody valign="top"><row><entry /><entry>Node 1</entry><entry>Node 5</entry><entry>Node 7</entry><entry>Node 8</entry><entry>Level 1</entry></row><row><entry /><entry>Node 2 Node 4</entry><entry>Node 6</entry><entry /><entry /><entry>Level 2</entry></row><row><entry /><entry>Node 3</entry><entry /><entry /><entry /><entry>Level 3</entry></row><row><entry /><entry namest="OFFSET" nameend="5" align="center" rowsep="1" /></row></tbody></tgroup></table></tables>
[0174] Since each strip boundary and level represents a time unit, and since the position of the nodes in a strip and on a level represents the distribution that may occur in that time slot, adding an additional node 8 in Table 10 increases the overall processing time (as compared to Table 9) because of the additional time unit. Unlike Table 10, where the additional node 8 was added to the top-most level, Table 11 illustrates how additional nodes may be added without the additional time cost. <tables id="TABLE-US-00012" num="12"><table frame="none" colsep="0" rowsep="0"><tgroup align="left" colsep="0" rowsep="0" cols="1"><colspec colname="1" colwidth="217PT" align="center" /><thead><row><entry namest="1" nameend="1" align="center">TABLE 11</entry></row></thead><tbody valign="top"><row><entry /></row><row><entry namest="1" nameend="1" align="center" rowsep="1" /></row><row><entry>Four-Level, Four-Strip HC</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="4"><colspec colname="1" colwidth="77PT" align="left" /><colspec colname="2" colwidth="56PT" align="left" /><colspec colname="3" colwidth="28PT" align="left" /><colspec colname="4" colwidth="56PT" align="left" /><tbody valign="top"><row><entry>Strip 1</entry><entry>Strip 2</entry><entry>Strip 3</entry><entry>Strip 4</entry></row><row><entry namest="1" nameend="4" align="center" rowsep="1" /></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="5"><colspec colname="1" colwidth="77PT" align="left" /><colspec colname="2" colwidth="56PT" align="left" /><colspec colname="3" colwidth="28PT" align="left" /><colspec colname="4" colwidth="28PT" align="left" /><colspec colname="5" colwidth="28PT" align="center" /><tbody valign="top"><row><entry>Node1</entry><entry>Node9</entry><entry>Node13</entry><entry>Node15</entry><entry>Level 1</entry></row><row><entry>Node2 Node5 Node8</entry><entry>Node 10 Node12</entry><entry>Node14</entry><entry /><entry>Level 2</entry></row><row><entry>Node3 Node6 Node7</entry><entry>Node11</entry><entry /><entry /><entry>Level 3</entry></row><row><entry>Node4</entry><entry /><entry /><entry /><entry>Level 4</entry></row><row><entry namest="1" nameend="5" align="center" rowsep="1" /></row></tbody></tgroup></table></tables>
[0175] One way to balance the HC is to fill node spaces by level, to ensure that the nodes are aligned in time. Table 12 illustrates one technique for adding nodes with reference to the HC of Table 11: <tables id="TABLE-US-00013" num="13"><table frame="none" colsep="0" rowsep="0"><tgroup align="left" colsep="0" rowsep="0" cols="1"><colspec colname="1" colwidth="217PT" align="center" /><thead><row><entry namest="1" nameend="1" align="center">TABLE 12</entry></row></thead><tbody valign="top"><row><entry /></row><row><entry namest="1" nameend="1" align="center" rowsep="1" /></row><row><entry>Adding Nodes to a Balanced HC</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="2"><colspec colname="1" colwidth="56PT" align="center" /><colspec colname="2" colwidth="161PT" align="left" /><tbody valign="top"><row><entry>Number of</entry><entry /></row><row><entry>Added Nodes</entry><entry>Which Nodes are added</entry></row><row><entry namest="1" nameend="2" align="center" rowsep="1" /></row><row><entry>1</entry><entry>Node 15 (A)</entry></row><row><entry>2</entry><entry>Node 15 (A) and Node 8 (B)</entry></row><row><entry>3</entry><entry>Node 15 (A) and Node 8 (B) and Node 12 (B)</entry></row><row><entry>4</entry><entry>Node 15 (A) and Node 8 (B) and Node 12 (B) and</entry></row><row><entry /><entry>Node 14 (B)</entry></row><row><entry>5</entry><entry>Node 15 (A) and Node 8 (B) and Node 12 (B) and</entry></row><row><entry /><entry>Node 14 (B) and Node 6 (C)</entry></row><row><entry>6</entry><entry>Node 15 (A) and Node 8 (B) and Node 12 (B) and</entry></row><row><entry /><entry>Node 14 (B) and Node 6 (C) and Node 11 (C)</entry></row><row><entry>7</entry><entry>Node 15 (A) and Node 8 (B) and Node 12 (B) and</entry></row><row><entry /><entry>Node 14 (B) and Node 6 (C) and Node 11 (C) and</entry></row><row><entry /><entry>Node 7 (D)</entry></row><row><entry>8</entry><entry>Node 15 (A) and Node 8 (B) and Node 12 (B) and</entry></row><row><entry /><entry>Node 14 (B) and Node 6 (C) and Node 11 (C) and</entry></row><row><entry /><entry>Node 7 (D) and Node 4 (D)</entry></row><row><entry namest="1" nameend="2" align="center" rowsep="1" /></row></tbody></tgroup></table></tables>
[0176] In Table 12, any of the nodes shown in the same type (A; B, C or D, respectively ) can replace any other node of the same type in placement order. By ensuring that the nodes are added in a time efficient manner, HC system overhead is reduced.
[0177] Pre-Building Agglomeration Communication Paths
[0178] Typically in the prior art, implementations using multiple network interfaced cards (“NICs”) require that the NICs are bonded at the device driver level. This makes off-the-shelf device driver upgrades unavailable, and removes the ability to use the NICs independently of each other to create multiple communication channels. By allowing the software application level to control the NMCs independently, multiple communication channels can be created and used in different ways, giving greater communication flexibility. For example, multiple communication channels can be made between two nodes for increasing the communication bandwidth, or the communication channels can be used independently to allow one node to communicate with many other nodes concurrently. Another advantage is obtained as a result of the requirement to have a physically separate switch network for each parallel communication channel. This requirement provides channel redundancy commensurate with the number of parallel communication channels implemented. FIG. 70 shows an example of connectivity between two processing nodes, <b>112</b> and <b>114</b>, each with two NICs <b>244</b><i>a</i>, <b>244</b><i>b</i>, <b>244</b><i>c </i>and <b>244</b><i>d</i>, respectively. Each parallel communication channel has an independent network switch, <b>260</b> and <b>262</b>. Processing node <b>112</b> has two possible communication channels with processing node <b>114</b>; a) using NIC <b>244</b><i>a </i>that connects to network switch <b>260</b> and thus to NIC <b>244</b><i>c </i>in processing node <b>114</b>, and b) using NIC <b>244</b><i>b </i>that connects to network switch <b>262</b> and thus to NIC <b>244</b><i>d </i>in processing node <b>114</b>. It is preferable, but not necessary, that each node in the cluster have the same number of NICs, and hence parallel communication channels, and thereby connection to all switch networks for optimal communication throughput.
[0179]FIG. 71 illustrates how multiple communication channels, <b>448</b> and <b>448</b>′ in this example, may be implemented on a processing node <b>112</b>. A software application, (for example, control software <b>94</b> of FIG. 5), may consist of multiple threads. Each thread uses a multiple channel software API <b>446</b>, to facilitate communication with other nodes. API <b>446</b> consists of a library of thread-safe subroutines that coordinate use of communication channels <b>448</b> and <b>448</b>′ in a channel resource pool <b>447</b>. In this example, each channel, <b>448</b>, <b>448</b>′, utilizes a command/control thread, <b>450</b>, <b>450</b>′, respectively, to communicate with the specific network interface device drivers, <b>452</b> and <b>452</b>′. Command/control threads <b>450</b>, <b>450</b>′, use network interface device drivers <b>452</b>, <b>452</b>′, respectively, for handling specific communication protocols. API <b>446</b> decouples application threads <b>442</b> and <b>444</b> from specific network hardware and protocol knowledge. API <b>446</b> allows application threads to use one or more channels for a communication, and manages the data distribution and reconstruction across selected channels as necessary. API <b>446</b> may also handle channel protocols, detecting and recovering from channel failures.
[0180]FIG. 72 illustrates how multiple communication channels may be used to gain additional efficiency when utilized on an HCAS. In this example, HC <b>100</b> consists of one home node <b>110</b> and seven processing nodes <b>112</b>, <b>114</b>, <b>116</b>, <b>118</b>, <b>120</b>, <b>122</b>, <b>124</b> and <b>126</b>. Each processing node, <b>112</b>, <b>114</b>, <b>116</b>, <b>118</b>, <b>120</b>, <b>122</b>, <b>124</b> and <b>126</b>, has two parallel communication channels. One channel is used for problem expansion on HC <b>100</b>, and the other channel is used for agglomeration on HC <b>100</b>. The connecting lines in FIG. 72 represent messages being passed between nodes. The style of the line indicates the time unit during which the message is passed, as shown in key <b>456</b>. Home node <b>110</b> sends a processing request to processing node 1, <b>112</b>, during time unit <b>1</b>. During time unit <b>2</b>, home node <b>110</b> sends the processing request to processing node 5, <b>120</b>. Processing node 1, <b>112</b>, sends the processing request to processing node 2, <b>114</b>, and configures its second communication channel back to home node <b>110</b> ready for agglomeration. During time unit <b>3</b> home node <b>110</b> sends the processing request to processing node 7, <b>124</b>. Processing node 1, <b>112</b>, sends the processing request to processing node 4, <b>118</b>, using its first communication channel. Processing node 2, <b>114</b>, sends the processing request to processing node 3, <b>116</b>, using its first communication channel and configures its second communication channel back to processing node 2, <b>114</b>, ready for agglomeration. Processing node 5, <b>120</b>, sends the processing request to processing node 6, <b>122</b>, using its first communication channel and configures its second communication channel to processing node 1, <b>112</b>, ready for agglomeration. During time unit <b>4</b> processing node 3 configures its second communication channel to processing node 2, <b>114</b>, ready for agglomeration. Processing node 4, <b>118</b>, configures its second communication channel to processing node 1, <b>112</b>, ready for agglomeration. Processing node 6, <b>112</b>, configures its second communication channel to processing node 5, <b>120</b>, ready for agglomeration. Processing node 7, <b>124</b>, configures its second communication channel to processing node 5, <b>120</b>, ready for agglomeration.
[0181] As shown, after 3 time units the processing request has been sent to all 7 processing nodes. During time unit <b>4</b> the processing nodes, if expecting data, configure their first communication channel to receive the data broadcast. After 4 time units the full data agglomeration communication path has been established, thus saving channel configuration overhead normally incurred prior to data agglomeration as in the case when only one communication channel is available on processing nodes.
[0182] Processing Nodes as a Home Node
[0183] If each processing node contains an additional NIC associated with the home node switch network, then a processing node can be used in place of a home node. For example, if a home node fails, either processing nodes or other, home nodes will detect the lack of communication. In one embodiment, the lowest numbered communicating home node selects one of its processing nodes to reconfigure as a new home node by terminating the active processing node software and by starting new home node software through a remote procedure call. The failed home node's assigned processing nodes are then reassigned to the new home node, and processing continues. This is discussed further in connection with FIG. 18, FIG. 21-FIG. 25.
[0184] Multiple Home Nodes In Cascading Cluster Systems
[0185] A HC may have multiple home nodes operating in tandem. Such a HC may further automate detection and mitigation of failed home nodes. In particular, each HC cluster within a HCAS may reconfigure itself in the event of a dropped out node, and without human intervention. Since the replacement home node operates like the failed node, the HC functions fully for parallel processing. Moreover, multiple home nodes facilitate access by multiple remote hosts.
[0186] In order to have multiple home nodes allocating processing nodes at the same time, the home nodes share data. FIG. 18 shows the communication paths for multiple gateway nodes in one HC <b>100</b>C. FIG. 18 illustratively has nine processing nodes <b>112</b>, <b>114</b>, <b>116</b>, <b>118</b>, <b>120</b>, <b>122</b>, <b>124</b>, <b>126</b>, <b>128</b> with three home nodes <b>110</b>C, <b>224</b>, <b>226</b>. Each home node can communicate with each of the other home nodes. At system startup, each home node has a list of processing nodes for which it is responsible; this is indicated in FIG. 18 by identifying the home node on the communication line between a processing node switch network <b>220</b> and the processing nodes. For example, processing nodes <b>118</b>, <b>120</b>, <b>122</b> are under control of home node <b>224</b> as part of the same cluster network. Table 13 sets forth all associations of HC <b>100</b>C.
[0187] If any of home nodes <b>110</b>C, <b>224</b>, <b>226</b> require additional processing nodes, then it looks to the list of processing nodes for any other home node to see if a free processing node exists; switch <b>220</b> then reconfigures to bring the new processing node under control of the requesting home node. In another embodiment, a home node needing processing power can broadcast a request for additional processing nodes to each of the other home nodes; these nodes respond and, if a free processing node exists, HC <b>100</b>C reconfigures to adjust the nodes processing for a particular home node. <tables id="TABLE-US-00014" num="14"><table frame="none" colsep="0" rowsep="0"><tgroup align="left" colsep="0" rowsep="0" cols="1"><colspec colname="1" colwidth="217PT" align="center" /><thead><row><entry namest="1" nameend="1" align="center">TABLE 13</entry></row></thead><tbody valign="top"><row><entry /></row><row><entry namest="1" nameend="1" align="center" rowsep="1" /></row><row><entry>One Multi-Home Node HC with Processing Node Associations</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="2"><colspec colname="1" colwidth="56PT" align="center" /><colspec colname="2" colwidth="161PT" align="left" /><tbody valign="top"><row><entry>home node</entry><entry>processing node list</entry></row><row><entry namest="1" nameend="2" align="center" rowsep="1" /></row><row><entry>home node 1</entry><entry>Processing Node 1, Processing Node 2, Processing</entry></row><row><entry /><entry>Node 3</entry></row><row><entry>home node 2</entry><entry>Processing Node 4, Processing Node 5, Processing</entry></row><row><entry /><entry>Node 6</entry></row><row><entry>home node 3</entry><entry>Processing Node 7, Processing Node 8, Processing</entry></row><row><entry /><entry>Node 9</entry></row><row><entry namest="1" nameend="2" align="center" rowsep="1" /></row></tbody></tgroup></table></tables>
[0188] In one example of operation of HC <b>100</b>C, FIG. 18, all nodes in home node 1 (<b>110</b>C) and home node 2 (<b>224</b>) are idle and home node 3 (<b>226</b>) needs seven processing nodes. Home node 3 issues a request for four more processing nodes in a single multi-cast message to both home node 2 and home node 1. The request also contains the total number of home nodes addressed by the message. When home node 1 receives the multi-cast message and determines that the total number of home nodes addressed is two, it calculates that it should return two free processing nodes by dividing the total number of nodes requested by the total number of home nodes addressed by the message, and sends a return message identifying processing node 1 and processing node 2. When home node 2 receives the multi-cast message, it sends a return message identifying processing nodes 4 and 5.
[0189] In another example of operation of HC <b>100</b>C, home node 3 requires only one additional processing node. Home node 3 issues a multi-cast message request for one processing node to home node 1 and home node 2. The message also contains the total number of home nodes addressed by the multi-cast message. When home node 1 receives the multi-cast message, it sends a return message identifying only processing node 1, as it is the lowest numbered home node and the request was for a single processing node. When home node 2 receives the multi-cast message, it recognizes that it does not need to send a return message to home node 3 because it will recognize that home nodes with lower numbers have fulfilled the request.
[0190]FIG. 19 shows one NIC configuration for a home node <b>110</b>D. Home node <b>110</b>D has two network interface cards, NIC <b>240</b> and NIC <b>242</b>. NIC <b>240</b> is connected to processing node switch network <b>220</b>, FIG. 18, and NIC <b>242</b> is connected to home node switch network <b>222</b>, FIG. 18. FIG. 20 shows further connectivity of the NIC configuration relative to a processing node <b>112</b>. Processing node <b>112</b> has a single network interface card NIC <b>244</b> connected to processing node switch network <b>220</b>. One difference between the home node configuration of FIG. 19 and the processing node configuration of FIG. 20, in terms of network connections, is that processing node <b>112</b> does not have an NIC connected to home node switch network <b>222</b>. Home nodes and processing nodes contain different software; however, a remote procedure call (“RPC”) may be used to change the software configuration of a node. Therefore, if processing node <b>112</b> contains a second network interface card connected to home node switch network <b>222</b>, it may reconfigure to operate as a home node by the RPC.
[0191] The Use of Overlapped Data to Decrease I/O Accesses in Clustered Computers
[0192] In one embodiment, the HC serves to eliminate the large number of data transfers common in shared memory clusters of the prior art; such a HC provides computing efficiencies in the 40% to 55% range, as compared to the 3% to 7% range of the prior art. Coupling transfers, the HC efficiency is in the 80% to 95% range. In one example, an algorithm runs on two processing nodes of the HC, processing node 1 and processing node 2. A 5×6 element matrix is divided into two 5×3 element matrices for parallel processing by the two processing nodes. The data for the processing node 1 is shown in Table 14, and the data for processing node 2 is shown in Table 15. <tables id="TABLE-US-00015" num="15"><table frame="none" colsep="0" rowsep="0"><tgroup align="left" colsep="0" rowsep="0" cols="1"><colspec colname="1" colwidth="217PT" align="center" /><thead><row><entry namest="1" nameend="1" align="center">TABLE 14</entry></row><row><entry /></row><row><entry namest="1" nameend="1" align="center" rowsep="1" /></row><row><entry>Example 5 × 3 Matrix</entry></row><row><entry namest="1" nameend="1" align="center" rowsep="1" /></row></thead><tbody valign="top"><row><entry /></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="5"><colspec colname="1" colwidth="63PT" align="char" char="." /><colspec colname="2" colwidth="14PT" align="char" char="." /><colspec colname="3" colwidth="63PT" align="char" char="." /><colspec colname="4" colwidth="14PT" align="char" char="." /><colspec colname="5" colwidth="63PT" align="char" char="." /><tbody valign="top"><row><entry>01</entry><entry>02</entry><entry>03</entry><entry>04</entry><entry>05</entry></row><row><entry>11</entry><entry>12</entry><entry>13</entry><entry>14</entry><entry>15</entry></row><row><entry>21</entry><entry>22</entry><entry>23</entry><entry>24</entry><entry>25</entry></row><row><entry namest="1" nameend="5" align="center" rowsep="1" /></row></tbody></tgroup></table></tables>
[0193]<tables id="TABLE-US-00016" num="16"><table frame="none" colsep="0" rowsep="0"><tgroup align="left" colsep="0" rowsep="0" cols="1"><colspec colname="1" colwidth="217PT" align="center" /><thead><row><entry namest="1" nameend="1" align="center">TABLE 15</entry></row><row><entry /></row><row><entry namest="1" nameend="1" align="center" rowsep="1" /></row><row><entry>Example 5 × 3 Matrix</entry></row><row><entry namest="1" nameend="1" align="center" rowsep="1" /></row></thead><tbody valign="top"><row><entry /></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="5"><colspec colname="1" colwidth="63PT" align="char" char="." /><colspec colname="2" colwidth="14PT" align="char" char="." /><colspec colname="3" colwidth="63PT" align="char" char="." /><colspec colname="4" colwidth="14PT" align="char" char="." /><colspec colname="5" colwidth="63PT" align="char" char="." /><tbody valign="top"><row><entry>06</entry><entry>07</entry><entry>08</entry><entry>09</entry><entry>10</entry></row><row><entry>16</entry><entry>17</entry><entry>18</entry><entry>19</entry><entry>20</entry></row><row><entry>26</entry><entry>27</entry><entry>28</entry><entry>29</entry><entry>30</entry></row><row><entry namest="1" nameend="5" align="center" rowsep="1" /></row></tbody></tgroup></table></tables>
[0194] In this example it is assumed that processing node 1 and processing node 2 need to share data items 5, 6, 15, 16, 25, and 26. The sequence for processing and transferring the shared data items is shown in Table 16. <tables id="TABLE-US-00017" num="17"><table frame="none" colsep="0" rowsep="0"><tgroup align="left" colsep="0" rowsep="0" cols="1"><colspec colname="1" colwidth="217PT" align="center" /><thead><row><entry namest="1" nameend="1" align="center">TABLE 16</entry></row></thead><tbody valign="top"><row><entry /></row><row><entry namest="1" nameend="1" align="center" rowsep="1" /></row><row><entry>Sequence and Transfers in Example HC</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="3"><colspec colname="OFFSET" colwidth="14PT" align="left" /><colspec colname="1" colwidth="84PT" align="left" /><colspec colname="2" colwidth="119PT" align="left" /><tbody valign="top"><row><entry /><entry>Processing Node 1</entry><entry>Processing Node 2</entry></row><row><entry /><entry namest="OFFSET" nameend="2" align="center" rowsep="1" /></row><row><entry /><entry>Process Data Item - 05</entry><entry>Process Data Item - 06</entry></row><row><entry /><entry>Transfer Data Item - 05</entry><entry>Receive Data Item - 05 from Node 1</entry></row><row><entry /><entry>to Node 2</entry></row><row><entry /><entry>Receive Data Item - 06</entry><entry>Transfer Data Item - 06 to Node 1</entry></row><row><entry /><entry>from Node 2</entry></row><row><entry /><entry>Process Data Item - 06</entry><entry>Process Data Item - 05</entry></row><row><entry /><entry>Process Data Item - 15</entry><entry>Process Data Item - 16</entry></row><row><entry /><entry>Transfer Data Item - 15</entry><entry>Receive Data Item - 15 from Node 1</entry></row><row><entry /><entry>to Node 2</entry></row><row><entry /><entry>Receive Data Item - 16</entry><entry>Transfer Data Item - 16 to Node 1</entry></row><row><entry /><entry>from Node 2</entry></row><row><entry /><entry>Process Data Item - 16</entry><entry>Process Data Item - 16</entry></row><row><entry /><entry>Process Data Item - 25</entry><entry>Process Data Item - 26</entry></row><row><entry /><entry>Transfer Data Item - 25</entry><entry>Receive Data Item - 25 from Node 1</entry></row><row><entry /><entry>to Node 2</entry></row><row><entry /><entry>Receive Data Item - 26</entry><entry>Transfer Data Item - 26 to Node 1</entry></row><row><entry /><entry>from Node 2</entry></row><row><entry /><entry>Process Data Item - 26</entry><entry>Process Data Item - 25</entry></row><row><entry /><entry namest="OFFSET" nameend="2" align="center" rowsep="1" /></row></tbody></tgroup></table></tables>
[0195] As can be seen in Table 16, the processing by the HC generates twelve data transfer/receives for only six boundary data processes. By changing the boundary data such that the shared data is overlapped (to provide the shared data on the required nodes when needed), the number of required data transfers decreases. This is important as processing speed is compromised with a large number of data transfers. Using the data from the above example, the 5×6 element matrix is divided into two 6×3 element matrices, as shown in Tables 17 and 18. <tables id="TABLE-US-00018" num="18"><table frame="none" colsep="0" rowsep="0"><tgroup align="left" colsep="0" rowsep="0" cols="1"><colspec colname="1" colwidth="217PT" align="center" /><thead><row><entry namest="1" nameend="1" align="center">TABLE 17</entry></row></thead><tbody valign="top"><row><entry /></row><row><entry namest="1" nameend="1" align="center" rowsep="1" /></row><row><entry>Example 6x3 Matrix</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="2"><colspec colname="OFFSET" colwidth="175PT" align="left" /><colspec colname="1" colwidth="42PT" align="center" /><tbody valign="top"><row><entry /><entry>a</entry></row><row><entry /><entry namest="OFFSET" nameend="1" align="center" rowsep="1" /></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="7"><colspec colname="OFFSET" colwidth="14PT" align="left" /><colspec colname="1" colwidth="21PT" align="center" /><colspec colname="2" colwidth="49PT" align="center" /><colspec colname="3" colwidth="21PT" align="center" /><colspec colname="4" colwidth="49PT" align="center" /><colspec colname="5" colwidth="21PT" align="center" /><colspec colname="6" colwidth="42PT" align="center" /><tbody valign="top"><row><entry /><entry>01</entry><entry>02</entry><entry>03</entry><entry>04</entry><entry>05</entry><entry>06</entry></row><row><entry /><entry>11</entry><entry>12</entry><entry>13</entry><entry>14</entry><entry>15</entry><entry>16</entry></row><row><entry /><entry>21</entry><entry>22</entry><entry>23</entry><entry>24</entry><entry>25</entry><entry>26</entry></row><row><entry /><entry namest="OFFSET" nameend="6" align="center" rowsep="1" /></row></tbody></tgroup></table></tables>
[0196]<tables id="TABLE-US-00019" num="19"><table frame="none" colsep="0" rowsep="0"><tgroup align="left" colsep="0" rowsep="0" cols="1"><colspec colname="1" colwidth="217PT" align="center" /><thead><row><entry namest="1" nameend="1" align="center">TABLE 18</entry></row></thead><tbody valign="top"><row><entry /></row><row><entry namest="1" nameend="1" align="center" rowsep="1" /></row><row><entry>Example 6x3 Matrix</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="3"><colspec colname="OFFSET" colwidth="14PT" align="left" /><colspec colname="1" colwidth="21PT" align="center" /><colspec colname="2" colwidth="182PT" align="center" /><tbody valign="top"><row><entry /><entry>b</entry></row><row><entry /><entry namest="OFFSET" nameend="2" align="center" rowsep="1" /></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="7"><colspec colname="OFFSET" colwidth="14PT" align="left" /><colspec colname="1" colwidth="21PT" align="center" /><colspec colname="2" colwidth="49PT" align="center" /><colspec colname="3" colwidth="21PT" align="center" /><colspec colname="4" colwidth="49PT" align="center" /><colspec colname="5" colwidth="21PT" align="center" /><colspec colname="6" colwidth="42PT" align="center" /><tbody valign="top"><row><entry /><entry>05</entry><entry>06</entry><entry>07</entry><entry>08</entry><entry>09</entry><entry>10</entry></row><row><entry /><entry>15</entry><entry>16</entry><entry>17</entry><entry>18</entry><entry>19</entry><entry>20</entry></row><row><entry /><entry>25</entry><entry>26</entry><entry>27</entry><entry>28</entry><entry>29</entry><entry>30</entry></row><row><entry /><entry namest="OFFSET" nameend="6" align="center" rowsep="1" /></row></tbody></tgroup></table></tables>
[0197] Column a of Table 17 and column b of Table 18 represent overlapped data. The overlapping area is a one-dimensional overlap area. The overlap area is treated as an overlap but without the current need to transfer data. The resulting processing sequence of processing node 1 and processing node 2 in the HC is shown in Table 19. <tables id="TABLE-US-00020" num="20"><table frame="none" colsep="0" rowsep="0"><tgroup align="left" colsep="0" rowsep="0" cols="1"><colspec colname="1" colwidth="217PT" align="center" /><thead><row><entry namest="1" nameend="1" align="center">TABLE 19</entry></row></thead><tbody valign="top"><row><entry /></row><row><entry namest="1" nameend="1" align="center" rowsep="1" /></row><row><entry>Example Processing Sequence of the HC</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="3"><colspec colname="OFFSET" colwidth="28PT" align="left" /><colspec colname="1" colwidth="91PT" align="left" /><colspec colname="2" colwidth="98PT" align="left" /><tbody valign="top"><row><entry /><entry>Processing Node 1</entry><entry>Processing Node 2</entry></row><row><entry /><entry namest="OFFSET" nameend="2" align="center" rowsep="1" /></row><row><entry /><entry>Process Data Item - 05</entry><entry>Process Data Item - 06</entry></row><row><entry /><entry>Process Data Item - 06</entry><entry>Process Data Item - 05</entry></row><row><entry /><entry>Process Data Item - 15</entry><entry>Process Data Item - 16</entry></row><row><entry /><entry>Process Data Item - 16</entry><entry>Process Data Item - 16</entry></row><row><entry /><entry>Process Data Item - 25</entry><entry>Process Data Item - 26</entry></row><row><entry /><entry>Process Data Item - 26</entry><entry>Process Data Item - 25</entry></row><row><entry /><entry namest="OFFSET" nameend="2" align="center" rowsep="1" /></row></tbody></tgroup></table></tables>
[0198] By way of comparison, consider the number of parallel activities. Where there is an overlap area, processing node 1 and processing node 2 process data in parallel. Thus, six time units are consumed. Where the data is not overlapped, the processing is serialized whenever data is transferred between the processing nodes, thus twelve time units are consumed (assuming that both the data transfer and the data processing consume one time unit). This effect is exacerbated when multiple overlapping dimensions exist. For example, Table 20 shows a 9×15 element matrix divided between nine processing nodes, with overlap between the data items for each node. <tables id="TABLE-US-00021" num="21"><table frame="none" colsep="0" rowsep="0" pgwide="1"><tgroup align="left" colsep="0" rowsep="0" cols="1"><colspec colname="1" colwidth="399PT" align="center" /><thead><row><entry namest="1" nameend="1" align="center">TABLE 20</entry></row><row><entry /></row><row><entry namest="1" nameend="1" align="center" rowsep="1" /></row><row><entry>Example 9x15 Matrix Over Nine Processing Nodes, with Overlap</entry></row><row><entry namest="1" nameend="1" align="center" rowsep="1" /></row></thead><tbody valign="top"><row><entry /></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="19"><colspec colname="1" colwidth="21PT" align="char" char="." /><colspec colname="2" colwidth="21PT" align="char" char="." /><colspec colname="3" colwidth="21PT" align="char" char="." /><colspec colname="4" colwidth="21PT" align="char" char="." /><colspec colname="5" colwidth="21PT" align="char" char="." /><colspec colname="6" colwidth="21PT" align="char" char="." /><colspec colname="7" colwidth="21PT" align="char" char="." /><colspec colname="8" colwidth="21PT" align="char" char="." /><colspec colname="9" colwidth="21PT" align="char" char="." /><colspec colname="10" colwidth="21PT" align="char" char="." /><colspec colname="11" colwidth="21PT" align="char" char="." /><colspec colname="12" colwidth="21PT" align="char" char="." /><colspec colname="13" colwidth="21PT" align="char" char="." /><colspec colname="14" colwidth="21PT" align="char" char="." /><colspec colname="15" colwidth="21PT" align="char" char="." /><colspec colname="16" colwidth="21PT" align="char" char="." /><colspec colname="17" colwidth="21PT" align="char" char="." /><colspec colname="18" colwidth="21PT" align="char" char="." /><colspec colname="19" colwidth="21PT" align="char" char="." /><tbody valign="top"><row><entry>001</entry><entry>002</entry><entry>003</entry><entry>004</entry><entry>005</entry><entry>006</entry><entry>005</entry><entry>006</entry><entry>007</entry><entry>008</entry><entry>009</entry><entry>010</entry><entry>011</entry><entry>010</entry><entry>011</entry><entry>012</entry><entry>013</entry><entry>014</entry><entry>015</entry></row><row><entry>016</entry><entry>017</entry><entry>018</entry><entry>019</entry><entry>020</entry><entry>021</entry><entry>020</entry><entry>021</entry><entry>022</entry><entry>023</entry><entry>024</entry><entry>025</entry><entry>026</entry><entry>025</entry><entry>026</entry><entry>027</entry><entry>028</entry><entry>029</entry><entry>030</entry></row><row><entry>031</entry><entry>032</entry><entry>033</entry><entry>034</entry><entry>035</entry><entry>036</entry><entry>035</entry><entry>036</entry><entry>037</entry><entry>038</entry><entry>039</entry><entry>040</entry><entry>041</entry><entry>040</entry><entry>041</entry><entry>042</entry><entry>043</entry><entry>044</entry><entry>045</entry></row><row><entry>046</entry><entry>047</entry><entry>048</entry><entry>049</entry><entry>050</entry><entry /><entry /><entry>051</entry><entry>052</entry><entry>053</entry><entry>054</entry><entry>055</entry><entry /><entry /><entry>056</entry><entry>057</entry><entry>058</entry><entry>059</entry><entry>060</entry></row><row><entry>031</entry><entry>032</entry><entry>033</entry><entry>034</entry><entry>035</entry><entry /><entry /><entry>036</entry><entry>037</entry><entry>038</entry><entry>039</entry><entry>040</entry><entry /><entry /><entry>041</entry><entry>042</entry><entry>043</entry><entry>044</entry><entry>045</entry></row><row><entry>046</entry><entry>047</entry><entry>048</entry><entry>049</entry><entry>050</entry><entry>051</entry><entry>050</entry><entry>051</entry><entry>052</entry><entry>053</entry><entry>054</entry><entry>055</entry><entry>056</entry><entry>055</entry><entry>056</entry><entry>057</entry><entry>058</entry><entry>059</entry><entry>060</entry></row><row><entry>061</entry><entry>062</entry><entry>063</entry><entry>064</entry><entry>065</entry><entry>066</entry><entry>065</entry><entry>066</entry><entry>067</entry><entry>068</entry><entry>069</entry><entry>070</entry><entry>071</entry><entry>070</entry><entry>071</entry><entry>072</entry><entry>073</entry><entry>074</entry><entry>075</entry></row><row><entry>076</entry><entry>077</entry><entry>078</entry><entry>079</entry><entry>080</entry><entry>081</entry><entry>080</entry><entry>081</entry><entry>082</entry><entry>083</entry><entry>084</entry><entry>085</entry><entry>086</entry><entry>085</entry><entry>086</entry><entry>087</entry><entry>088</entry><entry>089</entry><entry>090</entry></row><row><entry>091</entry><entry>092</entry><entry>093</entry><entry>094</entry><entry>095</entry><entry /><entry /><entry>096</entry><entry>097</entry><entry>098</entry><entry>099</entry><entry>100</entry><entry /><entry /><entry>101</entry><entry>102</entry><entry>103</entry><entry>104</entry><entry>105</entry></row><row><entry>07</entry><entry>07</entry><entry>07</entry><entry>07</entry><entry>08</entry><entry /><entry /><entry>08</entry><entry>08</entry><entry>08</entry><entry>08</entry><entry>08</entry><entry /><entry /><entry>08</entry><entry>08</entry><entry>08</entry><entry>08</entry><entry>09</entry></row><row><entry>6</entry><entry>7</entry><entry>8</entry><entry>9</entry><entry>0</entry><entry /><entry /><entry>1</entry><entry>2</entry><entry>3</entry><entry>4</entry><entry>5</entry><entry /><entry /><entry>6</entry><entry>7</entry><entry>8</entry><entry>9</entry><entry>0</entry></row><row><entry>09</entry><entry>09</entry><entry>09</entry><entry>09</entry><entry>09</entry><entry>09</entry><entry>09</entry><entry>09</entry><entry>09</entry><entry>09</entry><entry>09</entry><entry>10</entry><entry>10</entry><entry>10</entry><entry>10</entry><entry>10</entry><entry>10</entry><entry>10</entry><entry>10</entry></row><row><entry>1</entry><entry>2</entry><entry>3</entry><entry>4</entry><entry>5</entry><entry>6</entry><entry>5</entry><entry>6</entry><entry>7</entry><entry>8</entry><entry>9</entry><entry>0</entry><entry>1</entry><entry>0</entry><entry>1</entry><entry>2</entry><entry>3</entry><entry>4</entry><entry>5</entry></row><row><entry>10</entry><entry>10</entry><entry>10</entry><entry>10</entry><entry>11</entry><entry>11</entry><entry>11</entry><entry>11</entry><entry>11</entry><entry>11</entry><entry>11</entry><entry>11</entry><entry>11</entry><entry>11</entry><entry>11</entry><entry>11</entry><entry>11</entry><entry>11</entry><entry>12</entry></row><row><entry>6</entry><entry>7</entry><entry>8</entry><entry>9</entry><entry>0</entry><entry>1</entry><entry>0</entry><entry>1</entry><entry>2</entry><entry>3</entry><entry>4</entry><entry>5</entry><entry>6</entry><entry>5</entry><entry>6</entry><entry>7</entry><entry>8</entry><entry>9</entry><entry>0</entry></row><row><entry>12</entry><entry>12</entry><entry>12</entry><entry>12</entry><entry>12</entry><entry>12</entry><entry>12</entry><entry>12</entry><entry>12</entry><entry>12</entry><entry>12</entry><entry>13</entry><entry>13</entry><entry>13</entry><entry>13</entry><entry>13</entry><entry>13</entry><entry>13</entry><entry>13</entry></row><row><entry>1</entry><entry>2</entry><entry>3</entry><entry>4</entry><entry>5</entry><entry>6</entry><entry>5</entry><entry>6</entry><entry>7</entry><entry>8</entry><entry>9</entry><entry>0</entry><entry>1</entry><entry>0</entry><entry>1</entry><entry>2</entry><entry>3</entry><entry>4</entry><entry>5</entry></row><row><entry namest="1" nameend="19" align="center" rowsep="1" /></row></tbody></tgroup></table></tables>
[0199] As can be seen with the underlined data entries in Table 20, the processing node with that dataset has all sides of a 2-dimensional matrix containing an overlap area (thus representing a 2-dimensional overlap.) This may be extended from 1- to N-dimensions analogously. The overlap methodology decreases the overhead of cluster transfers and allows for efficient use of separate processors in problem solving, as compared to continuously sharing data between nodes.
[0200] Processing Node Dropout Detection and Replacement
[0201] As previously discussed, the HC may be configured to detect non-functioning processing nodes and to reassign associated work for the non-functioning processing node to another processing node without human intervention. The clusters of the prior art have enormous difficulty in both detecting and ameliorating failed nodes; this difficulty is a function of how computer nodes are assigned to a problem. Typically, in the prior art (e.g., as shown in FIG. 1-FIG. 4), a remote host or some other computer external to the cluster compiles the application code that is to run on the cluster. A compile time parallel processing communication tool is used, such as MPI or PVM, known in the art, whereby the communication relationships between processing nodes is established at compile time. When the processing node relationships are formed at compile time, it is difficult to establish run time re-allocation of a parallel problem to other functioning processing nodes. If in the course of processing a job a node loses communication with other nodes in the cluster, the non-communication condition is not corrected without either specialized hardware and/or human intervention.
[0202] In one embodiment of the HC, on the other hand, node communications are determined at run time rather than at compile time. Further, the geometric nature of the HC fixes the run-time node communication relationships for the duration of the job. This makes it possible to both detect processing node communication failures and to reallocate problem sets to other processing nodes, thereby correcting for node failures at run time.
[0203]FIG. 21 illustrates one HC <b>100</b>E with seven processing nodes <b>112</b>, <b>114</b>, <b>116</b>, <b>118</b>, <b>120</b>, <b>122</b> and <b>124</b> configured in the cascade, and an eighth unallocated processing node <b>126</b>. Home node <b>110</b>E communicates with processing nodes <b>112</b>, <b>114</b>, <b>116</b>, <b>118</b>, <b>120</b>, <b>122</b> and <b>124</b> during normal processing. If however home node <b>110</b>E fails to contact processing node <b>112</b>, for example after an appropriate number of retries, it places the address of processing node <b>112</b> onto the “not communicating” list and immediately allocates another processing node to take its place in HC <b>100</b>E. FIG. 22 illustrates the connectivity of cluster HC <b>100</b>E after reconfiguration. In this example, home node <b>110</b>E detects the failure of communication path <b>250</b>. Home node <b>110</b>E then communicates with processing node <b>126</b> via data path <b>252</b>, to inform processing node <b>126</b> of its position in HC <b>100</b>E. Processing node <b>126</b> then establishes communication with its down level nodes (processing node <b>114</b> and processing node <b>118</b>) via communication paths <b>254</b> and <b>256</b>, respectively. Once completed, HC <b>100</b>E is repaired and processing resumes using processing node <b>126</b> in place of processing node <b>112</b>.
[0204] Non-communication is detectable because the communication protocol (e.g., TCP/IP) returns an error code when one networked node attempts to communicate with another node on the same network, but cannot. By way of example, FIG. 23 shows physical connections between the nodes of FIG. 21 and a network switch <b>260</b>. If processing node <b>112</b> is no longer communicating with network switch <b>260</b>, home node <b>110</b>E selects the next available node, in this example node <b>126</b>, as replacement. The physical connection topology of HC <b>100</b>E allows any node to be positioned anywhere in the cascade without problem or overhead.
[0205] In another embodiment, FIG. 24 illustrates a HC <b>100</b>F in a state where parallel processing through processing nodes <b>112</b>, <b>114</b>, <b>116</b>, <b>118</b>, <b>120</b>, <b>122</b> and <b>124</b> has completed and HC <b>100</b>F has configured to return agglomerated results to home node <b>110</b>F. FIG. 25 illustrates an autonomous error recovery process resulting from the failure of processing node <b>112</b>. When processing node <b>114</b> tries to communicate with processing node <b>112</b>, via communication path <b>280</b>, it receives an error. Processing node <b>114</b> then sends a message directly to home node <b>110</b>F to inform home node <b>110</b>F that processing node <b>112</b> is not communicating, via communication path <b>282</b>. Home node <b>110</b>F uses communication path <b>284</b> to verify that processing node <b>110</b> is no longer communicating. Home node <b>110</b>F then allocates the next available processing node, in this case processing node <b>126</b>, using communication path <b>286</b>. All information required to allow processing node <b>126</b> to take the place of processing node <b>112</b> is transmitted to processing node <b>126</b>. Home node <b>110</b>F then informs processing node <b>114</b> and processing node <b>118</b> that a new up-level node (i.e., node <b>126</b>) exists via communication paths <b>288</b> and <b>290</b>, respectively. Processing node <b>114</b> and processing node <b>118</b> then send respective results to processing node <b>126</b> via communication paths <b>292</b> and <b>294</b>, respectively. Processing node <b>126</b> then sends its agglomerated results upstream to home node <b>110</b>F, via communication path <b>296</b>.
[0206] In another example, a HC <b>100</b>G is configured similarly to HC <b>100</b>E of FIG. 21. FIG. 26 illustrates the error recovery sequence that occurs when processing node <b>112</b> attempts to cascade the algorithm processing request to processing node <b>114</b> and communication path <b>300</b> fails. Processing node <b>112</b> informs home node <b>110</b>G of the failure via communication path <b>302</b>. Home node <b>110</b>G then selects the next available processing node, processing node <b>126</b> in this example, and informs processing node <b>112</b> of the identity of the new node, via communication path <b>304</b>. Processing node <b>112</b> then informs processing node <b>126</b> of its new position in HC <b>100</b>G via communication path <b>306</b>, to communicate the algorithm processing request originally destined for processing node <b>114</b>. Processing node <b>126</b> then continues the cascade to processing node <b>116</b> via communication path <b>308</b>.
[0207] In another error recovery example, a HC <b>100</b>H, FIG. 27, recovers during agglomeration of results. FIG. 27 shows the error recovery sequence of HC <b>100</b>H that occurs when communication path <b>320</b> between processing node <b>116</b> and processing node <b>114</b> fails. Processing node <b>116</b> informs home node <b>110</b>H of the failure via communication path <b>322</b>. Home node <b>110</b>H selects the next available processing node, in this example processing node <b>126</b>, via communication path <b>324</b>. Home node <b>110</b>H then informs processing node <b>116</b> of the identity of the new node via communication path <b>326</b>. Processing node <b>116</b> then sends its results to processing node <b>126</b> via communication path <b>328</b>. Processing node <b>126</b> then sends its agglomerated results to processing node <b>112</b> via communication path <b>330</b>.
[0208] In the event there are no spare processing nodes, the HC may suspend the current processing cascades and recasts the algorithm processing request at the next lower cascade level, where additional, free processing nodes may be used to restart processing. In this example, a HC <b>100</b>E is configured as in FIG. 21, except that it is assumed that node <b>126</b> is unavailable. FIG. 28 illustrates the error recovery sequence that occurs when processing node <b>120</b> fails to communicate with processing node <b>122</b> via communication path <b>340</b>. Processing node <b>120</b> informs home node <b>110</b>E of the failure via communication path <b>342</b>. Home node <b>110</b>E determines that there are no spare processing nodes, and sends a command, via communication path <b>344</b>, to stop processing of the current algorithm processing request on processing nodes <b>112</b>, <b>120</b>, <b>124</b>. Processing node <b>112</b> stops algorithm processing on processing node <b>114</b>, which in turn stops algorithm processing on processing node <b>116</b>. Home node <b>110</b>E adds processing node <b>122</b> to the failed list. FIG. 29 illustrates home node <b>110</b>E retransmitting the current algorithm processing request to a smaller group of nodes via communication paths <b>360</b>, leaving three unallocated processing nodes, processing node <b>118</b>, processing node <b>120</b> and processing node <b>124</b>. The algorithm processing request is thus completed among replacement nodes without human intervention and may be used at any inner-cluster location.
[0209] In one embodiment, each node in a HC may have the same computational power. Since the connection between a remote host and the home node does not stay open during job processing, it is possible to switch out and replace a failed primary home node without the need to inform the remote host. FIG. 30 shows information received from a remote host <b>329</b> being shared with spare nodes (node <b>331</b>) in a HC <b>100</b>I via additional communication paths <b>380</b>. FIG. 31 further illustrates that if primary home node <b>110</b>I fails, (e.g., communication path <b>400</b> fails), then the next available spare node is used to replace failed home node <b>110</b>I; in this example, the spare processing node <b>331</b> is reconfigured by a RPC to become the new home node.
[0210] HC <b>100</b>I need not utilize a client-server model, known to those skilled in the art, and thus the connection to the remote host may be rebuilt.
[0211] Hidden Function API for Blind Function Calling
[0212] A HCSA may also be constructed with interface methodology that avoids third-party access to proprietary algorithms, overcoming one problem of the prior art. FIG. 34 illustrates two interfaces between the application code <b>460</b> and a HCAS. The first interface is function request interface <b>476</b>; the second interface is the data request interface <b>466</b>. Through the function request interface <b>476</b>, application code <b>460</b> may invoke a function in computationally intensive algorithm library <b>99</b>, FIG. 5, with a parameter list without having a complex application program interface (“API”). Prior to calling message creation software using the API, a template number <b>468</b>, a function number <b>470</b>, a data definition of buffer size <b>472</b>, and a buffer <b>474</b> are created. Template number <b>468</b> defines the data template for the function (in computationally intensive algorithm library <b>99</b>) selected by function number <b>470</b>. This information is stored in a parameter definition file <b>482</b>, FIG. 35. Parameter definition file <b>482</b> is read into message creation software <b>480</b> at run time as part of the initialization sequence of application code <b>460</b>. The message creation software is a static dynamic link library (“DLL”) called by application code <b>460</b>. When application code <b>460</b> uses a function (of computationally intensive algorithm library <b>99</b>), the requirement is triggered and template number <b>468</b> and function number <b>470</b> are sent to the message creation software along with buffer size <b>472</b> and buffer <b>474</b>, containing function parameters. Template number <b>468</b> and function number <b>470</b> are used to find the definition of the parameters, contained in buffer <b>474</b>, and of the output data structure definition of application code <b>460</b>. Message creation software <b>480</b> uses all of this data to create an output message definition <b>484</b> that is sent to the home node of the HCAS. The home node uses the template number <b>468</b> and parameter data contained in buffer <b>474</b>, which defines the data types, number of dimensions, and the size of each dimension, to determine the number of processing nodes to use.
[0213] More particularly, FIG. 36 shows one home node <b>110</b>J initiating distribution of the algorithm processing request to processing node <b>112</b> via output message <b>502</b>. Only one processing node <b>112</b> is shown for purposes of illustration. Processing node <b>112</b> then waits for data as required by the invoked function. Home node <b>110</b>J then requests and receives the data from application code <b>460</b> via the data request interface <b>466</b>. Home node <b>110</b>J is shown receiving data from application <b>504</b> in FIG. 36. Home node <b>110</b>J then broadcasts the data to processing node <b>112</b> using a data broadcast message <b>500</b>. Control software <b>94</b>, FIG. 5, within processing node <b>112</b>, then makes a function call to the selected function within the computationally intensive algorithm library <b>99</b> using parameters translated from the output message definition <b>484</b> contained in output message <b>502</b>. The invoked function processes the data, returning the results. The agglomeration process then returns the results to home node <b>110</b>, which sends the results back to application code <b>460</b> as data to application <b>506</b>. Home node <b>110</b>J then sends the results back to application code <b>460</b>.
[0214] Some parameter types passed to the message creation software <b>480</b> by application code <b>460</b> may contain special values. One such value is a storage type for pointers to specific data areas used for input parameters to the requested function, where the pointer value can only be resolved within the processing node control software at run time. This allows flexibility and efficiency in parameter passing.
[0215] In summary, the function request interface <b>476</b> and data request interface <b>466</b> of FIG. 34 allow functions to be “parallelized” at the data level without application code <b>460</b> requiring access to sensitive functions. It also provides for a simple yet highly extensible API to application programs.
[0216] Complex Algorithms on an HCAS
[0217] A complex algorithm is defined as an algorithm that contains one or more branching statements. A branching statement is a statement that compares one or more variables and constants, and then selects an execution path based upon the comparison. For example: <tables id="TABLE-US-00022" num="22"><table frame="none" colsep="0" rowsep="0"><tgroup align="left" colsep="0" rowsep="0" cols="2"><colspec colname="OFFSET" colwidth="63PT" align="left" /><colspec colname="1" colwidth="154PT" align="left" /><thead><row><entry /><entry /></row><row><entry /><entry namest="OFFSET" nameend="1" align="center" rowsep="1" /></row></thead><tbody valign="top"><row><entry /><entry>IF (variable_a > 10) then</entry></row><row><entry /><entry>{</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="2"><colspec colname="OFFSET" colwidth="77PT" align="left" /><colspec colname="1" colwidth="140PT" align="left" /><tbody valign="top"><row><entry /><entry>Execution-Path-A</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="2"><colspec colname="OFFSET" colwidth="63PT" align="left" /><colspec colname="1" colwidth="154PT" align="left" /><tbody valign="top"><row><entry /><entry>}</entry></row><row><entry /><entry>else</entry></row><row><entry /><entry>{</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="2"><colspec colname="OFFSET" colwidth="77PT" align="left" /><colspec colname="1" colwidth="140PT" align="left" /><tbody valign="top"><row><entry /><entry>Execution-Path-B</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="2"><colspec colname="OFFSET" colwidth="63PT" align="left" /><colspec colname="1" colwidth="154PT" align="left" /><tbody valign="top"><row><entry /><entry>}</entry></row><row><entry /><entry namest="OFFSET" nameend="1" align="center" rowsep="1" /></row></tbody></tgroup></table></tables>
[0218] In the above branching statement one of two execution paths is selected based upon whether or not the contents of variable a is greater than 10. The execution paths may contain further branching statements and computational algorithms.
[0219] Computational algorithms can always be represented as one or more subroutines or library function calls. In the HCAS concept, code that can be represented as either a subroutine or a library function can be installed into algorithm library <b>99</b> on an HCAS. Thus the computational algorithm parts of a branching statement can be processed on a HCAS.
[0220] Each branching statement represents a serial activity. All parallel processes executing a complex algorithm are required to make the same decision at the same point in their processing. There are primarily two ways to accomplish this. A first method is for each process executing the complex algorithm (e.g., a processing node in an HCAS) to send messages to every other process executing the complex algorithm, such that each process has sufficient information to make the same decision. This first method maximizes the amount of cross-communication required and thus is unacceptable as a solution. A second method uses a central location to evaluate the branching statement. For example, this central location may be on the remote host or on a home node in an HCAS. The conditional variable data is the only information that has to be transmitted from each process executing the complex algorithm to the central location, thus keeping data transfers to a minimum.
[0221] A branching statement may be represented by a conditional function with the following attributes: <tables id="TABLE-US-00023" num="23"><table frame="none" colsep="0" rowsep="0"><tgroup align="left" colsep="0" rowsep="0" cols="1"><colspec colname="1" colwidth="217PT" align="left" /><thead><row><entry /></row><row><entry namest="1" nameend="1" align="center" rowsep="1" /></row></thead><tbody valign="top"><row><entry>FUNCTION_NAME ((Variable|constant)<sub>1</sub> <comparison attribute></entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="2"><colspec colname="OFFSET" colwidth="42PT" align="left" /><colspec colname="1" colwidth="175PT" align="left" /><tbody valign="top"><row><entry /><entry>(Variable|constant)<sub>2</sub>)</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="1"><colspec colname="1" colwidth="217PT" align="left" /><tbody valign="top"><row><entry>{true path}</entry></row><row><entry>else</entry></row><row><entry>{false path}</entry></row><row><entry>where:</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="3"><colspec colname="OFFSET" colwidth="14PT" align="left" /><colspec colname="1" colwidth="77PT" align="left" /><colspec colname="2" colwidth="126PT" align="left" /><tbody valign="top"><row><entry /><entry>FUNCTION_NAME</entry><entry>= the condition type, for example:</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="2"><colspec colname="OFFSET" colwidth="105PT" align="left" /><colspec colname="1" colwidth="112PT" align="left" /><tbody valign="top"><row><entry /><entry>IF, While, Until, etc.</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="3"><colspec colname="OFFSET" colwidth="14PT" align="left" /><colspec colname="1" colwidth="77PT" align="left" /><colspec colname="2" colwidth="126PT" align="left" /><tbody valign="top"><row><entry /><entry>Variable|constant</entry><entry>= either the name of a variable or a</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="2"><colspec colname="OFFSET" colwidth="91PT" align="left" /><colspec colname="1" colwidth="126PT" align="left" /><tbody valign="top"><row><entry /><entry>constant value</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="3"><colspec colname="OFFSET" colwidth="14PT" align="left" /><colspec colname="1" colwidth="77PT" align="left" /><colspec colname="2" colwidth="126PT" align="left" /><tbody valign="top"><row><entry /><entry>Comparison attribute</entry><entry>= a logical or mathematical statement,</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="2"><colspec colname="OFFSET" colwidth="91PT" align="left" /><colspec colname="1" colwidth="126PT" align="left" /><tbody valign="top"><row><entry /><entry>for example:</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="2"><colspec colname="OFFSET" colwidth="105PT" align="left" /><colspec colname="1" colwidth="112PT" align="left" /><tbody valign="top"><row><entry /><entry>AND, OR, <, >, =, NOR, NAND, etc</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="2"><colspec colname="OFFSET" colwidth="91PT" align="left" /><colspec colname="1" colwidth="126PT" align="left" /><tbody valign="top"><row><entry /><entry>Which compares (Variable|constant)<sub>1</sub> with</entry></row><row><entry /><entry>(Variable|constant)<sub>2</sub></entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="2"><colspec colname="OFFSET" colwidth="14PT" align="left" /><colspec colname="1" colwidth="203PT" align="left" /><tbody valign="top"><row><entry /><entry>Note: If the comparison is true then the true path is taken</entry></row><row><entry /><entry>otherwise the false path is taken.</entry></row><row><entry /><entry namest="OFFSET" nameend="1" align="center" rowsep="1" /></row></tbody></tgroup></table></tables>
[0222] Thus, each branching statement in a complex algorithm becomes a function call. The only difference between a conditional function and all other HCAS functions is that the conditional data is sent to a central location for evaluation at each branching statement in the complex algorithm.
[0223]FIG. 63 is a flow chart illustrating one process <b>1180</b> as an example of a complex algorithm containing a branching statement. Process <b>1180</b> starts at step <b>1182</b>, and continues with step <b>1184</b>. Step <b>1184</b> represents a section of code, section A, that computes a coefficient using a <b>2</b><i>d </i>correlation function on image data. Step <b>1186</b> represents a branching statement that uses the coefficient computed in step <b>1184</b>. If the coefficient is greater than 0.9, process <b>1180</b> continues with step <b>1190</b>; otherwise process <b>1180</b> continues at step <b>1188</b>. Step <b>1188</b> represents a second section of code, section B, that performs additional image processing in this example. After computing section B, process <b>1180</b> continues with step <b>1190</b>. Step <b>1190</b> gets the results of the complex algorithm computation and process <b>1180</b> terminates at step <b>1192</b>. As can be seen in this example, the complex algorithm may be broken down into computable sections, section A <b>1184</b> and section B <b>1188</b>, and a conditional function, as in step <b>1186</b>.
[0224]FIG. 64 illustrates one embodiment for implementing complex algorithms on an HCAS. FIG. 64 shows three flow charts for three interacting processes; host process <b>1200</b>, control process <b>1300</b> and computing process <b>1400</b>. Host process <b>1200</b> represents a process running on a remote host computer that initiates a complex algorithm processing request and receives the results. Control process <b>1300</b> represents a process that arbitrates conditional branching statements in the complex algorithm. Control process <b>1300</b> may run on a remote host computer or within an HCAS. Computing process <b>1400</b> represents a process for computing complex algorithm sections on an HCAS.
[0225] Host process <b>1200</b> starts at step <b>1202</b> and continues with step <b>1204</b>. Step <b>1204</b> sends a complex algorithm processing request control process <b>1300</b>. Host process <b>1200</b> continues at step <b>1206</b> where it waits to receive the computed results.
[0226] Control process <b>1300</b> starts at step <b>1302</b> and continues with step <b>1304</b>. Step <b>1304</b> receives the complex algorithm processing request from remote process <b>1200</b> and continues with step <b>1306</b>. Step <b>1306</b> sends the complex algorithm processing sections to computing process <b>1400</b>, indicating the first section to be computed. Control process <b>1300</b> continues with step <b>1308</b>. Step <b>1308</b> waits for computing process <b>1400</b> to return results after completing a section of the complex algorithm.
[0227] Computing process <b>1400</b> starts at step <b>1402</b> and continues with step <b>1404</b>. Step <b>1404</b> receives the complex algorithm processing sections from control process <b>1300</b>. Computing process <b>1400</b> continues with step <b>1406</b> where the indicated section of the complex algorithm processing request is computed. Computing process <b>1400</b> continues with step <b>1408</b>. Step <b>1408</b> agglomerates the results of the computed section and sends them to control process <b>1300</b>. Computing process <b>1400</b> continues with step <b>1410</b> where it waits for direction from control process <b>1300</b>.
[0228] Control process <b>1300</b> receives the results of the computed section from computing process <b>1400</b> in step <b>1308</b> and continues with step <b>1310</b>. If no further computing of the complex algorithm processing request is required, control process <b>1300</b> continues with step <b>1314</b>; otherwise, step <b>1310</b> evaluates the respective conditional branch in the complex algorithm processing request using results returned from computing process <b>1400</b> and continues with step <b>1312</b>.
[0229] Step <b>1312</b> initiates the computation of the next section of complex algorithm processing request by sending a direction messages to computing process <b>1400</b>. Control process <b>1300</b> then continues with step <b>1308</b> where it waits for the computation of the initiated section to be completed by computing process <b>1400</b>.
[0230] Step <b>1410</b> of computing process <b>1400</b> receives direction from control process <b>1300</b> indicating if, and which, section to compute next. If there is no section to compute next, computing process <b>1400</b> terminates at step <b>1414</b>; otherwise computing process <b>1400</b> continues with step <b>1406</b> where the next section is computed.
[0231] In step <b>1314</b> of control process <b>1300</b>, a direction message is sent to computing process <b>1400</b> indicating that there are no further sections for computing. Control process <b>1300</b> continues with step <b>1316</b> where it sends the final results received from computing process <b>1400</b> in step <b>1308</b> to host process <b>1200</b>. Control process <b>1300</b> then terminates at step <b>1318</b>.
[0232] Host process <b>1200</b> receives the results from control process <b>1300</b> in step <b>1206</b>. Host process <b>1200</b> then terminates at step <b>1208</b>.
[0233] As can be seen in the above example, the decomposition of complex algorithm into computable sections and conditional functions, data is only moved as necessary. By evaluating the branching statement in a central location the inherent serial nature of the branching statement is maintained, and a HCAS is able to handle complex algorithms efficiently.
[0234] Algorithm Development Toolkit
[0235] This section (Algorithm Development Toolkit) describes select processes for implementing algorithms within a HC. Such processes may be automated as a matter of design choice.
[0236]FIG. 37 is a block schematic illustrating how an algorithm development toolkit <b>522</b> is used to add a new algorithm <b>520</b> to algorithm library <b>99</b> such that new algorithm <b>520</b> will operate on HCAS <b>80</b>A. Algorithm development toolkit <b>522</b> includes a set of routines that may be added to new algorithm <b>520</b> to enable it to operate in parallel on processing nodes <b>112</b>, <b>114</b>, <b>116</b> and <b>118</b> of HCAS <b>80</b> with a minimum amount of development work. The added routines make new algorithm <b>520</b> aware of the processing node on which it is running, and the data structure that it has to process.
[0237] Home node <b>110</b>K, processing node <b>112</b>, processing node <b>114</b>, processing node <b>116</b>, and processing node <b>118</b> contain data template definitions <b>528</b> that define the data and parameters for functions in the computationally intensive algorithm library <b>99</b>. All processing nodes in HCAS <b>80</b>A contain identical copies of the data template definitions <b>528</b> and computationally intensive algorithm library <b>99</b>. Algorithm development toolkit <b>522</b> facilitates the addition of new algorithm <b>520</b> to data definitions <b>528</b> and computationally intensive algorithm library <b>99</b>, via data path <b>524</b>, and to a parallel interface library <b>530</b> used by the application running on remote host <b>82</b>, via data path <b>526</b>.
[0238]FIG. 38 is a flow chart illustrating one process <b>550</b> of new algorithm <b>520</b> as augmented by routines from algorithm development toolkit <b>522</b>. Process <b>550</b> starts at step <b>552</b> and continues with step <b>554</b>.
[0239] Step <b>554</b> is a function call to the mesh tool function, described in FIG. 39, which extracts the input and output data descriptions for the new algorithm. Process <b>550</b> continues with step <b>556</b>.
[0240] Step <b>556</b> is a function call to acquire input data function described in FIG. 51, which acquires the data for the new algorithm. Process <b>550</b> continues with step <b>558</b>.
[0241] Step <b>558</b> is a function call to compute the results for a processing node, and is the invocation of the new algorithm on the data set acquired in step <b>556</b> for the processing node on which this function is to run. Process <b>550</b> continues with step <b>560</b>.
[0242] Step <b>560</b> is a function call to an agglomerate result function described in FIG. 54, which receives results from down-stream processing nodes during the agglomeration phase of the HC. Process <b>550</b> continues with step <b>562</b>.
[0243] Step <b>562</b> is a function call to a return results function described in FIG. 58, which sends the local and agglomerated result to the correct processing node, or home node as necessary. Process <b>550</b> terminates at step <b>564</b>.
[0244]FIG. 39 is a flow chart illustrating the mesh tool sub-process <b>570</b> as invoked in step <b>554</b> of process <b>550</b> in FIG. 38. Mesh tool sub-process <b>570</b> extracts the input and output data descriptions from the algorithm request message and computes the mesh size for local computation, using the processing node position on which it runs. The mesh is a description of the input and output data to be used, and defines how much of the result is computed by this processing node and which input values are used to compute the result. The mesh is dependant on the type of algorithm, the position of the processing node in the HC, and the total number of processing nodes in the HC. Sub-process <b>570</b> begins at step <b>572</b> and continues with step <b>574</b>.
[0245] Step <b>574</b> is a function call to a find an input data size function described in FIG. 40, which extracts the input data size from the algorithm processing request. Sub-process <b>570</b> continues with step <b>576</b>.
[0246] Step <b>576</b> is a function call to a compute output data size sub-process described in FIG. 41, which calculates the output data set size based on the input data set size and the computation type of the new algorithm. Sub-process <b>570</b> continues with step <b>578</b>.
[0247] Step <b>578</b> is a function call to a compute mesh parameters sub-process described in FIG. 45, which determines the mesh parameters that define which part of the results are to be computed by this processing node. Sub-process <b>570</b> terminates at step <b>580</b>, returning to where it was invoked.
[0248]FIG. 40 is a flow chart describing one find input data size sub-process <b>590</b> that starts at step <b>592</b> and continues with step <b>594</b>.
[0249] Step <b>594</b> is a decision. If the new algorithm is a series expansion, sub-process <b>590</b> continues with step <b>596</b>; otherwise sub-process <b>590</b> continues with step <b>598</b>.
[0250] Step <b>596</b> sets the data input rows and data input size to the number of terms in the series. Sub-process <b>590</b> continues with step <b>600</b>.
[0251] Step <b>598</b> sets the first data input rows, columns and element size to match the input image size. Sub-process <b>590</b> continues with step <b>602</b>.
[0252] Step <b>600</b> sets the data input element size to the number of digits provided per element. Sub-process <b>590</b> continues with step <b>604</b>.
[0253] Step <b>602</b> sets the first data input size to the data input rows×columns×element size defined in step <b>598</b>. Sub-process <b>590</b> continues with step <b>606</b>.
[0254] Step <b>604</b> sets the data input columns to one. Sub-process <b>590</b> terminates at step <b>612</b>, returning control to the invoking process.
[0255] Step <b>606</b> is a decision. If there is a second input image, sub-process <b>590</b> continues with step <b>608</b>; otherwise sub-process <b>590</b> terminates at step <b>612</b>, returning control to the invoking process.
[0256] Step <b>608</b> sets the second data input rows, columns and element size to match the second input image. Sub-process <b>590</b> continues with step <b>610</b>.
[0257] Step <b>610</b> sets the second data input size to the data input rows×columns×element size defined in step <b>608</b>. Sub-process <b>590</b> terminates at step <b>612</b>, returning control to the invoking process.
[0258]FIG. 41 is a flowchart illustrating one sub-process <b>620</b> to find output data size as invoked by process <b>550</b> in FIG. 38. Sub-process <b>620</b> starts at step <b>622</b> and continues with step <b>624</b>.
[0259] Step <b>624</b> is a decision. If the new algorithm is a series expansion, sub-process <b>620</b> continues with step <b>628</b>; otherwise sub-process <b>620</b> continues with step <b>626</b>.
[0260] Step <b>626</b> is a decision. If there is a second image, sub-process <b>620</b> continues with step <b>632</b>; otherwise sub-process <b>620</b> continues with step <b>630</b>.
[0261] Step <b>628</b> invokes sub-process compute output series size for a series expansion defined in FIG. 42, which calculates the output data size for the series calculated by the new algorithm. Sub-process <b>620</b> terminates at step <b>634</b>, returning control to the invoking process.
[0262] Step <b>630</b> invokes sub-process compute output data size for a single image defined in FIG. 43, which calculates the output data size for a single image. Sub-process <b>620</b> terminates at step <b>634</b>, returning control to the invoking process.
[0263] Step <b>632</b> invokes sub-process compute output data size for two images defined in FIG. 44, which calculates the output data size for two images. Sub-process <b>620</b> terminates at step <b>634</b>, returning control to the invoking process.
[0264]FIG. 42 illustrates one sub-process <b>640</b> for computing the output data size for the series expansion used in the new algorithm. Sub-process <b>640</b> starts at step <b>642</b> and continues with step <b>644</b>.
[0265] Step <b>644</b> sets the data output rows to the number of terms in the series. Sub-process <b>640</b> continues with step <b>646</b>.
[0266] Step <b>646</b> sets the data output element size to one. Sub-process <b>640</b> continues with step <b>648</b>.
[0267] Step <b>648</b> sets the data output columns to one. Sub-process <b>640</b> continues with step <b>650</b>.
[0268] Step <b>650</b> is a decision. If the new algorithm is an e<sup>x </sup>expansion, sub-process <b>640</b> continues with step <b>654</b>; otherwise sub-process <b>640</b> continues with step <b>652</b>.
[0269] Step <b>652</b> is a decision. If the new algorithm is a Sigma square root expansion, sub-process <b>640</b> continues with step <b>658</b>; otherwise sub-process <b>640</b> continues with step <b>656</b>.
[0270] Step <b>654</b> sets the output data size to the size of the agglomeration structure for the ex expansion. Sub-process <b>640</b> terminates at step <b>660</b>, returning control to the invoking process.
[0271] Step <b>656</b> sets the data output size to the number of terms in the series plus one. Sub-process <b>640</b> terminates at step <b>660</b>, returning control to the invoking process.
[0272] Step <b>658</b> sets the data output size to the number of digits in an ASCII floating point number. Sub-process <b>640</b> terminates at step <b>660</b>, returning control to the invoking process.
[0273]FIG. 43 illustrates one sub-process <b>670</b> for calculating the output data size when the input data is a single image. Sub-process <b>670</b> starts at step <b>672</b> and continues with step <b>674</b>.
[0274] Step <b>674</b> sets the data output rows equal to the data input rows. Sub-process <b>670</b> continues with step <b>676</b>.
[0275] Step <b>676</b> sets the data output columns equal to the data input columns. Sub-process <b>670</b> continues with step <b>678</b>.
[0276] Step <b>678</b> is a decision. If the new algorithm is a FFT computation, sub-process <b>670</b> continues with step <b>682</b>; otherwise sub-process <b>670</b> continues with step <b>680</b>.
[0277] Step <b>680</b> sets the data output size equal to the data input element size. Sub-process <b>670</b> continues with step <b>684</b>;
[0278] Step <b>682</b> sets the data output element size to the size of two double precision floating-point numbers. Sub-process <b>670</b> continues with step <b>684</b>.
[0279] Step <b>684</b> sets the data output size to the data output rows×columns×element size. Sub-process <b>670</b> terminates at step <b>686</b>, returning control to the invoking process.
[0280]FIG. 44 illustrates one sub-process <b>690</b> for calculating the size of the output data when the new algorithm has two input images. Sub-process <b>690</b> starts at step <b>692</b> and continues with step <b>694</b>.
[0281] Step <b>694</b> sets the data output rows equal to the data input first rows, minus the data input second rows, plus one. Sub-process <b>690</b> continues with step <b>696</b>.
[0282] Step <b>696</b> sets the data output columns to the data input first columns, minus the data input second columns, plus one. Sub-process <b>690</b> continues with step <b>698</b>.
[0283] Step <b>698</b> is a decision. If the new algorithm is a convolve computation, sub-process <b>690</b> continues with step <b>700</b>; otherwise sub-process <b>690</b> continues with step <b>702</b>.
[0284] Step <b>700</b> sets the data output element size to the size of a single precision floating-point number. Sub-process <b>690</b> continues with step <b>704</b>.
[0285] Step <b>702</b> sets the data output element size to the size of a double precision floating-point number. Sub-process <b>690</b> continues with step <b>704</b>.
[0286] Step <b>704</b> sets the output data size to the data output rows×columns×element size. Sub-process <b>690</b> terminates at step <b>706</b>, returning control to the invoking process.
[0287]FIG. 45 illustrates one sub-process <b>710</b> for calculating parameters for the mesh. Sub-process <b>710</b> starts at step <b>712</b> and continues with step <b>714</b>.
[0288] Step <b>714</b> is a decision. If the new algorithm is a Pi, Sigma square root or LN(X) computation, sub-process <b>710</b> continues with step <b>716</b>; otherwise sub-process <b>710</b> continues with step <b>718</b>.
[0289] Step <b>716</b> invokes a sub-process to calculate a single dimensional mesh for every value of N, described in FIG. 46. Sub-process <b>710</b> terminates with step <b>732</b>, returning control to the invoking process.
[0290] Step <b>718</b> is a decision. If the new algorithm is a e<sup>x </sup>computation, sub-process <b>710</b> continues with step <b>720</b>; otherwise sub-process <b>710</b> continues with step <b>722</b>.
[0291] Step <b>720</b> invokes a sub-process to calculate a 1-dimensional continuous block mesh, defined in FIG. 47. Sub-process <b>710</b> terminates with step <b>732</b>, returning control to the invoking process.
[0292] Step <b>722</b> is a decision. If the new algorithm is a convolution, normalized cross correlate or edge detection algorithm, sub-process <b>710</b> continues with step <b>724</b>; otherwise sub-process <b>710</b> continues with step <b>726</b>.
[0293] Step <b>724</b> invokes a sub-process to calculate a 2-dimensional continuous block row mesh, defined in FIG. 47. Sub-process <b>710</b> terminates with step <b>732</b>, returning control to the invoking process.
[0294] Step <b>726</b> is a decision. If the new algorithm is an FFT calculation, sub-process <b>710</b> continues with step <b>728</b>; otherwise sub-process <b>710</b> continues with step <b>730</b>.
[0295] Step <b>728</b> invokes a sub-process to calculate a 2-dimensional continuous block row and column mesh, defined in FIG. 49. Sub-process <b>710</b> terminates with step <b>732</b>, returning control to the invoking process.
[0296] Step <b>730</b> invokes a sub-process to calculate a 2-dimensional continuous block column mesh, defined in FIG. 48. Sub-process <b>710</b> terminates with step <b>732</b>, returning control to the invoking process.
[0297]FIG. 46 illustrates one sub-process to <b>740</b> to calculate a single dimensional every N mesh. Sub-process <b>740</b> starts with step <b>742</b> and continues with step <b>744</b>.
[0298] Step <b>744</b> sets the mesh input size to the data input rows. Sub-process <b>740</b> continues with step <b>746</b>.
[0299] Step <b>746</b> sets the mesh input offset to the data input element size x the processing node's cascade position. Sub-process <b>740</b> continues with step <b>748</b>.
[0300] Step <b>748</b> sets the mesh input step equal to the number of processing nodes working on the algorithm processing request. Sub-process <b>740</b> continues with step <b>750</b>.
[0301] Step <b>750</b> is a decision. If the mesh is to start at term 0, sub-process <b>740</b> continues with step <b>752</b>; otherwise sub-process <b>740</b> continues with step <b>754</b>.
[0302] Step <b>752</b> sets the mesh start equal to the processing node's cascade position. Sub-process <b>740</b> continues with step <b>756</b>;
[0303] Step <b>754</b>. sets the mesh input start to the processing node's cascade position, plus one. Sub-process <b>740</b> continues with step <b>756</b>.
[0304] Step <b>756</b> sets the mesh output size to the data input rows. Sub-process <b>740</b> continues with step <b>758</b>.
[0305] Step <b>758</b> sets the mesh output offset to the data input element size x the processing node's cascade position. Sub-process <b>740</b> continues with step <b>760</b>.
[0306] Step <b>760</b> sets the mesh output step equal to the count of processing nodes working on the algorithm processing request. Sub-process <b>740</b> continues with step <b>762</b>.
[0307] Step <b>762</b> is a decision. If the mesh is to start at term zero, sub-process <b>740</b> continues with step <b>764</b>; otherwise sub-process <b>740</b> continues with step <b>766</b>.
[0308] Step <b>764</b> set the mesh output start equal to the processing node's cascade position. Sub-process <b>740</b> terminates at step <b>768</b>, returning control to the invoking process.
[0309] Step <b>766</b> sets the mesh output start to the processing node's cascade position, plus one. Sub-process <b>740</b> terminates at step <b>768</b>, returning control to the invoking process.
[0310]FIG. 47 illustrates one sub-process <b>780</b> to compute a single dimensional continuous block mesh, which is also the same sub-process <b>780</b> for computing a two dimensional continuous block row mesh. Sub-process <b>780</b> starts at step <b>782</b> and continues with step <b>784</b>.
[0311] Step <b>784</b> invokes a sub-process to calculate a linear mesh based on data input rows, as defined in FIG. 50. Sub-process <b>780</b> continues with step <b>786</b>.
[0312] Step <b>786</b> invokes a sub-process to calculate a linear mesh based on data output rows, defined in FIG. 51. Sub-process <b>780</b> terminates at step <b>788</b>, returning control to the invoking process.
[0313]FIG. 48 is a flow chart illustrating one sub-process <b>800</b> for calculating a 2-dimensional continuous block mesh. Sub-process <b>800</b> starts at step <b>802</b> and continues with step <b>804</b>.
[0314] Step <b>804</b> invokes a sub-process to compute a linear mesh on data input columns, defined in FIG. 50. Sub-process <b>800</b> continues with step <b>806</b>.
[0315] Step <b>806</b> invokes a sub-process to calculate a linear mesh based on data output columns, defined in FIG. 51. Sub-process <b>800</b> terminates at step <b>808</b>, returning control to the invoking process.
[0316]FIG. 49 is a flow chart illustrating one sub-process <b>820</b> for calculating a 2-dimensional continuous clock row and column mesh. Sub-process <b>820</b> starts at step <b>822</b> and continues with step <b>824</b>.
[0317] Step <b>824</b> invokes a sub-function to calculate a linear mesh on data input rows, defined in FIG. 50. Sub-process <b>820</b> continues with step <b>826</b>.
[0318] Step <b>826</b> invokes a sub-process to calculate a linear mesh based on data output rows. Sub-process <b>820</b> continues with step <b>828</b>.
[0319] Step <b>828</b> invokes a sub-function to calculate a linear mesh on data input columns, defined in FIG. 50. Sub-process <b>820</b> continues with step <b>830</b>.
[0320] Step <b>830</b> invokes a sub-process to calculate a linear mesh based on data output columns. Sub-process <b>820</b> terminate at step <b>832</b>, returning control to the invoking process.
[0321]FIG. 50 is a flow chart illustrating one sub-process <b>840</b> for computing a linear mesh on rows. Sub-process <b>840</b> may be invoked for both data input and data output calculations. Sub-process <b>840</b> starts at step <b>842</b> and continues with step <b>844</b>.
[0322] Step <b>844</b> is a decision. If a second image is present with a kernel, sub-process <b>840</b> continues with step <b>848</b>; otherwise sub-process <b>840</b> continues with step <b>846</b>.
[0323] Step <b>846</b> sets the input size to rows. Sub-process <b>840</b> continues with step <b>850</b>.
[0324] Step <b>848</b> sets the input size to the first rows minus the second rows. Sub-process <b>840</b> continues with step <b>850</b>.
[0325] Step <b>850</b> sets the mesh size to the input size divided by the count of processing nodes working on the algorithm processing request. Sub-process <b>840</b> continues with step <b>852</b>.
[0326] Step <b>852</b> sets the mesh index to the cascade position of the processing node on which the algorithm is running x mesh size. Sub-process <b>840</b> continues with step <b>854</b>,
[0327] Step <b>854</b> is a decision. If a second image is present with a kernel, sub-process <b>840</b> continues with step <b>856</b>; otherwise sub-process <b>840</b> continues with step <b>858</b>.
[0328] Step <b>856</b> adds second rows minus one to the mesh index. Sub-process <b>840</b> continues with step <b>858</b>.
[0329] Step <b>858</b> sets the mesh remainder to the remainder of the input size divided by the count of processing nodes working on the algorithm processing request. Sub-process <b>840</b> continues with step <b>860</b>.
[0330] Step <b>860</b> is a decision. If the cascade position of the processing node is less than the mesh remainder calculated in step <b>858</b>, sub-process <b>840</b> continues with step <b>864</b>; otherwise sub-process <b>840</b> continues with step <b>862</b>.
[0331] Step <b>862</b> adds the cascade position of the processing node to the mesh index. Sub-process <b>840</b> continues with step <b>866</b>.
[0332] Step <b>864</b> adds the mesh remainder calculated in step <b>858</b> to the mesh index. Sub-process <b>840</b> terminates at step <b>868</b>, returning control to the invoking process.
[0333] Step <b>866</b> increments the mesh size. Sub-process <b>840</b> terminates at step <b>868</b>, returning control to the invoking process.
[0334]FIG. 51 is a flow chart illustrating one sub-process <b>880</b> for computing a linear mesh on columns. Sub-process may be invoked for both data input and data output calculations. Sub-process <b>880</b> starts at step <b>882</b> and continues with step <b>884</b>.
[0335] Step <b>884</b> is a decision. If a second image is present with a kernel sub-process <b>880</b> continues with step <b>888</b>; otherwise sub-process <b>880</b> continues with step <b>886</b>.
[0336] Step <b>886</b> sets the input size to columns. Sub-process <b>880</b> continues with step <b>890</b>.
[0337] Step <b>888</b> sets the input size to the first columns minus second rows. Sub-process <b>880</b> continues with step <b>890</b>.
[0338] Step <b>890</b> sets the mesh size to the input size divided by the count of processing nodes working on the algorithm processing request. Sub-process <b>880</b> continues with step <b>892</b>.
[0339] Step <b>892</b> sets the mesh index to the cascade position of the processing node on which the algorithm is running x mesh size. Sub-process <b>880</b> continues with step <b>894</b>,
[0340] Step <b>894</b> is a decision. If a second image is present with a kernel, sub-process <b>880</b> continues with step <b>896</b>; otherwise sub-process <b>880</b> continues with step <b>898</b>.
[0341] Step <b>896</b> adds second columns minus one to the mesh index. Sub-process <b>880</b> continues with step <b>898</b>.
[0342] Step <b>898</b> sets the mesh remainder to the remainder of the input size divided by the count of processing nodes working on the algorithm processing request. Sub-process <b>880</b> continues with step <b>900</b>.
[0343] Step <b>900</b> is a decision. If the cascade position of the processing node is less than the mesh remainder calculated in step <b>898</b>, sub-process <b>880</b> continues with step <b>904</b>; otherwise sub-process <b>880</b> continues with step <b>902</b>.
[0344] Step <b>902</b> adds the cascade position of the processing node to the mesh index. Sub-process <b>880</b> continues with step <b>906</b>.
[0345] Step <b>904</b> adds the mesh remainder calculated in step <b>898</b> to the mesh index. Sub-process <b>880</b> terminates at step <b>908</b>, returning control to the invoking process.
[0346] Step <b>906</b> increments the mesh size. Sub-process <b>880</b> terminates at step <b>908</b>, returning control to the invoking process.
[0347]FIG. 52 is a flow chart illustrating one sub-process <b>920</b> to acquire input data needed by the processing node to perform the algorithm processing request. Sub-process <b>920</b> starts at step <b>922</b> and continues with step <b>924</b>.
[0348] Step <b>924</b> is a decision. If the algorithm processing request expects input data, sub-process <b>920</b> continues with step <b>926</b>; otherwise sub-process <b>920</b> terminates at step <b>930</b>, returning control to the invoking process.
[0349] Step <b>926</b> is a decision. If the data is sent by the HC as a data broadcast, sub-process <b>920</b> continues with step <b>928</b>; otherwise sub-process <b>920</b> terminates at step <b>930</b>, returning control to the invoking process.
[0350] Step <b>928</b> invokes a sub-process to use multi-cast tools to receive the broadcast message, defined in FIG. 53. Sub-process <b>920</b> terminates at step <b>930</b>, returning control to the invoking process.
[0351]FIG. 53 is a flowchart illustrating one sub-process <b>940</b> for using multicast tools to receive the broadcast data. Sub-process <b>940</b> starts at step <b>929</b> and continues with step <b>944</b>.
[0352] Step <b>944</b> opens the multicast socket to receive the broadcast. Sub-process <b>940</b> continues with step <b>946</b>.
[0353] Step <b>946</b> receives the multicast data. Sub-process <b>940</b> continues with step <b>948</b>.
[0354] Step <b>948</b> is a decision. If there is more data to receive, sub-process <b>940</b> continues with step <b>946</b>; otherwise sub-process <b>940</b> continues with step <b>950</b>.
[0355] Step <b>950</b> closes the multicast socket opened in step <b>944</b>. Sub-process <b>940</b> terminates at step <b>952</b>, returning control to the invoking process.
[0356]FIG. 54 is a flowchart illustrating one sub-process <b>960</b> to receive results from downstream processing nodes, in agglomeration. Sub-process <b>960</b> starts at step <b>962</b> and continues with step <b>964</b>.
[0357] Step <b>964</b> is a decision. If the agglomeration is of type with multi-home nodes, sub-process <b>960</b> terminates at step <b>980</b>, returning control to the invoking process; otherwise sub-process <b>960</b> continues with step <b>966</b>.
[0358] Step <b>966</b> determines the number of message to expect from downstream processing nodes. Sub-process <b>960</b> continues with step <b>968</b>.
[0359] Step <b>968</b> is a decision. If there are no expected messages, sub-process <b>960</b> terminates at step <b>980</b>, returning control to the invoking process; otherwise sub-process <b>960</b> continues with step <b>970</b>.
[0360] Step <b>970</b> invokes a sub-process to set the result pointers that are used for storing the received results. Sub-process <b>960</b> continues with step <b>972</b>.
[0361] Step <b>972</b> receives a message with data attached. Sub-process <b>960</b> continues with step <b>974</b>.
[0362] Step <b>974</b> invokes a sub-process to combine results with prior results. Sub-process <b>960</b> continues with step <b>976</b>.
[0363] Step <b>976</b> is a decision. If more messages with attached data are expected, sub-process <b>960</b> continues with step <b>970</b>; otherwise sub-process <b>960</b> continues with step <b>978</b>.
[0364] Step <b>978</b> invokes a sub-process to clean up the storage after agglomeration is complete. Sub-process <b>960</b> terminates at step <b>980</b>, returning control to the invoking process.
[0365]FIG. 55 is a flowchart illustrating one sub-process <b>990</b> for setting the results pointers ready for the agglomeration results. Sub-process <b>990</b> start at step <b>992</b> and continues with step <b>994</b>.
[0366] Step <b>994</b> is a decision. If the agglomeration type is a row mesh, the existing memory allocated for the input data can be used, and sub-process <b>990</b> continues with step <b>1000</b>; otherwise sub-process <b>990</b> continues with step <b>996</b>.
[0367] Step <b>996</b> is a decision. If the agglomeration type is result list, sub-process <b>990</b> continues with step <b>1002</b>; otherwise sub-process <b>990</b> continues with step <b>998</b>.
[0368] Step <b>998</b> is a decision. If the pointer is currently in use, sub-process <b>990</b> terminates at step <b>1006</b>, returning control to the invoking process; otherwise sub-process <b>990</b> continues with step <b>1002</b>.
[0369] Step <b>1000</b> sets the data pointer to point at the input data space. Sub-process <b>990</b> terminates at step <b>1006</b>, returning control to the invoking process.
[0370] Step <b>1002</b> allocated more memory for agglomeration. Sub-process <b>990</b> continues with step <b>1004</b>.
[0371] Step <b>1004</b> sets the pointer to the memory space allocated by step <b>1002</b>. Sub-process <b>990</b> terminates at step <b>1006</b>, returning control to the invoking process.
[0372]FIG. 56 is a flowchart illustrating one sub-process <b>1020</b> for processing the results during agglomeration. Sub-process <b>1020</b> starts at step <b>1022</b> and continues with step <b>1024</b>.
[0373] Step <b>1024</b> is a decision. If the agglomeration type is arbitrary precision addition, then sub-process <b>1020</b> continues with step <b>1026</b>; otherwise sub-process <b>1020</b> continues with step <b>1028</b>.
[0374] Step <b>1026</b> converts the received agglomeration data to an APFLOAT number and adds it to the accumulated result. Sub-process <b>1020</b> terminates at step <b>1040</b>, returning control to the invoking process.
[0375] Step <b>1028</b> is a decision. If the agglomeration type is floating point addition, sub-process <b>1020</b> continues with step <b>1030</b>; otherwise sub-process <b>1020</b> continues with step <b>1032</b>.
[0376] Step <b>1030</b> converts the received agglomeration data to a floating point number and adds it to the accumulated result. Sub-process <b>1020</b> terminates at step <b>1040</b>, returning control to the invoking process.
[0377] Step <b>1032</b> is a decision. If the agglomeration type is save largest, sub-process <b>1020</b> continues with step <b>1034</b>; otherwise sub-process <b>1020</b> continues with step <b>1036</b>.
[0378] Step <b>1034</b> compares the received agglomeration result with a stored value, and, if larger, replaces the stored value with the received agglomeration data. Sub-process <b>1020</b> terminates at step <b>1040</b>, returning control to the invoking process.
[0379] Step <b>1036</b> is a decision. If the agglomeration type is result list, sub-process <b>1020</b> continues with step <b>1038</b>; otherwise sub-process <b>1020</b> terminates at step <b>1040</b>, returning control to the invoking process.
[0380] Step <b>1038</b> adds the result pointer to the result list and increments the result counter. Sub-process <b>1020</b> terminates at step <b>1040</b>, returning control to the invoking process.
[0381]FIG. 57 is a flowchart illustrating one sub-process <b>1050</b> for cleaning up the used result space after agglomeration is complete. Sub-process <b>1050</b> starts at step <b>1052</b> and continues with step <b>1054</b>.
[0382] Step <b>1054</b> is a decision. If the agglomeration method was ROWMESH, sub-process <b>1050</b> terminates at step <b>1058</b>, retuning control to the invoking process; otherwise sub-process <b>1050</b> continues with step <b>1056</b>.
[0383] Step <b>1056</b> frees the allocated memory space. Sub-process <b>1050</b> terminates at step <b>1058</b>, retuning control to the invoking process.
[0384]FIG. 58 is a flowchart illustrating one sub-process <b>1070</b> for returning the local or agglomerated results to the correct processing node or the home node. Sub-process <b>1070</b> starts at step <b>1072</b> and continues with step <b>1074</b>.
[0385] Step <b>1074</b> invokes a sub-process to get the address of the destination node to receive the results, defined in FIG. 59. Sub-process <b>1070</b> continues with step <b>1076</b>.
[0386] Step <b>1076</b> invokes a sub-process to get the format for the result message, defined in FIG. 60. Sub-process <b>1070</b> continues with step <b>1078</b>.
[0387] Step <b>1078</b> invokes a sub-process to build the result message, defined in FIG. 61. Sub-process <b>1070</b> continues with step <b>1080</b>.
[0388] Step <b>1080</b> invokes a sub-process to send the result message to the destination node, defined in FIG. 62. Sub-process <b>1070</b> terminates at step <b>1082</b>, returning control to the invoking process.
[0389]FIG. 59 is a flowchart illustrating one sub-process <b>1090</b> for determining the address of the node to receive the agglomeration results. Sub-process <b>1090</b> starts at step <b>1092</b> and continues with step <b>1094</b>.
[0390] Step <b>1094</b> is a decision. If the results are to be sent to the home node, sub-process <b>1090</b> continues with step <b>1098</b>; otherwise sub-process <b>1090</b> continues with step <b>1096</b>.
[0391] Step <b>1096</b> gets the address for the upstream processing node. Sub-process <b>1090</b> terminates at step <b>1100</b>, returning control to the invoking process.
[0392] Step <b>1098</b> gets the address of the home node. Sub-process <b>1090</b> terminates at step <b>1100</b>, returning control to the invoking process.
[0393]FIG. 60 is a flowchart illustrating one sub-process <b>1110</b> for getting the format of the results message. Sub-process <b>1110</b> starts at step <b>1112</b> and continues with step <b>1114</b>.
[0394] Step <b>1114</b> is a decision. If the result is for the home node, sub-process <b>1110</b> continues with step <b>1118</b>, otherwise sub-process <b>1110</b> continues with step <b>1116</b>.
[0395] Step <b>1116</b> gets the message format for the upstream processing node. Sub-process <b>1110</b> terminates at step <b>1120</b>, returning control to the invoking process.
[0396] Step <b>1118</b> gets the message format for the home node. Sub-process <b>1110</b> terminates at step <b>1120</b>, returning control to the invoking process.
[0397]FIG. 61 is a flowchart illustrating one sub-process <b>1130</b> for building the results message. Sub-process <b>1130</b> starts at step <b>1132</b> and continues with step <b>1134</b>.
[0398] Step <b>1134</b> is a decision. If the message is to be returned without a data header, sub-process <b>1030</b> continues with step <b>1136</b>; otherwise sub-process <b>1130</b> continues with step <b>1138</b>.
[0399] Step <b>1136</b> builds the message without a data header. Sub-process <b>1130</b> terminates at step <b>1144</b>, returning control to the invoking process.
[0400] Step <b>1138</b> is a decision. If the message it to be built with a data header, sub-process <b>1130</b> continues with step <b>1140</b>; otherwise sub-process <b>1130</b> continues with step <b>1142</b>.
[0401] Step <b>1142</b> builds a processing node results message. Sub-process <b>1130</b> terminates at step <b>1144</b>, returning control to the invoking process.
[0402]FIG. 62 is a flowchart illustrating one sub-process <b>1160</b> for sending the results message. Sub-process <b>1160</b> starts at step <b>1162</b> and continues with step <b>1164</b>.
[0403] Step <b>1164</b> opens the stream socket to the destination node. Sub-process <b>1160</b> continues with step <b>1166</b>.
[0404] Step <b>1166</b> sends the data down the stream. Sub-process <b>1160</b> continues with step <b>1168</b>.
[0405] Step <b>1168</b> is a decision. If there is more data to send, sub-process <b>1160</b> continues with step <b>1166</b>; otherwise sub-process <b>1160</b> continues with step <b>1170</b>.
[0406] Step <b>1170</b> closes the stream socket. Sub-process <b>1160</b> terminates at step <b>1172</b>, returning control to the invoking process.
[0407] Using Heterogeneous Computer Systems and Communication Channels to Build an HCAS.
[0408] Whilst it is preferred to use computer systems and communication channels of equal specification and performance in an HCAS, the HCAS may be constructed from systems and channels of varying specifications without any significant loss in system efficiency. For example, computer systems with differing processor speeds, single or multiple processor motherboards, and varying numbers of NICs can be utilized as home nodes and processing nodes in the same HCAS.
[0409] In parallel processing systems of the prior art, such imbalance of node specification would cause processing imbalances and hence significant efficiency losses in the cluster. On an HCAS, however, load balancing can be performed automatically.
[0410] As appreciated by those skilled in the art, many techniques are available for load balancing on parallel processing clusters. One technique that may be utilized on an HCAS is to proportionately allocate the amount of processing required by each processing node based on its processing and communication capability. For example, a processing node with a slower processor clock may be allocated less data to process, or elements to calculate in a series expansion, than a processing node with a faster processor clock. In another example, two processing nodes have identical processor clock speeds, but the first processing node has a communication channel with twice the bandwidth of the second processing node. The first processing node would be allocated more data than the second processing node as it would be able to receive more data in the time taken for the second processing node to receive data.
[0411] The expansion of the algorithm processing request in a HC is less influenced by the communication channel speed due to the small message size. Any imbalances in system performance during problem expansion due to communication channel bandwidth imbalances are insignificant.
FURTHER EXAMPLES
[0412] The following sections provide further examples of algorithms that may be run on the HCAS.
[0413] Decomposition of Two-Dimensional Object Transformation Data for Parallel Processing
[0414] Geometric transformations of a two-dimensional object include translation, scaling, rotation, and shearing. This section describes how the geometric transformation of a 2D object may be implemented on a HC, in one embodiment. The data partitioning handles an arbitrary number of parallel processing nodes and arbitrarily large objects.
[0415] A 2D object may be defined by its endpoints, expressed as coordinates in the (x,y) plane. Let the matrix M<sub>XY </sub>represent these N endpoints as column vectors: <maths id="MATH-US-00003" num="3"><math overflow="scroll"><mrow><msub><mi>M</mi><mi>XY</mi></msub><mo>=</mo><mrow><mo>[</mo><mtable><mtr><mtd><msub><mi>x</mi><mn>1</mn></msub></mtd><mtd><msub><mi>x</mi><mn>1</mn></msub></mtd><mtd><msub><mi>x</mi><mn>1</mn></msub></mtd><mtd><mi>…</mi></mtd><mtd><msub><mi>x</mi><mi>N</mi></msub></mtd></mtr><mtr><mtd><msub><mi>y</mi><mn>1</mn></msub></mtd><mtd><msub><mi>y</mi><mn>1</mn></msub></mtd><mtd><msub><mi>y</mi><mn>1</mn></msub></mtd><mtd><mi>…</mi></mtd><mtd><msub><mi>y</mi><mi>N</mi></msub></mtd></mtr></mtable><mo>]</mo></mrow></mrow></math><img file="US20030195938A1-20031016-M00003.TIF" id="EMI-M00003" he="21.12075" wi="216.027" img-format="tif" img-content="mf" /><attachments><attachment idref="MATHEMATICA-00003" attachment-type="nb" file="US20030195938A1-20031016-M00003.NB" /></attachments></maths>
[0416] The endpoints are converted to homogeneous coordinates, setting the third coordinate to 1, to create a new matrix M<sub>H</sub>: <maths id="MATH-US-00004" num="4"><math overflow="scroll"><mrow><msub><mi>M</mi><mi>H</mi></msub><mo>=</mo><mrow><mo>[</mo><mtable><mtr><mtd><msub><mi>x</mi><mn>1</mn></msub></mtd><mtd><msub><mi>x</mi><mn>1</mn></msub></mtd><mtd><msub><mi>x</mi><mn>1</mn></msub></mtd><mtd><mi>…</mi></mtd><mtd><msub><mi>x</mi><mi>N</mi></msub></mtd></mtr><mtr><mtd><msub><mi>y</mi><mn>1</mn></msub></mtd><mtd><msub><mi>y</mi><mn>1</mn></msub></mtd><mtd><msub><mi>y</mi><mn>1</mn></msub></mtd><mtd><mi>…</mi></mtd><mtd><msub><mi>y</mi><mi>N</mi></msub></mtd></mtr><mtr><mtd><mn>1</mn></mtd><mtd><mn>1</mn></mtd><mtd><mn>1</mn></mtd><mtd><mi>…</mi></mtd><mtd><mn>1</mn></mtd></mtr></mtable><mo>]</mo></mrow></mrow></math><img file="US20030195938A1-20031016-M00004.TIF" id="EMI-M00004" he="31.9221" wi="216.027" img-format="tif" img-content="mf" /><attachments><attachment idref="MATHEMATICA-00004" attachment-type="nb" file="US20030195938A1-20031016-M00004.NB" /></attachments></maths>
[0417] This conversion to homogeneous coordinates reduces the 2D transformation problem to the matrix multiplication problem: M<sub>T</sub>=T×M<sub>H</sub>, where T represents a 3×3 transform matrix and M<sub>H </sub>is the 2D object expressed as a 3×N matrix of homogeneous coordinates. The resulting product M<sub>T </sub>represents the transformation of the original 2D object.
[0418] The transform matrix T may equal one of the following matrices, depending on the type of transformation: <maths id="MATH-US-00005" num="5"><math overflow="scroll"><mrow><msub><mi>T</mi><mi>translate</mi></msub><mo>=</mo><mrow><mo>[</mo><mtable><mtr><mtd><mn>1</mn></mtd><mtd><mn>0</mn></mtd><mtd><msub><mi>t</mi><mi>x</mi></msub></mtd></mtr><mtr><mtd><mn>0</mn></mtd><mtd><mn>1</mn></mtd><mtd><msub><mi>t</mi><mi>y</mi></msub></mtd></mtr><mtr><mtd><mn>0</mn></mtd><mtd><mn>0</mn></mtd><mtd><mn>1</mn></mtd></mtr></mtable><mo>]</mo></mrow></mrow></math><math overflow="scroll"><mrow><msub><mi>t</mi><mi>x</mi></msub><mo>=</mo><mrow><mrow><mi>translation</mi><mo></mo><mstyle><mtext> </mtext></mstyle><mo></mo><mi>offset</mi><mo></mo><mstyle><mtext> </mtext></mstyle><mo></mo><mi>in</mi><mo></mo><mstyle><mtext> </mtext></mstyle><mo></mo><mi>x</mi></mrow><mo>-</mo><mi>direction</mi></mrow></mrow></math><math overflow="scroll"><mrow><msub><mi>t</mi><mi>y</mi></msub><mo>=</mo><mrow><mrow><mi>translation</mi><mo></mo><mstyle><mtext> </mtext></mstyle><mo></mo><mi>offset</mi><mo></mo><mstyle><mtext> </mtext></mstyle><mo></mo><mi>in</mi><mo></mo><mstyle><mtext> </mtext></mstyle><mo></mo><mi>y</mi></mrow><mo>-</mo><mi>direction</mi></mrow></mrow></math><math overflow="scroll"><mrow><msub><mi>T</mi><mi>scale</mi></msub><mo>=</mo><mrow><mo>[</mo><mtable><mtr><mtd><msub><mi>s</mi><mi>x</mi></msub></mtd><mtd><mn>0</mn></mtd><mtd><mn>0</mn></mtd></mtr><mtr><mtd><mn>0</mn></mtd><mtd><msub><mi>s</mi><mi>y</mi></msub></mtd><mtd><mn>0</mn></mtd></mtr><mtr><mtd><mn>0</mn></mtd><mtd><mn>0</mn></mtd><mtd><mn>1</mn></mtd></mtr></mtable><mo>]</mo></mrow></mrow></math><math overflow="scroll"><mrow><msub><mi>s</mi><mi>x</mi></msub><mo>=</mo><mrow><mrow><mi>scale</mi><mo></mo><mstyle><mtext> </mtext></mstyle><mo></mo><mi>factor</mi><mo></mo><mstyle><mtext> </mtext></mstyle><mo></mo><mi>in</mi><mo></mo><mstyle><mtext> </mtext></mstyle><mo></mo><mi>x</mi></mrow><mo>-</mo><mi>direction</mi></mrow></mrow></math><math overflow="scroll"><mrow><msub><mi>s</mi><mi>y</mi></msub><mo>=</mo><mrow><mrow><mi>scale</mi><mo></mo><mstyle><mtext> </mtext></mstyle><mo></mo><mi>factor</mi><mo></mo><mstyle><mtext> </mtext></mstyle><mo></mo><mi>in</mi><mo></mo><mstyle><mtext> </mtext></mstyle><mo></mo><mi>y</mi></mrow><mo>-</mo><mi>direction</mi></mrow></mrow></math><math overflow="scroll"><mrow><msub><mi>T</mi><mi>rotate</mi></msub><mo>=</mo><mrow><mo>[</mo><mtable><mtr><mtd><mrow><mi>cos</mi><mo></mo><mstyle><mtext> </mtext></mstyle><mo></mo><mi>θ</mi></mrow></mtd><mtd><mo>-</mo></mtd><mtd><mn>0</mn></mtd></mtr><mtr><mtd><mrow><mi>sin</mi><mo></mo><mstyle><mtext> </mtext></mstyle><mo></mo><mi>θ</mi></mrow></mtd><mtd><mrow><mi>cos</mi><mo></mo><mstyle><mtext> </mtext></mstyle><mo></mo><mi>θ</mi></mrow></mtd><mtd><mn>0</mn></mtd></mtr><mtr><mtd><mn>0</mn></mtd><mtd><mn>0</mn></mtd><mtd><mn>1</mn></mtd></mtr></mtable><mo>]</mo></mrow></mrow></math><math overflow="scroll"><mrow><mi>θ</mi><mo>=</mo><mrow><mi>angle</mi><mo></mo><mstyle><mtext> </mtext></mstyle><mo></mo><mi>of</mi><mo></mo><mstyle><mtext> </mtext></mstyle><mo></mo><mi>rotation</mi><mo></mo><mstyle><mtext> </mtext></mstyle><mo></mo><mi>about</mi><mo></mo><mstyle><mtext> </mtext></mstyle><mo></mo><mi>the</mi><mo></mo><mstyle><mtext> </mtext></mstyle><mo></mo><mi>orgin</mi></mrow></mrow></math><math overflow="scroll"><mrow><msub><mi>T</mi><mi>shear</mi></msub><mo>=</mo><mrow><mo>[</mo><mtable><mtr><mtd><mn>1</mn></mtd><mtd><msub><mi>sh</mi><mi>x</mi></msub></mtd><mtd><mn>0</mn></mtd></mtr><mtr><mtd><msub><mi>sh</mi><mi>y</mi></msub></mtd><mtd><mn>1</mn></mtd><mtd><mn>0</mn></mtd></mtr><mtr><mtd><mn>0</mn></mtd><mtd><mn>0</mn></mtd><mtd><mn>1</mn></mtd></mtr></mtable><mo>]</mo></mrow></mrow></math><math overflow="scroll"><mrow><msub><mi>sh</mi><mi>x</mi></msub><mo>=</mo><mrow><mrow><mi>shear</mi><mo></mo><mstyle><mtext> </mtext></mstyle><mo></mo><mi>constant</mi><mo></mo><mstyle><mtext> </mtext></mstyle><mo></mo><mi>in</mi><mo></mo><mstyle><mtext> </mtext></mstyle><mo></mo><mi>x</mi></mrow><mo>-</mo><mi>direction</mi></mrow></mrow></math><math overflow="scroll"><mrow><msub><mi>sh</mi><mi>y</mi></msub><mo>=</mo><mrow><mrow><mi>shear</mi><mo></mo><mstyle><mtext> </mtext></mstyle><mo></mo><mi>constant</mi><mo></mo><mstyle><mtext> </mtext></mstyle><mo></mo><mi>in</mi><mo></mo><mstyle><mtext> </mtext></mstyle><mo></mo><mi>y</mi></mrow><mo>-</mo><mi>direction</mi></mrow></mrow></math><img file="US20030195938A1-20031016-M00005.TIF" id="EMI-M00005" he="226.1196" wi="216.027" img-format="tif" img-content="mf" /><attachments><attachment idref="MATHEMATICA-00005" attachment-type="nb" file="US20030195938A1-20031016-M00005.NB" /></attachments></maths>
[0419] Now the 2D object transformation problem has been reduced to a matrix multiplication problem to be implemented on the HC. An arbitrary matrix is shown below in Table 21 an array of numbers arranged in J rows and K columns. Data partitioning evenly distributes the elements of a matrix over the parallel processing nodes. <tables id="TABLE-US-00024" num="24"><table frame="none" colsep="0" rowsep="0"><tgroup align="left" colsep="0" rowsep="0" cols="1"><colspec colname="1" colwidth="217PT" align="center" /><thead><row><entry namest="1" nameend="1" align="center">TABLE 21</entry></row></thead><tbody valign="top"><row><entry /></row><row><entry namest="1" nameend="1" align="center" rowsep="1" /></row><row><entry>A Matrix for the HC</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="2"><colspec colname="OFFSET" colwidth="42PT" align="left" /><colspec colname="1" colwidth="175PT" align="center" /><tbody valign="top"><row><entry /><entry>columns</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="7"><colspec colname="OFFSET" colwidth="14PT" align="left" /><colspec colname="1" colwidth="28PT" align="center" /><colspec colname="2" colwidth="42PT" align="center" /><colspec colname="3" colwidth="21PT" align="center" /><colspec colname="4" colwidth="42PT" align="center" /><colspec colname="5" colwidth="21PT" align="center" /><colspec colname="6" colwidth="49PT" align="center" /><tbody valign="top"><row><entry /><entry>rows</entry><entry>1</entry><entry>2</entry><entry>3</entry><entry>. . .</entry><entry>K</entry></row><row><entry /><entry namest="OFFSET" nameend="6" align="center" rowsep="1" /></row><row><entry /><entry>1</entry><entry>1,1</entry><entry>1,2</entry><entry>1,3</entry><entry>. . .</entry><entry>1,K</entry></row><row><entry /><entry>2</entry><entry>2,1</entry><entry>2,2</entry><entry>2,3</entry><entry>. . .</entry><entry>2,K</entry></row><row><entry /><entry>3</entry><entry>3,1</entry><entry>3,2</entry><entry>3,3</entry><entry>. . .</entry><entry>3,K</entry></row><row><entry /><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry></row><row><entry /><entry>J</entry><entry>J,1</entry><entry>J,2</entry><entry>J,3</entry><entry>. . .</entry><entry>J,K</entry></row><row><entry /><entry namest="OFFSET" nameend="6" align="center" rowsep="1" /></row></tbody></tgroup></table></tables>
[0420] Consider the matrix multiplication problem M<sub>T</sub>=T×M, where matrix T has dimensions 3×3, matrix M has dimensions 3×N, and the resulting product M<sub>T </sub>also has dimensions 3×N. Let 3×N represent the size of the product M<sub>T</sub>. In the HC implementation of matrix multiplication, each of the P parallel processing nodes receives a copy of T and M. Assuming 3×N>P, the solution size 3×N is divided by the number of nodes P to obtain an integer quotient W and a remainder R. W elements are assigned to each of the P nodes. Any remainder R is distributed, one element per node, to each of the first R nodes. Thus the first R nodes are assigned an element count of W+1, and each and the remaining P−R nodes are assigned an element count of W. The total number of elements assigned equals: R(W+1)+(P−R)W=PW+R=3×N.
[0421] The described approach maintains the node computational load balanced to within one element.
[0422] Consider the distribution of a product matrix M<sub>T </sub>consisting of 3 rows and 50 columns on a 7-node HC. In this case, the integer quotient W is 150/7=21, and the remainder R is 3. Processing nodes P<sub>1 </sub>through P<sub>3 </sub>are assigned 22 elements each; nodes P<sub>4 </sub>through P<sub>7 </sub>are assigned 21 each. The matrix data partitioning is shown in Table 22, computed with 150 total elements. <tables id="TABLE-US-00025" num="25"><table frame="none" colsep="0" rowsep="0"><tgroup align="left" colsep="0" rowsep="0" cols="1"><colspec colname="1" colwidth="217PT" align="center" /><thead><row><entry namest="1" nameend="1" align="center">TABLE 22</entry></row></thead><tbody valign="top"><row><entry /></row><row><entry namest="1" nameend="1" align="center" rowsep="1" /></row><row><entry>Matrix Data Partitioning</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="5"><colspec colname="OFFSET" colwidth="14PT" align="left" /><colspec colname="1" colwidth="35PT" align="center" /><colspec colname="2" colwidth="56PT" align="center" /><colspec colname="3" colwidth="35PT" align="center" /><colspec colname="4" colwidth="77PT" align="center" /><tbody valign="top"><row><entry /><entry /><entry /><entry>Column</entry><entry>Number of</entry></row><row><entry /><entry /><entry /><entry>Indices of</entry><entry>Elements</entry></row><row><entry /><entry>Processing</entry><entry>Row</entry><entry>Elements</entry><entry>Computed Per</entry></row><row><entry /><entry>Node</entry><entry>Index</entry><entry>Computed</entry><entry>Row</entry></row><row><entry /><entry namest="OFFSET" nameend="4" align="center" rowsep="1" /></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="5"><colspec colname="OFFSET" colwidth="14PT" align="left" /><colspec colname="1" colwidth="35PT" align="center" /><colspec colname="2" colwidth="56PT" align="center" /><colspec colname="3" colwidth="35PT" align="center" /><colspec colname="4" colwidth="77PT" align="char" char="." /><tbody valign="top"><row><entry /><entry>P<sub>1</sub></entry><entry>1</entry><entry> 1 → 22</entry><entry>22</entry></row><row><entry /><entry>P<sub>2</sub></entry><entry>1</entry><entry>23 → 44</entry><entry>22</entry></row><row><entry /><entry>P<sub>3</sub></entry><entry>1</entry><entry>45 → 50</entry><entry>6</entry></row><row><entry /><entry /><entry>2</entry><entry> 1 → 16</entry><entry>16</entry></row><row><entry /><entry>P<sub>4</sub></entry><entry>2</entry><entry>17 → 37</entry><entry>21</entry></row><row><entry /><entry>P<sub>5</sub></entry><entry>2</entry><entry>38 → 50</entry><entry>13</entry></row><row><entry /><entry /><entry>3</entry><entry>1 → 8</entry><entry>8</entry></row><row><entry /><entry>P<sub>6</sub></entry><entry>3</entry><entry> 9 → 29</entry><entry>21</entry></row><row><entry /><entry>P<sub>7</sub></entry><entry>3</entry><entry>30 → 50</entry><entry>21</entry></row><row><entry /><entry namest="OFFSET" nameend="4" align="center" rowsep="1" /></row></tbody></tgroup></table></tables>
[0423] Processing node P<sub>1 </sub>multiplies row 1 of transform matrix T with columns 1 through 22 of object matrix M to compute its part of the result. Processing node P<sub>2 </sub>multiplies row 1 of transform matrix T with columns 23 through 44 of object matrix M to compute its part of the result. Processing node P<sub>3 </sub>uses rows 1 and 2 of transform matrix T and columns 1 through 16 and 45 through 50 of object matrix M to compute its result. The distribution of work is spread similarly for processing nodes P<sub>4 </sub>through P<sub>7</sub>.
[0424] The home node performs the matrix data partitioning. Messages describing the command (in this case a matrix multiply) and data partitioning are sent out to the processing nodes. Once the nodes have received their command messages, they wait for the home node to send the input matrix data. The data may be broadcast such that all nodes in the cascade receive it at the same time. Each node thus receives both of the input matrices, which may be more efficient than sending each individual node a separate message with just its piece of input data. This is especially important when considering large numbers of parallel processing nodes. Once a node receives the input matrices, it can proceed with computing the product independent of the other nodes. When the matrix multiply results are ready, they are accumulated up to the home node and merged into the final result. At this point the process is complete.
[0425] As this example shows, the 2D object transformation data distribution applied to the HC accommodates arbitrary sized objects in a simple and efficient manner.
[0426] Decomposition of Matrix Multiplication Data for Parallel Processing
[0427] This section describes how matrix multiplication may be implemented on a HC. The described matrix partitioning can handle an arbitrary number of parallel processing nodes and arbitrarily large matrices.
[0428] An arbitrary matrix is shown in Table 23 as an array of numbers arranged in M rows and N columns. Data partitioning evenly distributes the elements of the matrix over the parallel processing nodes. <tables id="TABLE-US-00026" num="26"><table frame="none" colsep="0" rowsep="0"><tgroup align="left" colsep="0" rowsep="0" cols="1"><colspec colname="1" colwidth="217PT" align="center" /><thead><row><entry namest="1" nameend="1" align="center">TABLE 23</entry></row></thead><tbody valign="top"><row><entry /></row><row><entry namest="1" nameend="1" align="center" rowsep="1" /></row><row><entry>Matrix Distribution on a HC</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="2"><colspec colname="OFFSET" colwidth="35PT" align="left" /><colspec colname="1" colwidth="182PT" align="center" /><tbody valign="top"><row><entry /><entry>columns</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="7"><colspec colname="OFFSET" colwidth="14PT" align="left" /><colspec colname="1" colwidth="21PT" align="center" /><colspec colname="2" colwidth="49PT" align="center" /><colspec colname="3" colwidth="21PT" align="center" /><colspec colname="4" colwidth="49PT" align="center" /><colspec colname="5" colwidth="14PT" align="center" /><colspec colname="6" colwidth="49PT" align="center" /><tbody valign="top"><row><entry /><entry>rows</entry><entry>1</entry><entry>2</entry><entry>3</entry><entry>. . .</entry><entry>N</entry></row><row><entry /><entry namest="OFFSET" nameend="6" align="center" rowsep="1" /></row><row><entry /><entry>1</entry><entry>1,1</entry><entry>1,2</entry><entry>1,3</entry><entry>. . .</entry><entry>1,N</entry></row><row><entry /><entry>2</entry><entry>2,1</entry><entry>2,2</entry><entry>2,3</entry><entry>. . .</entry><entry>2,N</entry></row><row><entry /><entry>3</entry><entry>3,1</entry><entry>3,2</entry><entry>3,3</entry><entry>. . .</entry><entry>3,N</entry></row><row><entry /><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry></row><row><entry /><entry>M</entry><entry>M,1</entry><entry>M,2</entry><entry>M,3</entry><entry>. . .</entry><entry>M,N</entry></row><row><entry /><entry namest="OFFSET" nameend="6" align="center" rowsep="1" /></row></tbody></tgroup></table></tables>
[0429] Consider the matrix multiplication problem A×B=C, where input matrix A has dimensions M×N, input matrix B has dimensions N×Q, and the resulting product C has dimensions M×Q. Let E=M×Q. represent the size of the product C. In the matrix multiplication implementation, each of the P parallel processing nodes receives a copy of A and B. Assuming E>P, solution size E is divided by the number of nodes P to obtain an integer quotient W and a remainder R. W elements are assigned to each of the P nodes. Any remainder R is distributed, one element per node, to each of the first R nodes. Thus the first R nodes are assigned an element count of E<sub>i</sub>=W1+1 each and the remaining P−R nodes will be assigned an element count of E<sub>i</sub>=W. The total number of elements assigned equals R(W+1)+(P−R)W=PW+R=E.
[0430] Once again, the node computational load is balanced to within one element.
[0431] Consider the distribution of a product matrix consisting of 50 rows and 50 columns in a 7-node HC. In this case, the integer quotient Q is 2500÷7=357 and the remainder R is 1. Processing node P<sub>1 </sub>is assigned 358 elements; nodes P<sub>2 </sub>through P<sub>7 </sub>are assigned <b>357</b> each. The matrix data partitioning is shown in Table 24. <tables id="TABLE-US-00027" num="27"><table frame="none" colsep="0" rowsep="0"><tgroup align="left" colsep="0" rowsep="0" cols="1"><colspec colname="1" colwidth="217PT" align="center" /><thead><row><entry namest="1" nameend="1" align="center">TABLE 24</entry></row></thead><tbody valign="top"><row><entry /></row><row><entry namest="1" nameend="1" align="center" rowsep="1" /></row><row><entry>50 × 50 Matrix Partitioning</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="4"><colspec colname="1" colwidth="35PT" align="center" /><colspec colname="2" colwidth="42PT" align="center" /><colspec colname="3" colwidth="63PT" align="center" /><colspec colname="4" colwidth="77PT" align="center" /><tbody valign="top"><row><entry>Processing</entry><entry>Row</entry><entry>Comlumn Indices of</entry><entry>Number of Elements</entry></row><row><entry>Node</entry><entry>Index</entry><entry>Elements Computed</entry><entry>Computed Per Row</entry></row><row><entry namest="1" nameend="4" align="center" rowsep="1" /></row><row><entry>P<sub>1</sub></entry><entry> 1</entry><entry>1 → 50</entry><entry>50</entry></row><row><entry /><entry> 2</entry><entry>1 → 50</entry><entry>50</entry></row><row><entry /><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry></row><row><entry /><entry> 7</entry><entry>1 → 50</entry><entry>50</entry></row><row><entry /><entry> 8</entry><entry>1 → 8</entry><entry>8</entry></row><row><entry namest="1" nameend="4" align="center" rowsep="1" /></row><row><entry namest="1" nameend="4" align="left"><!--footnotes removed--></entry></row><row><entry>P<sub>2</sub></entry><entry> 8</entry><entry>9 → 50</entry><entry>42</entry></row><row><entry /><entry> 9</entry><entry>1 → 50</entry><entry>50</entry></row><row><entry /><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry></row><row><entry /><entry>14</entry><entry>1 → 50</entry><entry>50</entry></row><row><entry /><entry>15</entry><entry>1 → 15</entry><entry>15</entry></row><row><entry namest="1" nameend="4" align="center" rowsep="1" /></row><row><entry>P<sub>i</sub></entry><entry>7i-6</entry><entry>7i − → 50</entry><entry>51 − (7i − 5)</entry></row><row><entry /><entry>7i-5</entry><entry>1 → 50</entry><entry>50</entry></row><row><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry></row><row><entry /><entry>7i</entry><entry>1 → 50</entry><entry>50</entry></row><row><entry /><entry>7i + 1</entry><entry>1 → 7i + 1</entry><entry>7i + 1</entry></row><row><entry>P<sub>7</sub></entry><entry>43</entry><entry>44 → 50 </entry><entry> 7</entry></row><row><entry /><entry>44</entry><entry>1 → 50</entry><entry>50</entry></row><row><entry /><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry></row><row><entry /><entry>49</entry><entry>1 → 50</entry><entry>50</entry></row><row><entry /><entry>50</entry><entry>1 → 50</entry><entry>50</entry></row><row><entry namest="1" nameend="4" align="center" rowsep="1" /></row><row><entry namest="1" nameend="4" align="left"><!--footnotes removed--></entry></row></tbody></tgroup></table></tables>
[0432] Processing node P<sub>1 </sub>computes rows 1 through 7 for all columns, then computes row 8 for columns 1 through 8 only, for a total of 358 elements computed. P<sub>1 </sub>uses the first <b>8</b> rows of matrix A and all of matrix B to compute its result.
[0433] Processing node P<sub>2 </sub>computes row 8 for columns 9 through 50, rows 9 through 14 for all columns, and row 15 for columns 1 through 15, totaling 357 elements computed. P<sub>2 </sub>uses rows 8 through 15 of matrix A and all of matrix B to compute its result.
[0434] Processing node P<sub>i</sub>, for 1<i≦7, computes elements 357(i−1)+2 through 357i+1. Dividing by column width 50 produces a row index range of 7i−6 through 7i+1. The first row, 7i−6, is computed for columns 7i−5 through 50. The last row, 7i+1, is computed for columns 1 through 7i+1. The remaining rows are computed for all columns. The total number of elements computed equals 51−(7i−5)+6(50)+7i+1=56−7i+300+7i+1=357. P<sub>i </sub>uses rows 7i−6 through 7i+1 of matrix A and all of matrix B to compute its result.
[0435] The home node performs the matrix data partitioning. Messages describing the command (in this case a matrix multiply) and data partitioning are sent out to the processing nodes. Once the nodes have received their command messages, they wait for the home node to send the input matrix data. The data is broadcast such that all nodes in the cascade receive it at the same time. Each node receives both of the input matrices, which is much more efficient than sending each individual node a separate message with just its piece of input data.
[0436] Once a node receives the input matrices, it can proceed with computing the product independent of the other nodes. When the matrix multiply results are ready, they are accumulated up to the home node and merged into the final result. At this point the process is complete. As this example shows, matrix data distribution applied to the HC accommodates arbitrary sized matrices in a simple and efficient manner.
[0437] Decomposition of Parallel Processing Data for Two-Dimensional Convolution Using Fast Fourier Transforms
[0438] This section describes how a two-dimensional convolution using Fast Fourier Transforms may be partitioned on a HC. The partitioning can handle an arbitrary number of parallel processing nodes and arbitrarily large images.
[0439] Consider an image containing M rows and N columns of pixel values and a smaller kernel image containing J rows and K columns (J≦M, K≦N). An efficient method for performing 2D convolution on the image and kernel involves the use of 2D Fast Fourier Transforms. First, the kernel is padded to match the size of the image. Then, a 2D FFT is performed on the image and kernel separately. The results are multiplied, per element, and then an inverse 2D FFT is applied to the product. The final result is equivalent to computing the 2D convolution directly.
[0440] An arbitrary input image and kernel are shown below in Tables 25 and 26, respectively, as arrays of pixels. Each parallel processing node receives a copy of the entire kernel. The image is evenly distributed by rows or columns over the processing nodes; the data partitioning here is described by row distribution. <tables id="TABLE-US-00028" num="28"><table frame="none" colsep="0" rowsep="0"><tgroup align="left" colsep="0" rowsep="0" cols="1"><colspec colname="1" colwidth="217PT" align="center" /><thead><row><entry namest="1" nameend="1" align="center">TABLE 25</entry></row></thead><tbody valign="top"><row><entry /></row><row><entry namest="1" nameend="1" align="center" rowsep="1" /></row><row><entry>Image Matrix</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="2"><colspec colname="OFFSET" colwidth="35PT" align="left" /><colspec colname="1" colwidth="182PT" align="center" /><tbody valign="top"><row><entry /><entry>columns</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="7"><colspec colname="OFFSET" colwidth="14PT" align="left" /><colspec colname="1" colwidth="21PT" align="center" /><colspec colname="2" colwidth="49PT" align="center" /><colspec colname="3" colwidth="21PT" align="center" /><colspec colname="4" colwidth="49PT" align="center" /><colspec colname="5" colwidth="14PT" align="center" /><colspec colname="6" colwidth="49PT" align="center" /><tbody valign="top"><row><entry /><entry>rows</entry><entry>1</entry><entry>2</entry><entry>3</entry><entry>. . .</entry><entry>N</entry></row><row><entry /><entry namest="OFFSET" nameend="6" align="center" rowsep="1" /></row><row><entry /><entry>1</entry><entry>1,1</entry><entry>1,2</entry><entry>1,3</entry><entry>. . .</entry><entry>1,N</entry></row><row><entry /><entry>2</entry><entry>2,1</entry><entry>2,2</entry><entry>2,3</entry><entry>. . .</entry><entry>2,N</entry></row><row><entry /><entry>3</entry><entry>3,1</entry><entry>3,2</entry><entry>3,3</entry><entry>. . .</entry><entry>3,N</entry></row><row><entry /><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry></row><row><entry /><entry>M</entry><entry>M,1</entry><entry>M,2</entry><entry>M,3</entry><entry>. . .</entry><entry>M,N</entry></row><row><entry /><entry namest="OFFSET" nameend="6" align="center" rowsep="1" /></row></tbody></tgroup></table></tables>
[0441]<tables id="TABLE-US-00029" num="29"><table frame="none" colsep="0" rowsep="0"><tgroup align="left" colsep="0" rowsep="0" cols="1"><colspec colname="1" colwidth="217PT" align="center" /><thead><row><entry namest="1" nameend="1" align="center">TABLE 26</entry></row></thead><tbody valign="top"><row><entry /></row><row><entry namest="1" nameend="1" align="center" rowsep="1" /></row><row><entry>Kernel Matrix</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="3"><colspec colname="OFFSET" colwidth="70PT" align="left" /><colspec colname="1" colwidth="126PT" align="center" /><colspec colname="2" colwidth="21PT" align="center" /><tbody valign="top"><row><entry /><entry>columns</entry><entry /></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="5"><colspec colname="1" colwidth="70PT" align="center" /><colspec colname="2" colwidth="14PT" align="center" /><colspec colname="3" colwidth="56PT" align="center" /><colspec colname="4" colwidth="14PT" align="center" /><colspec colname="5" colwidth="63PT" align="center" /><tbody valign="top"><row><entry>rows</entry><entry>1</entry><entry>2</entry><entry>. . .</entry><entry>K</entry></row><row><entry namest="1" nameend="5" align="center" rowsep="1" /></row><row><entry>1</entry><entry>1,1</entry><entry>1,2</entry><entry>. . .</entry><entry>1,K</entry></row><row><entry>2</entry><entry>2,1</entry><entry>2,2</entry><entry>. . .</entry><entry>2,K</entry></row><row><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry></row><row><entry>J</entry><entry>J,1</entry><entry>J,2</entry><entry>. . .</entry><entry>J,K</entry></row><row><entry namest="1" nameend="5" align="center" rowsep="1" /></row></tbody></tgroup></table></tables>
[0442] The M rows of the input image are evenly distributed over P parallel processing nodes. The rows assigned to a node are defined by the starting row, referred to as row index, IR<sub>i</sub>, and a row count, M<sub>i</sub>. Rows are not split across nodes, so the row count is constrained to whole -numbers. Assuming M>P, M rows is divided by P nodes to obtain an integer quotient Q and a remainder R. Q contiguous rows are assigned to each of the P nodes. Any remainder R is distributed, one row per node, to each of the first R nodes. Thus the first R nodes are assigned a row count of M<sub>i</sub>=Q+1 each and the remaining P−R nodes are assigned a row count of M<sub>i</sub>=Q. The total number of rows assigned equals R(Q+1)+(P−R)Q=PQ+R=M. The row index, IR<sub>i</sub>, for the first R nodes equals (i−1)(Q+1) +1. For the remaining P−R nodes, IR<sub>i </sub>equals (i−1)Q+R+1.
[0443] This described approach of this section maintains the node computational load balanced to within one row. The mapping of rows to nodes is illustrated below in Table 27. <tables id="TABLE-US-00030" num="30"><table frame="none" colsep="0" rowsep="0"><tgroup align="left" colsep="0" rowsep="0" cols="1"><colspec colname="1" colwidth="217PT" align="center" /><thead><row><entry namest="1" nameend="1" align="center">TABLE 27</entry></row></thead><tbody valign="top"><row><entry /></row><row><entry namest="1" nameend="1" align="center" rowsep="1" /></row><row><entry>Row Data Partitioning</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="6"><colspec colname="1" colwidth="28PT" align="center" /><colspec colname="2" colwidth="42PT" align="right" /><colspec colname="3" colwidth="49PT" align="right" /><colspec colname="4" colwidth="42PT" align="right" /><colspec colname="5" colwidth="14PT" align="center" /><colspec colname="6" colwidth="42PT" align="right" /><tbody valign="top"><row><entry>Node</entry><entry>1</entry><entry>2</entry><entry>3</entry><entry>. . .</entry><entry>N</entry></row><row><entry namest="1" nameend="6" align="center" rowsep="1" /></row><row><entry>P<sub>1</sub></entry><entry>IR<sub>1</sub>,1</entry><entry>IR<sub>1</sub>,2</entry><entry>IR<sub>1</sub>,3</entry><entry>. . .</entry><entry>IR<sub>1</sub>,N</entry></row><row><entry /><entry>IR<sub>1</sub> + 1,1</entry><entry>IR<sub>1</sub> + 1,2</entry><entry>IR<sub>1</sub> + 1,3</entry><entry>. . .</entry><entry>IR<sub>1</sub> + 1,N</entry></row><row><entry /><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry></row><row><entry /><entry>IR<sub>1</sub> + M<sub>1</sub>,1</entry><entry>IR<sub>1</sub> + M<sub>1</sub>,2</entry><entry>IR<sub>1</sub> + M<sub>1</sub>,3</entry><entry>. . .</entry><entry>IR<sub>1</sub> + M<sub>1</sub>,N</entry></row><row><entry>P<sub>2</sub></entry><entry>IR<sub>2</sub>,1</entry><entry>IR<sub>2</sub>,2</entry><entry>IR<sub>2</sub>,3</entry><entry>. . .</entry><entry>IR<sub>2</sub>,N</entry></row><row><entry /><entry>IR<sub>2</sub> + 1,1</entry><entry>IR<sub>2</sub> + 1,2</entry><entry>IR<sub>2</sub> + 1,3</entry><entry>. . .</entry><entry>IR<sub>2</sub> + 1,N</entry></row><row><entry /><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry></row><row><entry /><entry>IR<sub>2</sub> + M<sub>2</sub>,1</entry><entry>IR<sub>2</sub> + M<sub>2</sub>,2</entry><entry>IR<sub>2</sub> + M<sub>2</sub>,3</entry><entry>. . .</entry><entry>IR<sub>2</sub> + M<sub>2</sub>,N</entry></row><row><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry></row><row><entry>P<sub>P</sub></entry><entry>IR<sub>P</sub>,1</entry><entry>IR<sub>P</sub>,2</entry><entry>IR<sub>P</sub>,3</entry><entry>. . .</entry><entry>IR<sub>P</sub>,N</entry></row><row><entry /><entry>IR<sub>P</sub> + 1,1</entry><entry>IR<sub>P</sub> + 1,2</entry><entry>IR<sub>P</sub> + 1,3</entry><entry>. . .</entry><entry>IR<sub>P</sub> + 1,N</entry></row><row><entry /><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry></row><row><entry /><entry>IR<sub>P</sub> + M<sub>P</sub>,1</entry><entry>IR<sub>P</sub> + M<sub>P</sub>,2</entry><entry>IR<sub>P</sub> + M<sub>P</sub>,3</entry><entry>. . .</entry><entry>IR<sub>P</sub> + M<sub>P</sub>,N</entry></row><row><entry namest="1" nameend="6" align="center" rowsep="1" /></row><row><entry namest="1" nameend="6" align="left"><!--footnotes removed--></entry></row><row><entry namest="1" nameend="6" align="left"><!--footnotes removed--></entry></row><row><entry namest="1" nameend="6" align="left"><!--footnotes removed--></entry></row><row><entry namest="1" nameend="6" align="left"><!--footnotes removed--></entry></row></tbody></tgroup></table></tables>
[0444] Consider the distribution of an image consisting of 1024 rows and 700 columns on a 7-node HC. In this case, the integer quotient Q is 1024÷7=146 and the row remainder R is 2. Processing nodes P<sub>1 </sub>and P<sub>2 </sub>are assigned 147 rows each; nodes P<sub>3 </sub>through P<sub>7 </sub>are assigned 146 each. The row index for P<sub>1 </sub>is 1, the row index for P<sub>2 </sub>is 148, etc., up to P<sub>7</sub>, which has a row index of 879.
[0445] At the start of a HC computation, messages describing the command (in this case a 2D convolution) are sent out to the processing nodes. Once the nodes have received their command messages, they wait for the home node to send the input data. The data is broadcast such that all nodes in the cascade receive it at the same time. Each node receives the entire dataset, which is more efficient than sending each individual node a separate message with just its piece of data.
[0446] Once a node receives the input data, it proceeds with performing the 2D convolution in its assigned rows independent of the other nodes. When the individual results are ready, they are accumulated up to the home node and merged into the final result. At this point the process is complete.
[0447] As this example shows, the data distribution as applied to the HC accommodates arbitrary sized images and kernels in a simple and efficient manner.
[0448] Decomposition of a Linear System of Equations for Parallel Processing
[0449] This section describes a solution of a linear system of equations Ax=b partitioned on a HC. The following partitioning method handles an arbitrary number of parallel processing nodes and arbitrarily large systems of equations. Consider the solution of the equation Ax=b, or, more specifically: <maths id="MATH-US-00006" num="6"><math overflow="scroll"><mtable><mtr><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mrow><msub><mi>a</mi><mrow><mn>1</mn><mo>,</mo><mn>1</mn></mrow></msub><mo></mo><msub><mi>x</mi><mn>1</mn></msub></mrow></mtd><mtd><mo>+</mo></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mrow><msub><mi>a</mi><mrow><mn>1</mn><mo>,</mo><mn>2</mn></mrow></msub><mo></mo><msub><mi>x</mi><mn>2</mn></msub></mrow></mtd><mtd><mo>+</mo></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mrow><msub><mi>a</mi><mrow><mn>1</mn><mo>,</mo><mn>3</mn></mrow></msub><mo></mo><msub><mi>x</mi><mn>3</mn></msub></mrow></mtd><mtd><mo>+</mo></mtd><mtd><mi>…</mi></mtd></mtr><mtr><mtd><mo>+</mo></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mrow><msub><mi>a</mi><mrow><mn>1</mn><mo>,</mo><mi>N</mi></mrow></msub><mo></mo><msub><mi>x</mi><mi>N</mi></msub></mrow></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mrow><mo>=</mo><msub><mi>b</mi><mn>1</mn></msub></mrow></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd></mtr><mtr><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mrow><msub><mi>a</mi><mrow><mn>2</mn><mo>,</mo><mn>1</mn></mrow></msub><mo></mo><msub><mi>x</mi><mn>1</mn></msub></mrow></mtd><mtd><mrow><mo>+</mo><mstyle><mtext> </mtext></mstyle></mrow></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mrow><msub><mi>a</mi><mrow><mn>2</mn><mo>,</mo><mn>2</mn></mrow></msub><mo></mo><msub><mi>x</mi><mn>2</mn></msub></mrow></mtd><mtd><mo>+</mo></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mrow><msub><mi>a</mi><mrow><mn>2</mn><mo>,</mo><mn>3</mn></mrow></msub><mo></mo><msub><mi>x</mi><mn>3</mn></msub></mrow></mtd><mtd><mo>+</mo></mtd></mtr><mtr><mtd><mi>…</mi></mtd><mtd><mo>+</mo></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mrow><msub><mi>a</mi><mrow><mn>2</mn><mo>,</mo><mi>N</mi></mrow></msub><mo></mo><msub><mi>x</mi><mi>N</mi></msub></mrow></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mrow><mo>=</mo><msub><mi>b</mi><mn>2</mn></msub></mrow></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd></mtr><mtr><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mrow><msub><mi>a</mi><mrow><mn>3</mn><mo>,</mo><mn>1</mn></mrow></msub><mo></mo><msub><mi>x</mi><mn>1</mn></msub></mrow></mtd><mtd><mo>+</mo></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mrow><msub><mi>a</mi><mrow><mn>3</mn><mo>,</mo><mn>2</mn></mrow></msub><mo></mo><msub><mi>x</mi><mn>2</mn></msub></mrow></mtd><mtd><mo>+</mo></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mrow><msub><mi>a</mi><mrow><mn>3</mn><mo>,</mo><mn>3</mn></mrow></msub><mo></mo><msub><mi>x</mi><mn>3</mn></msub></mrow></mtd><mtd><mo>+</mo></mtd></mtr><mtr><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mi>…</mi></mtd><mtd><mo>+</mo></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mrow><msub><mi>a</mi><mrow><mn>3</mn><mo>,</mo><mi>N</mi></mrow></msub><mo></mo><msub><mi>x</mi><mi>N</mi></msub></mrow></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mrow><mo>=</mo><msub><mi>b</mi><mn>3</mn></msub></mrow></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd></mtr><mtr><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mi>…</mi></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd></mtr><mtr><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mrow><msub><mi>a</mi><mrow><mi>N</mi><mo>,</mo><mn>1</mn></mrow></msub><mo></mo><msub><mi>x</mi><mn>1</mn></msub></mrow></mtd><mtd><mo>+</mo></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mrow><msub><mi>a</mi><mrow><mi>N</mi><mo>,</mo><mn>2</mn></mrow></msub><mo></mo><msub><mi>x</mi><mn>2</mn></msub></mrow></mtd><mtd><mo>+</mo></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mrow><msub><mi>a</mi><mrow><mi>N</mi><mo>,</mo><mn>3</mn></mrow></msub><mo></mo><msub><mi>x</mi><mn>3</mn></msub></mrow></mtd><mtd><mo>+</mo></mtd></mtr><mtr><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mi>…</mi></mtd><mtd><mo>+</mo></mtd><mtd><mrow><msub><mi>a</mi><mrow><mi>N</mi><mo>,</mo><mi>N</mi></mrow></msub><mo></mo><msub><mi>x</mi><mi>N</mi></msub></mrow></mtd><mtd><mrow><mo>=</mo><msub><mi>b</mi><mi>N</mi></msub></mrow></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd><mtd><mstyle><mtext> </mtext></mstyle></mtd></mtr></mtable></math><img file="US20030195938A1-20031016-M00006.TIF" id="EMI-M00006" he="97.9209" wi="216.027" img-format="tif" img-content="mf" /><attachments><attachment idref="MATHEMATICA-00006" attachment-type="nb" file="US20030195938A1-20031016-M00006.NB" /></attachments></maths>
[0450] where A is an N×N matrix of coefficients, b is an N×1 column vector, and x is the N×1 solution being sought.
[0451] To solve for x, an LU factorization is performed on A. LU factorization results in a lower triangular matrix L and upper triangular matrix U such that A=L×U. Substituting forA in the original equation, LUx=b or Ux=L<sup>−1</sup>b, and x may also be solved.
[0452] First, A and b are combined into a single input matrix with N rows and N+1 columns, as shown in Table 25. Data partitioning evenly distributes the rows of the input matrix over the parallel processing nodes. <tables id="TABLE-US-00031" num="31"><table frame="none" colsep="0" rowsep="0"><tgroup align="left" colsep="0" rowsep="0" cols="1"><colspec colname="1" colwidth="217PT" align="center" /><thead><row><entry namest="1" nameend="1" align="center">TABLE 28</entry></row><row><entry /></row><row><entry namest="1" nameend="1" align="center" rowsep="1" /></row><row><entry>Input Matrix</entry></row><row><entry namest="1" nameend="1" align="center" rowsep="1" /></row></thead><tbody valign="top"><row><entry /></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="9"><colspec colname="OFFSET" colwidth="14PT" align="left" /><colspec colname="1" colwidth="14PT" align="center" /><colspec colname="2" colwidth="35PT" align="center" /><colspec colname="3" colwidth="14PT" align="center" /><colspec colname="4" colwidth="42PT" align="center" /><colspec colname="5" colwidth="21PT" align="center" /><colspec colname="6" colwidth="28PT" align="center" /><colspec colname="7" colwidth="14PT" align="center" /><colspec colname="8" colwidth="35PT" align="center" /><tbody valign="top"><row><entry /><entry /><entry /><entry /><entry /><entry /><entry>a<sub>1,</sub></entry><entry /><entry /></row><row><entry /><entry /><entry>a<sub>1,1</sub></entry><entry>a<sub>1,2</sub></entry><entry>a<sub>1,3</sub></entry><entry>. . . </entry><entry /><entry>b<sub>1</sub></entry></row><row><entry /><entry /><entry /><entry /><entry /><entry /><entry>N</entry></row><row><entry /><entry /><entry /><entry /><entry /><entry /><entry>a<sub>2,</sub></entry></row><row><entry /><entry /><entry>a<sub>2,1</sub></entry><entry>a<sub>2,2</sub></entry><entry>a<sub>2,3</sub></entry><entry /><entry /><entry>b<sub>2</sub></entry></row><row><entry /><entry /><entry /><entry /><entry /><entry /><entry>N</entry></row><row><entry /><entry> {open oversize bracket} </entry><entry /><entry /><entry /><entry /><entry> a<sub>3,</sub></entry><entry /><entry> {close oversize bracket} </entry></row><row><entry /><entry /><entry>a<sub>3,1</sub></entry><entry>a<sub>3,2</sub></entry><entry>a<sub>3,3</sub></entry><entry /><entry /><entry>b<sub>3</sub></entry></row><row><entry /><entry /><entry /><entry /><entry /><entry /><entry>N</entry></row><row><entry /><entry /><entry>. . . </entry><entry /><entry>. . . </entry></row><row><entry /><entry /><entry>a<sub>N,</sub></entry><entry>a<sub>N,</sub></entry><entry>a<sub>N,</sub></entry><entry /><entry>a<sub>N,</sub></entry></row><row><entry /><entry /><entry /><entry /><entry /><entry /><entry /><entry>b<sub>N</sub></entry></row><row><entry /><entry /><entry>1</entry><entry>2</entry><entry>3</entry><entry /><entry>N</entry></row><row><entry /><entry namest="OFFSET" nameend="8" align="center" rowsep="1" /></row></tbody></tgroup></table></tables>
[0453] The N rows of the input matrix are evenly distributed over P parallel processing nodes. The rows assigned to a node are defined by the starting row, referred to as row index, IR<sub>i</sub>, and a row count, N<sub>i</sub>. Rows are not split across nodes, so the row count is constrained to whole numbers. Assuming N>P, N rows is divided by P nodes to obtain an integer quotient Q and a remainder R. Q contiguous rows are assigned to each of the P nodes. Any remainder R is distributed, one row per node, to each of the first R nodes. Thus the first R nodes is assigned a row count of N<sub>i</sub>=Q+1 each and the remaining P−R nodes is assigned a row count of N<sub>i</sub>=Q. The total number of rows assigned equals R(Q+1)+(P−R)Q=PQ+R=N. The row index, IR<sub>i</sub>, for the first R nodes equals (i−1)(Q+1)+1. For the remaining P−R nodes, IR<sub>i </sub>equals (i−1)Q+R+1.
[0454] This described approach of this section again maintains node computational load balanced to within one row. The mapping of rows to nodes is illustrated in Table 29. <tables id="TABLE-US-00032" num="32"><table frame="none" colsep="0" rowsep="0"><tgroup align="left" colsep="0" rowsep="0" cols="1"><colspec colname="1" colwidth="217PT" align="center" /><thead><row><entry namest="1" nameend="1" align="center">TABLE 29</entry></row></thead><tbody valign="top"><row><entry /></row><row><entry namest="1" nameend="1" align="center" rowsep="1" /></row><row><entry>Mapping of Rows to Nodes</entry></row></tbody></tgroup><tgroup align="left" colsep="0" rowsep="0" cols="6"><colspec colname="1" colwidth="28PT" align="center" /><colspec colname="2" colwidth="42PT" align="right" /><colspec colname="3" colwidth="42PT" align="right" /><colspec colname="4" colwidth="42PT" align="right" /><colspec colname="5" colwidth="14PT" align="center" /><colspec colname="6" colwidth="49PT" align="right" /><tbody valign="top"><row><entry>Node</entry><entry>1</entry><entry>2</entry><entry>3</entry><entry>. . .</entry><entry>N + 1</entry></row><row><entry namest="1" nameend="6" align="center" rowsep="1" /></row><row><entry>P<sub>1</sub></entry><entry>IR<sub>1</sub>,1</entry><entry>IR<sub>1</sub>,2</entry><entry>IR<sub>1</sub>,3</entry><entry>. . .</entry><entry>IR<sub>1</sub>,N + 1</entry></row><row><entry /><entry>IR<sub>1</sub> + 1,1</entry><entry>IR<sub>1</sub> + 1,2</entry><entry>IR<sub>1</sub> + 1,3</entry><entry>. . .</entry><entry>IR<sub>1</sub> + 1,N + 1</entry></row><row><entry /><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry></row><row><entry /><entry>IR<sub>1</sub> + N<sub>1</sub>,1</entry><entry>IR<sub>1</sub> + N<sub>1</sub>,2</entry><entry>IR<sub>1</sub> + N<sub>1</sub>,3</entry><entry>. . .</entry><entry>IR<sub>1</sub> + N<sub>1</sub>,N + 1</entry></row><row><entry>P<sub>2</sub></entry><entry>IR<sub>2</sub>,1</entry><entry>IR<sub>2</sub>,2</entry><entry>IR<sub>2</sub>,3</entry><entry>. . .</entry><entry>IR<sub>2</sub>,N + 1</entry></row><row><entry /><entry>IR<sub>2</sub> + 1,1</entry><entry>IR<sub>2</sub> + 1,2</entry><entry>IR<sub>2</sub> + 1,3</entry><entry>. . .</entry><entry>IR<sub>2</sub> + 1,N + 1</entry></row><row><entry /><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry></row><row><entry /><entry>IR<sub>2</sub> + N<sub>2</sub>,1</entry><entry>IR<sub>2</sub> + N<sub>2</sub>,2</entry><entry>IR<sub>2</sub> + N<sub>2</sub>,3</entry><entry>. . .</entry><entry>IR<sub>2</sub> + N<sub>2</sub>,N + 1</entry></row><row><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry></row><row><entry>P<sub>P</sub></entry><entry>IR<sub>P</sub>,1</entry><entry>IR<sub>P</sub>,2</entry><entry>IR<sub>P</sub>,3</entry><entry>. . .</entry><entry>IR<sub>P</sub>,N + 1</entry></row><row><entry /><entry>IR<sub>P</sub> + 1,1</entry><entry>IR<sub>P</sub> + 1,2</entry><entry>IR<sub>P</sub> + 1,3</entry><entry>. . .</entry><entry>IR<sub>P</sub> + 1,N + 1</entry></row><row><entry /><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry><entry>. . .</entry></row><row><entry /><entry>IR<sub>P</sub> + N<sub>P</sub>,1</entry><entry>IR<sub>P</sub> + N<sub>P</sub>,2</entry><entry>IR<sub>P</sub> + N<sub>P</sub>,3</entry><entry>. . .</entry><entry>IR<sub>P</sub> + N<sub>P</sub>,N + 1</entry></row><row><entry namest="1" nameend="6" align="center" rowsep="1" /></row><row><entry namest="1" nameend="6" align="left"><!--footnotes removed--></entry></row><row><entry namest="1" nameend="6" align="left"><!--footnotes removed--></entry></row><row><entry namest="1" nameend="6" align="left"><!--footnotes removed--></entry></row></tbody></tgroup></table></tables>
[0455] Consider the distribution of an input matrix of 100 rows and 101 columns on a 7-node HC. In this case, the integer quotient Q is 10100÷7=1442 and the row remainder R is 6. Processing nodes P, through P<sub>6 </sub>are assigned 1443 rows each; node P<sub>7 </sub>is assigned 1442 rows. The row index for P<sub>1 </sub>is 1, the row index for P<sub>2 </sub>is 1444, etc., up to P<sub>7</sub>, which has a row index of 8659.
[0456] The home node performs the matrix data partitioning of this section. Messages describing the command (in this case a linear solve) and data partitioning are sent out to the processing nodes. Once the nodes have received their command messages, they wait for the home node to send the input matrix data. The data may be broadcast such that all nodes in the cascade receive it at the same time. Each node may thus receive the entire dataset, which may be more efficient than sending each individual node a separate message with just its piece of data.
[0457] Once a node receives the input matrix, it proceeds with solving the linear system in its assigned rows, and independent of the other nodes except when sharing data to determine the pivot row at each iteration. When the individual results are ready, they are accumulated up to the home node and merged into the final result. At this point the process is complete.
[0458] As this example shows, the data distribution applied to the HC accommodates arbitrary sized linear systems of equations in a simple and efficient manner.
[0459] Those skilled in the art will appreciate that variations from the specified embodiments disclosed above are contemplated herein. The description should not be restricted to the above embodiments, but should be measured by the following claims.
Contents6
60 sheets
Sheet 1 Sheet 2 Sheet 3 Sheet 4 Sheet 5 Sheet 6 Sheet 7 Sheet 8 Sheet 9 Sheet 10 Sheet 11 Sheet 12 Sheet 13 Sheet 14 Sheet 15 Sheet 16 Sheet 17 Sheet 18 Sheet 19 Sheet 20 Sheet 21 Sheet 22 Sheet 23 Sheet 24 Sheet 25 Sheet 26 Sheet 27 Sheet 28 Sheet 29 Sheet 30 Sheet 31 Sheet 32 Sheet 33 Sheet 34 Sheet 35 Sheet 36 Sheet 37 Sheet 38 Sheet 39 Sheet 40 Sheet 41 Sheet 42 Sheet 43 Sheet 44 Sheet 45 Sheet 46 Sheet 47 Sheet 48 Sheet 49 Sheet 50 Sheet 51 Sheet 52 Sheet 53 Sheet 54 Sheet 55 Sheet 56 Sheet 57 Sheet 58 Sheet 59 Sheet 60
Every citation, both ways
| Document | Relation | Office | Cited during |
|---|---|---|---|
| US10223763B2 | Cited by | United States of America | Search report |
| US2015163289A1 | Cited by | United States of America | Pre-grant |
| CN108241508A | Cited by | China | Search report |
| US2010122237A1 | Cited by | United States of America | Pre-grant |
| US10112606B2 | Cited by | United States of America | Search report |
| EP2616932A4 | Cited by | European Patent Office (EPO) | Search report |
| US2010251259A1 | Cited by | United States of America | Pre-grant |
| US7836118B1 | Cited by | United States of America | Applicant |
| US2017210376A1 | Cited by | United States of America | Pre-grant |
| US7506134B1 | Cited by | United States of America | Search report |
| US8843879B2 | Cited by | United States of America | Search report |
| US2010049941A1 | Cited by | United States of America | Pre-grant |
| US8140612B2 | Cited by | United States of America | Applicant |
| EP2195747A4 | Cited by | European Patent Office (EPO) | Search report |
| US2009222543A1 | Cited by | United States of America | Pre-grant |
| US9104503B2 | Cited by | United States of America | Applicant |
| US10148425B2 | Cited by | United States of America | Search report |
| US2008079724A1 | Cited by | United States of America | Pre-grant |
| US11563621B2 | Cited by | United States of America | Search report |
| US2005198634A1 | Cited by | United States of America | Pre-grant |
| US2017064333A1 | Cited by | United States of America | Pre-grant |
| US8832177B1 | Cited by | United States of America | Applicant |
| EP2195747A1 | Cited by | European Patent Office (EPO) | Search report |
| US8108512B2 | Cited by | United States of America | Search report |
| US9244729B1 | Cited by | United States of America | Applicant |
| WO2005111843A2 | Cited by | World Intellectual Property Organization (WIPO) | International search |
| FR3050848A1 | Cited by | France | Search report |
| US9876876B2 | Cited by | United States of America | Applicant |
| US8769034B2 | Cited by | United States of America | Search report |
| US2007288935A1 | Cited by | United States of America | Pre-grant |
| US2006130063A1 | Cited by | United States of America | Pre-grant |
| US2006187958A1 | Cited by | United States of America | Pre-grant |
| US2010325388A1 | Cited by | United States of America | Pre-grant |
| US10503557B2 | Cited by | United States of America | Search report |
| US7958194B2 | Cited by | United States of America | Search report |
| US2006187928A1 | Cited by | United States of America | Pre-grant |
| US2018367584A1 | Cited by | United States of America | Search report |
| WO2010056591A1 | Cited by | World Intellectual Property Organization (WIPO) | International search |
| WO2005111843A3 | Cited by | World Intellectual Property Organization (WIPO) | International search |
| US2017178281A1 | Cited by | United States of America | Pre-grant |
| US11321136B2 | Cited by | United States of America | Search report |
| US10318260B2 | Cited by | United States of America | Search report |
| US8040903B2 | Cited by | United States of America | Search report |
| WO2012006285A1 | Cited by | World Intellectual Property Organization (WIPO) | International search |
| US2009077483A9 | Cited by | United States of America | Pre-grant |
| US8201142B2 | Cited by | United States of America | Applicant |
| US7996458B2 | Cited by | United States of America | Search report |
| US8102855B2 | Cited by | United States of America | Applicant |
| US7474658B2 | Cited by | United States of America | Search report |
| EP2616932A2 | Cited by | European Patent Office (EPO) | Search report |
| US10333768B2 | Cited by | United States of America | Search report |
| US7844959B2 | Cited by | United States of America | Applicant |
| US9933838B2 | Cited by | United States of America | Search report |
| CN107506932A | Cited by | China | Search report |
| US2005044147A1 | Cited by | United States of America | Pre-grant |
| US9313087B2 | Cited by | United States of America | Applicant |
| US2012265835A1 | Cited by | United States of America | Pre-grant |
| US9424076B1 | Cited by | United States of America | Applicant |
| US2008082933A1 | Cited by | United States of America | Pre-grant |
| EP3379414A1 | Cited by | European Patent Office (EPO) | Search report |
| US2008184255A1 | Cited by | United States of America | Pre-grant |
| US10135944B2 | Cited by | United States of America | Applicant |
| US10673983B2 | Cited by | United States of America | Applicant |
| US8402080B2 | Cited by | United States of America | Applicant |
| EP3343370A1 | Cited by | European Patent Office (EPO) | Search report |
| US2008082644A1 | Cited by | United States of America | Pre-grant |
| US8161490B2 | Cited by | United States of America | Search report |
| US8082289B2 | Cited by | United States of America | Applicant |
| US2016085291A1 | Cited by | United States of America | Pre-grant |
| US2016006643A1 | Cited by | United States of America | Pre-grant |
| WO2009045526A1 | Cited by | World Intellectual Property Organization (WIPO) | Applicant |
| US8533722B2 | Cited by | United States of America | Search report |
| US2008148244A1 | Cited by | United States of America | Pre-grant |
| US8589468B2 | Cited by | United States of America | Applicant |
| US7912889B1 | Cited by | United States of America | Applicant |
| US8924654B1 | Cited by | United States of America | Search report |
| US9609082B2 | Cited by | United States of America | Applicant |
| US7394817B2 | Cited by | United States of America | Search report |
| US10516726B2 | Cited by | United States of America | Search report |
| US2007255782A1 | Cited by | United States of America | Pre-grant |
| US2006143608A1 | Cited by | United States of America | Pre-grant |
| US8813053B2 | Cited by | United States of America | Applicant |
| US10216692B2 | Cited by | United States of America | Applicant |
| US9626329B2 | Cited by | United States of America | Applicant |
| US8863145B2 | Cited by | United States of America | Applicant |
| US9860192B2 | Cited by | United States of America | Applicant |
| US2011320530A1 | Cited by | United States of America | Pre-grant |
| US7886294B2 | Cited by | United States of America | Applicant |
| US8533717B2 | Cited by | United States of America | Search report |
| US2010064033A1 | Cited by | United States of America | Pre-grant |
| US10026161B2 | Cited by | United States of America | Search report |
| US11570034B2 | Cited by | United States of America | Search report |
| US7689989B2 | Cited by | United States of America | Applicant |
| US2015163289A1 | Cited by | United States of America | Search report |
| US10630737B2 | Cited by | United States of America | Search report |
| JP2010541100A | Cited by | Japan | Search report |
| US8402083B2 | Cited by | United States of America | Applicant |
| US8302076B2 | Cited by | United States of America | Applicant |
| US2008271036A1 | Cited by | United States of America | Pre-grant |
| US2006143359A1 | Cited by | United States of America | Pre-grant |
39 members in 6 offices
Priority claims8
| Document | Office | Kind | Date |
|---|---|---|---|
| 60302000 | United States of America | A | |
| 34732502 | United States of America | P | |
| 34052403 | United States of America | A | |
| 09603020 | – | – | – |
| 60347325 | – | – | – |
| US20000603020 | – | – | – |
| US20020347325P | – | – | – |
| US20030340524 | – | – | – |
Members39
| Document | Office | Kind | |
|---|---|---|---|
| CA2378088A1 | Canada | A1 | |
| WO0101219A2 | World Intellectual Property Organization (WIPO) | A2 | |
| AU5891500A | Australia | A | |
| WO0101219A3 | World Intellectual Property Organization (WIPO) | A3 | |
| WO0101219A8 | World Intellectual Property Organization (WIPO) | A8 | |
| EP1203275A2 | European Patent Office (EPO) | A2 | |
| JP2003503787A | Japan | A | |
| CA2472442A1 | Canada | A1 | |
| WO03060748A2 | World Intellectual Property Organization (WIPO) | A2 | |
| AU2003217190A1 | Australia | A1 | |
| AU2003217190A8 | Australia | A8 | |
| US2003195938A1 | United States of America | A1 | |
| WO03060748A3 | World Intellectual Property Organization (WIPO) | A3 | |
| EP1203275A4 | European Patent Office (EPO) | A4 | |
| EP1502203A2 | European Patent Office (EPO) | A2 | |
| US6857004B1 | United States of America | B1 | |
| US2005038852A1 | United States of America | A1 | |
| JP2005515551A | Japan | A | |
| US7418470B2 | United States of America | B2 | |
| JP2008243216A | Japan | A | |
| US2009055625A1 | United States of America | A1 | |
| US2010049941A1 | United States of America | A1 | |
| US2010094924A1 | United States of America | A1 | |
| US7730121B2 | United States of America | B2 | |
| US2010183028A1 | United States of America | A1 | |
| US2010185719A1 | United States of America | A1 | |
| US2010251259A1 | United States of America | A1 | |
| JP2010277604A | Japan | A | |
| JP4596781B2 | Japan | B2 | |
| US7941479B2 | United States of America | B2 | |
| JP2011100487A | Japan | A | |
| US7958194B2 | United States of America | B2 | |
| JP4698700B2 | Japan | B2 | |
| US8325761B2 | United States of America | B2 | |
| JP2013061958A | Japan | A | |
| US8499025B2 | United States of America | B2 | |
| US2013311543A1 | United States of America | A1 | |
| JP5487128B2 | Japan | B2 | |
| US9626329B2 | United States of America | B2 |
8 legal events, as the office reported them to INPADOC
Over the term
Point at a mark for the eventEvents
| Event | Code | |
|---|---|---|
| Maintenance fee paymentMAFP | MAFP | |
| Fee paymentFPAY | FPAY | |
| RefundREFUND - PAYMENT OF MAINTENANCE FEE, 8TH YEAR, LARGE ENTITY (ORIGINAL EVENT CODE: R1552); ENTITY STATUS OF PATENT OWNER: SMALL ENTITYREFU | REFU | |
| Fee payment procedurePAT HOLDER CLAIMS SMALL ENTITY STATUS, ENTITY STATUS SET TO SMALL (ORIGINAL EVENT CODE: LTOS); ENTITY STATUS OF PATENT OWNER: SMALL ENTITYFEPP | FEPP | |
| Fee paymentFPAY | FPAY | |
| Certificate of correctionCC | CC | |
| Information on status: patent grantGrantedPATENTED CASESTCF | STCF | |
| AssignmentAS | AS |
Numbers
- Publication, DOCDB
- 2003195938
- Publication, EPODOC
- US2003195938
- Application
- 10340524
- Application, DOCDB
- 34052403
- Application, EPODOC
- US20030340524
Titles
- English
- Parallel processing systems and method
Classification
- CPC, 5
- G06F15/163
- G06F8/45
- G06F9/5044
- G06F9/5066
- G06F2209/509
- IPC, 2
- G06F9 44
- G06F9 50
- USPC, 1
- 709208000