Showing posts with label parallel programming. Show all posts
Showing posts with label parallel programming. Show all posts

Friday, November 23, 2007

matrix multiplication

0 comments;Click here for request info on this topic
Objective type questions

Code for matrix multiplication(question 1 to 6 based on this )

for (i=0 ; i<1 ; i++)
{
for (j=0 ; j
{
c[i][j] = 0.0;
for (k=0 ; k
c[i][j] + = a[i][k] * b[k][j];
}
}



Which loop does we parallize in matrix multiplication?
a) Outermost
b) Innermost
c) Both
d) All of above

Ans) a,
because in inner most loop data dependencies arises i.e. Two segment access/update a piece of common data, so outermost loop.
How the data dependency between line 3 be removed?
e) Cannot removed
f) By paralleling some other loop
g) By initializing the matrix C outside the loop
h) None of these
Ans) c,



If we parallelize j loop , the parallel algorithm executes n synchronization(per iteration of i )and grain size of parallel code is :
i) Θ(m n l / p)
j) Θ(l / p)
k) О(n l / p )
l) None of these.
Ans), c

If we parallelize i loop , the parallel algorithm executes only one synchronization, grain size of parallel code is :
m) О( m n l / p )
n) Θ( m n l / p )
o) О(n l / p )
p) None of these.
Ans), b

How many rows of each resultant matrix are calculated by each process in tightly coupled multiprocessor?
q) m n l / p
r) l / p
s) n l / p
t) none of these.
Ans), b

Time needed for computing a single row of matrix multiplication.
u) О m n
v) Θ m/n
w) Θ m n
x) None of these
Ans), c

All process needed to be synchronized once so synchronization over head is :
y) Θ l / p
z) Θ m / n
aa) Θ p
bb) None of these
Ans), c

What is the complexity of matrix multiplication parallel algorithm?
cc) Θ(n3 / p + p)
dd) Θ(n3 / p * p)
ee) Θ(n3 / p / p)
ff) None of these
Ans), a

As for each element of C , we must compute


so ,matrix-matrix multiplication involves, how many operations.


(N3)
N2
None of these

Ans), a

9. Assuming that = (n/k)2 processes are present, matrix multiplication is performed by dividing A and B into p blocks of size k*k. Each block multiplication requires, k3 additions, and k3 multiplications.
3k memory fetches
2k2 memory fetches
3k2 memory fetches
none of these

10.




Subjective type questions

Q 1. which loop we parallelize in matrix multiplications?
Ans ) Sequential code of matrix multilication
for (i=0 ; i<1 ; i++)
{
for (j=0 ; j
{
c[i][j] = 0.0;
for (k=0 ; k
c[i][j] + = a[i][k] * b[k][j];
}
}
In the above code there are three loop which can be parallelized but which one should be made parallelized.
Here, we can parallelize any loop either j or i without causing any problem because data dependence arises only in the innermost for loop that is in 3 lines over here.
(Data dependency is that when two segment access/update the common data. For example consider a= a+1and b=a+1 assume a=0,b=0 in order execution will give a=1,b=2,if opposite order result will be a=1,b=1.when executed parallelly we don’t expect any order and when the program is run again and again different output will be generated).
If we parallelize j loop, the parallel algorithm executes n synchronization (per iteration of i) and grain size of parallel code is О (n l / p).
If we parallelize i loop, the parallel algorithm executes only one synchronization, grain size of parallel code is Θ (m n l / p).
The rule should be always try to parallelize the outermost loop to maximize the grain size



Q 2. How the parallel multiplication is achieved and discusses the one
dimensional and two dimensional decompositions?

Ans. One dimensional decomposition: In particular, we consider the problem of developing a library to compute C = A.B , where A , B , and C are dense matrices of size N N . (A dense matrix is a matrix in which most of the entries are nonzero.) This matrix-matrix multiplication involves operations, since for each element of C , we must compute


Figure 4.10:
Matrix-matrix multiplication A.B=C with matrices A , B , and C decomposed in one dimension. The components of A , B , and C allocated to a single task are shaded black. During execution, this task requires all of matrix A (shown stippled).

We start by examining algorithms for various distributions of A , B , and C. We first consider a one-dimensional, columnwise decomposition in which each task encapsulates corresponding columns from A , B , and C . One parallel algorithm makes each task responsible for all computation associated with its . As shown in Figure 4.10, each task requires all of matrix A in order to compute its . data are required from each of P-1 other tasks, giving the following per-processor communication cost:

Note that as each task performs computation, if N P , then the algorithm will have to transfer roughly one word of data for each multiplication and addition performed. Hence, the algorithm can be expected to be efficient only when N is much larger than P or the cost of computation is much larger than .
Two dimensional decomposition

Figure 4.11
Matrix-matrix multiplication A.B=C with matrices A , B , and C decomposed in two dimensions. The components of A , B , and C allocated to a single task are shaded black. During execution, this task requires corresponding rows and columns of matrix A and B , respectively (shown stippled).
Less cost in two dimensional over one dimensional
We consider a two-dimensional decomposition of A , B , and C . As in the one-dimensional algorithm, we assume that a task encapsulates corresponding elements of A , B , and C and that each task is responsible for all computation associated with its . The computation of a single element requires an entire row and column of A and B , respectively. Hence, as shown in Figure 4.11, the computation performed within a single task requires the A and B submatrices allocated to tasks in the same row and column, respectively. This is a total of data, considerably less than in the one-dimensional algorithm.
Q3. A lgorithm of parallel matrix multiplication?
We represent the matrix multiplication algorithm for tightly coupled multiprocessor .
We know that it is not very expensive to share and exchange data across processes in tightly coupled machines.
The concern here is to achieve maximum speed up .
We choose to parallelize the outer most loop in this algorithm. We have Process working in parallel computing l/p rows of resultant matrix.
Each process in the program works to compute every p-th row (i.e., i, i+1, i+2p and soon) of the matrix.

Begin
/* we consider processor indexed from 1 to p addressed as p(1) top(p)*/
for all p(r),where 1 <= r <= p do begin
/* i,j,k,t are local to the p(r) */
for i= r to 1 step p do begin
for j= 1 to n do begin
t=0;
for k= 1 to m do begin
t + = a[I,k] * b[k,j];
end;
c[i,j] = t;
end;
end;
end (for all)
end.

We represent the matrix multiplication algorithm for loosely coupled multiprocessor
Where some matrix element may be much easier to access than others
Some reason for the algorithm needs to be redesigned.
A process must access l/ rows of matrix A and all elements of B(l/) times only a single addition and multiplication takes lace for every element in B
Ignoring the memory access time can be safe on tightly coupled multiprocessors where global memory is equally accessible to each processors this is not so with loosely coupled multi computer
On loosely coupled multiprocessors it is best to keep most memory references local as far as possible.
Block matrix multiplication
As over here Aand B both are n*n marix, where n=2k.
Then A and B can be thought of as conglomerates of four smaller matrices, each of size k*k.


A = A11 A12
A21 A22



B = B11 B12
B21 B22
C = C11 C 12
C 21 C 22

Then, C is written as:

A11 B 11 + A12B21 A11 B 12 + A12B22
A21 B 11 + A22B21 A21 B 12 + A22B22
C =

As seen from the above equations, the scheme partitions the problem into disjoint subtasks, which can be solved independently.
The multiplications of the smaller sub matrices can be done in parallel. These local multiplications can be further subdivided using the same scheme.


ANALYSIS:
We can assign processes to do the multiplication task local to the respective blocks.
Assuming that = (n/k)2 processes are present, matrix multiplication is performed by dividing A and B into p blocks of size k*k. Each block multiplication requires 2k2 memory fetches, k3 additions, and k3 multiplications.
The number of arithmetic operations per memory fetch has increased from 2 to k = n/√p, a significant improvement.
The block matrix multiplication algorithm performs better on the NUMA multiprocessors because it increases the number of computations performed per non local memory fetch.
For similar reasons a careful choice of block sizes to maximize the cache hit-rate can yield a sequential block-oriented matrix multiplication algorithm that executes faster than traditional algorithm illustrated earlier.
Matrix-matrix multiplication algorithm based on two-dimensional decompositions. Each step involves three stages:
(a) an A submatrix is broadcast to other tasks in the same row;
(b) local computation is performed; and
(c) the B submatrix is rotated upwards within each column.
To complete the second parallel algorithm, we need to design a strategy for communicating the submatrices between tasks. One approach is for each task to execute the following logic (Figure 4.12):

set
for j =0 to in each row i
, the th task broadcasts
to the other tasks in the row accumulate .
send to upward neighbor
endfor

Each of the steps in this algorithm involves a broadcast to tasks (for A' ) and a nearest-neighbor communication (for B' ). Both communications involve data. Because the broadcast can be accomplished in steps using a tree structure, the per-processor communication cost is
Read full story

Multi-processor systems OS

0 comments;Click here for request info on this topic
OBJECTIVE-TYPE QUESTIONS
Q1. _____ is a model of one address computer.
P-RAM
RAM
ROM
P-ROM

Ans: B

Q2. All algorithms for ___________________ can be expressed in terms of RAM model and its instruction set.

Parallel Machines
SIMD/ Vector Processing Systems
Sequential Machines
Every type of machines

Ans: C

Q3. Default behavior of P-RAM is running multiple instruction Streams on different processors. However, it possible to constraint all the processors to fetch and execute the same set of instructions? (True / False)

Ans: True.

Q4. A multi-computer system is a _____________ Whereas Multi-Processor System is a _____________

Collection of computer systems that include multiple processors, single autonomous computer that is interconnected with other systems for achieving parallelism.
Single autonomous computer that is interconnected with other systems for achieving parallelism, Collection of computer systems that include multiple processors.
Single Computer that include multiple processors, Collection of multiple interconnected autonomous computer systems.
Collection of multiple interconnected autonomous computer systems, Single Computer that include multiple processors.
Ans: D

Q5. The most suitable interconnection network for small multiprocessor systems (i.e. having upto 10 processing units) is:

Common Bus
Cross Bar Switch
Hierarchical Organization of Processors
None Of the Above
Ans: A

Q6. By organizing processors and memory modules in hierarchically, one can obtain _________ architecture.

UMA
NUMA
Symmetric
Asymmetric
Ans: B

Q7. _______________ , which is used for keeping latest updated value in cache for single processor systems, is neither necessary nor sufficient for achieving cache coherence in multi-processor systems.

Static Coherence Checking Algorithms
Dynamic Coherence Checking Algorithms
Write through fashion
Static & Dynamic Coherence Checking Algorithm and Write through mechanism
Ans: C

Q8. List possible organizations in design of operating systems for multi-processor systems:

Ans: 1. Master Slave Configuration
2. Separate Supervisor Configuration
3. Floating Supervisor Configuration

Q9. In which Message passing Model, the sender is also blocked till receiver receives and acknowledges the message?

Synchronous
Asynchronous
Both Synchronous and Asynchronous
In None of Message Passing Model sender get blocked

Ans: A

Q10. Does Message Passing Model require mechanism for mutual exclusion to access shared data? (Yes/NO)

Ans. Does not require (NO)

Q11. Multithreading is implementation of which parallel programming model?

Message Passing Model
Functional and Logical Models
Data Parallel Model
Shared Memory Model
Ans: D

Q12. Give an example (programming language) of Data parallel programming model for SIMD processors.

Ans: Fortran 90



SUBJECTIVE-TYPE QUESTIONS

Q1. Depending Upon the Sharing of resources Multiple Processor Systems can be classified into which categories? Explain each briefly. What are desirable properties of Multi-Processor Systems?
Ans: Depending Upon the Sharing of resources such as Memory, Clock, Computer bus and the way of Communication, Multi Processor Systems can be classified in following two categories:
Tightly Coupled Systems
Loosely Coupled Systems

· Tightly Coupled Systems:
The Tightly Coupled Multi Processor Systems are those in which the Processors are kept in close communication to share the processing tasks. These systems can further be divided in to two sub-categories:
1. Shared Memory Multi-Processor systems
2. Message Passing Multi-Computer Systems

Shared Memory Multi-Processor systems contain multiple processors that are connected at the system bus level. These systems allow 1000 CPU’s to communicate via a Shared Memory. Every CPU has equal access to the entire physical memory. These systems may also participate in a memory hierarchy with both local and shared memory.

Depending upon the way the Operating System used in Multiple Processors, two Organizations of Multi-Processor system are possible. These are:
Asymmetric Multiprocessor
Symmetric Multiprocessor

Asymmetric Multiprocessor: In an asymmetrical multiprocessor system, one processor, called Master, is dedicated to execute the Operating System. The remaining processors are usually identical and form a pool of computational processors. All the processors in this pool are called Slaves. The Master processor schedules the work and control the activities of the slave processors depending upon the resources available with them. Due to this arrangement of multi processors, it is also known as Master-Slave Multi Processing. In some asymmetrical multiprocessor system, the processors are assigned different roles such that they do not depend upon the master for their processing. For Example: One processor may handle all trivial requests and jobs like editing, calculating etc.


Symmetric Multi Processor: In a symmetrical multiprocessor system, all of the processors are essentially identical and perform identical functions. There is one copy of Operating System in Memory. The CPU, on which the system call was made, traps to the Kernel and processes the system calls. In such an arrangement there is said to be a floating Master as different processors execute operating system at different times.


Message Passing Multi-Computer Systems allow a number of CPU-Memory pairs (called node) to connect by some kind of high-speed interconnection. The Basic node of Multi Computer system consists of CPU, memory, a network interface and some times a hard disk but the Graphic adapter, monitor, keyboard and mouse are nearly always absent. Each memory is local to a single CPU and can only be directly accessed by that CPU. There is no shared memory in this design. So without a shared address space processors communicate by sending multiword messages over the interconnection (Message Passing) in order to share computational tasks. These systems are also known by a variety of names including Cluster Computers and COWS (Cluster Of Workstations).


These Multi-Computer systems can be either
Asymmetric Multi-Computers or
Symmetric Multi-Computers

1. In Asymmetric Multi-Computers, a front end computer interact with user as users may log into the front end computer which executes a full, multi-programmed operating System and provide all functions for program development. Other than this Front End Computer, other nodes of the Multi-Computer systems are Back End Computers, which are reserved for executing the parallel programs.
2. In Symmetric Multi-Computer operating system execute the same OS (Multi Programmed OS) and has identical functionality. User may log onto any computer to edit and compile their programs and any other computer may be called upon to execute a parallel program. This configuration of Multi Computer system solves the problem of single point of failure in Asymmetric Multi-Computer systems.

· Loosely Coupled Systems:
These types of Multi processor systems are also referred to as Distributed Systems. These are based on multiple standalone single or dual Processor Computers interconnected via a high-speed communication system (gigabit Ethernet is common). The intent of the distributed system is to turn a loosly connected bunch of machines into a cohrent system based upon one concept. This is done by having another layer of software on the top of the Operating System. This Layer is called Middleware.


Desireable Properties of MultiProcessor Systems:

Process Recoverability: - If aprocessor fails the process running on it should be recoverable. Another processor must be assigned to it. This can be achieved by maintaining a shared register file that keep a record of the state for each active process in the system.

Efficient Context Switching: - While executing parallel programs, it is needed to swap the processes in and out of the memory for efficient utilization of the resources. The processors must have efficient mechanisms to support Operating System in switching the context.

Large Virtual and Physical address space: - As the size of the problem being solved on the parallel machines, increases there is need for large memory space and addressing capacities.

Effective Synchronization Primitives: - Multiple processes that constitute a parallel program need to cooperate to compute results. This is possible by sharing data and providing access to shared data. To maintain the integrity of the data, access needs to be allowed only in exclusive mode.

Interprocess Communication Mechanism: - Multiprocessor systems must provide communication between cooperating processes in the form of signals, interrupts and messages.


Q2. What do you mean by programming model? Which are various Parallel Programming Models available?
Ans:
A programming model is a collection of program abstractions providing a programmer a simplified and transparent view of computer Hardware/Software system. Parallel programming models are specifically designed for multiprocessors, multi computer, SIMD or vector computers. There are five such models described bellow that differ in the way these processes share data, achieve synchronization and communication:

1. Shared Memory Programming: -
Multiprocessor programming is based on the use of shared variables in commonly accessible memory for communication and sharing data. Besides sharing variables in a common address space communication also takes place through software signals and interrupts. As multiple processors may attempt to access the shared data processes must ensure exclusive access to the critical sections (Part of the program accessing the shared resources or variables).
The implementation of the model is supported in various ways on different platforms.
· On UNIX platform it is available in the form of linkable libraries for creation of processes, shared memory blocks, semaphores and IPC mechanisms.
· This model is also available in the form of multithreading libraries.

2. Message Passing: -
Multi computer employ message passing as the mechanism for inter process communication. In this model, two processes residing on two different nodes communicate with each other by passing messages over communication channel. The message may be instruction, data, and synchronisation or interrupt signal.
Message passing model does not require mechanism for mutual exclusion for access to shared data as there is no way for processes to share each others address space.
The implementation of this model can be done by two ways as explained bellow:

a. Synchronous Message Passing: - It requires both the sender and receiver to be synchronized just like in telephone call communication. No buffering of message is done by communication channel.
The receiver is always blocked waiting for the message to arrive.
The sender is also blocked till receiver receives and acknowledges the message.
This mode of communication is suitable for tightly coupled message passing multi computer systems where communication delay overhead is sufficiently small.

b. Asynchronous Message Passing: - This mode of the model does not impose blocking on the sender. The outgoing messages get buffered in the communication channel and are delivered to target when target process chooses to look for it.
This mode of communication is suitable for loosely couples systems that are made up of networked autonomous machines.

3. Data-parallel Model: -
This model for SIMD processors is an extension of the sequential programming
The language compiler for a specific machine must translate the source code to executable object code, which will exploit the parallelism in the hardware. For Example: - Fortran is tailored for data parallelism.
The compilers must be aware of underlying interconnection topology to generate optimal code for array processor.
Synchronization of data parallel programs is done at compile time rather than at run time.

4. Object Oriented Model: -
In this model mapping of execution units to objects is achieved. These objects are dynamically created and manipulated. Sending and receiving messages among the objects achieve communication.
Concurrent programs are built from low level objects such as processes, threads, and semaphores into high level objects like monitors and program modules.



5. Functional and Logical Models: -
A functional programming language emphasizes the functionality of a program. There is no concept of storage, assignment and branching in functional programming. The evaluation of the function produces the same result regardless the order in which its arguments are evaluated.
Logical Programming is based upon predicate logic. This model is suitable for knowledge processing dealing with large database. This model adopts an implicit search strategy and support parallelism.
Both of these models are used in Artificial Intelligence Applications.



Q3. What are Abstract Models available for Sequential and Parallel Machines/Computers?
Ans:
The purposes of using such models are:
To get feel of the capacities of the machines without bothering about the specific constraints of the real life machines.
It is convenient to develop an algorithm for a general model and then to map it to the actual machine.

The abstract models available for Sequential as well as for parallel machines are:
RAM (Random Access Machine) à Abstract model of sequential computer.




P-RAM (Parallel RAM) à Abstract model of Parallel Computer.
Read full story

Assignment Issues in parallel programming

0 comments;Click here for request info on this topic
Q1. How parallelism can be introduced in sequential machines? Describe any two techniques briefly.

Answer: The sequential machines have been made faster by incorporating various schemes. The main idea in all these schemes is to match the speed of various components so as to utilize the resources to their peak performances.

There are various parallelism techniques that we can use in uniprocessor or sequential machines.
Some of them are:
1. Multiplicity of Functional Units.
2. Overlapped CPU and I/O operations.
3. Pipelining within the CPU.
4. Hierarchical Memory Systems.
5. Multiprogramming and Time Sharing.

Multiplicity of functional units:

Mostly computers have only one Arithmetic & Logic Unit (ALU). ALU could perform one function at a time.
The practical machines used today have multiple and specialized units that can operate in parallel.
In sequential machines now it is possible to have multiple functional units for addition, multiplication, division, increment, decrement, Boolean operations and shift instructions.



Example:
CDC-6600(designed in 1964) had 10 functional units.
A scoreboard was used to keep track of availability of the functional units and the registers being demanded.

With 10 functional units and 24 registers the machine could achieve a significant increase rate of instruction execution.

2. Overlapped CPU and I/O operations

As we know each instruction contains Computational phase as well as I/O phase. I/O operations can be performed simultaneously with the computational tasks by using separate I/O controllers, I/O processors.
Computers requiring multiple I/O devices employ multiple I/O subsystems which can operate concurrently. This type of multiprocessing speeds up the data transfer between the external devices and memory.

BUS Contention
The design of such systems depends upon Bus Contention.
Because if device interface and computational task require same bus to transfer data then there will be no speed up in data transfer rate.
To minimize bus contention, some systems employ redundant bus architecture.
Tradeoff is there,
Consequent complications in the design but at the same time the utilization of CPU and other resources maximize.

Processor Intervention
Data is transferred between memory and I/O devices or vice versa.
Processor interferes in this data transfer.
This interference decreases data transfer rate.
To avoid this we use Direct Memory Access (DMA) technique.
DMA is used to provide direct information transfer between the external devices and the primary memory




Q2. How pipelining within the CPU introduces parallelism in sequential machines?

Answer: Today every sequential processor manufactured is taking advantage of parallelism in the form of pipelining instructions.

The purpose of pipeline parallelism is to increase the speed of program and to decrease I/O operations.

Pipeline parallelism is when multiple steps depend on each other but the extension can overlap and the output of one step is streamed as input to the next step.

The different stages of pipelining are

Fetch Stage
Decode Stage
Issue stage
Execute stage
Write back Stage

The fetch stage (F) fetches instructions from cache memory, one per cycle.

The decode stage (D) reveals the instruction function to be performed and identifies the resources needed. Resources include general purpose registers, buses and functional units.

The issue stage (I) reserves resources. Pipeline interlock controls are maintained at this stage the operands are also read from registers during the issue stage.

The instructions are executed in one or several execute stages (E).

The last write back stage (W) is used to write results into the registers. Memory load and store operations are treated as part of execution.

To facilitate instruction execution through the pipe instruction prefetching is used.

The execution of multiple instructions is overlapped in time – even before an instruction gets completely executed, another instruction may be in the process of being decoded, yet another instruction may be getting fetched, and so on.
Pipelining of tasks gives us temporal parallelism.

Figure shows the flow of machine instructions through a typical sequential machine. These eight instructions are for execution of the statements X=Y+Z and A=B*C.






The shaded boxes correspond to idle cycles when instruction issues are blocked due to resources latency or conflicts or due to data dependences.

The first two load instructions issue on consecutive cycles. The add is dependent on both loads and must wait three cycles before the data (Y& Z) are loaded in.

Similarly, the store of the sum to memory location X must wait three cycles for add to finish due to a flow dependence. There are similar blockages during the calculation of A.

The total time required is 17 clock cycles. This time is measured beginning at cycle four when the first instruction starts execution until cycle 20 the last instruction starts execution.



Figure shows an improved time after the instruction issuing order is change to eliminate unnecessary delays due to dependence. The idea is to issue all four load operations in the beginning. Both add and multiply instructions are blocked fewer due to fewer cycles due to this data prefetching. The reordering should not change the end results. The time required is being reduced to 11 cycles, measured from 4 to 14.

Q3. What are the various issues in parallel programming?

Answer: The various issues in parallel programming are:

1. Load balancing:

The operating system must utilize the resources efficiently. This is commonly expressed in terms of achieving a uniform balance of loads across the processors. The operating system should schedule the subtasks such that there are no idle resources including the processors.

2. Scheduling cooperating processes:

Parallel programs which consist of concurrently executing tasks must be scheduled such that collectively they are able to use the resources in the machine required to solve the given problem. Scheduling half the subtask of one parallel program and half the subtasks of another parallel program which can not cooperate can be wasteful.
If a parallel program needs all its subtasks to be running and cooperating, then scheduling half of the tasks can lead to starvation. We can achieve this by maintaining a hierarchy of processes in the sense of parent-child relationship within the operating system’s data structure. Scheduling of processes then can be done such that the processes belonging to a single parallel program are scheduled for execution together.

3. Graceful degradation in case of failure of one of the resources:

Given that a parallel machine has multiple resources of the same kind, one must expect a higher degree of fault tolerance from it. Failure of one of its resources should not result in a catastrophic system crash. The operating system should be able to reschedule the task that had been running on the failed resource and continue the parallel program. Ideally, there should be only a fractional degradation in the performance of the parallel machine.

4. Communication schemes:

Most parallel programs need to share data and intermediate result across subtasks during the processing towards the solution of the problem. To achieve effective cooperation among the subtasks of a parallel program, the operating system must provide adequate facility for communication between tasks. These facilities vary depending on whether the machine is of a shared memory type or of a distributed memory type.

5. Synchronization mechanisms:

To ensure integrity of the shared data across the subtask of a parallel program, synchronization between tasks is required. The shared data must be accessed under mutual exclusion (implemented using a mechanism like semaphores). The task may need to wait till some state is reached across all the tasks of the parallel programs. The operating systems need to provide signaling mechanisms for such synchronization requirements.
Read full story

Assignment Parallel Programming

0 comments;Click here for request info on this topic
Long questions

Q1. Explain parallel reduction with an example.

Given a set of n values a0, a1,…….an-1 and an associative operator ○, reduction is the process of computing a0 ○ a1 ○ an-1. Addition, multiplication, and finding the maximum/minimum of a set are examples of associative operators.
The process of reduction is viewed as a tree-structured operation. The elements of input to the reduction problem are placed at the leaf nodes of a binary tree. The reduction operation is applied to children of each parent node, and the result is propagated towards the root of the tree.
Parallel summation is an example of a reduction operation.







PRAM algorithm to sum n elements using ën/2û processors

SUM :
/*Initial condition: List of n ³ 1 elements stored in A[0…(n-1)]
Final condition: Sum of elements stored in A[0] */
Global variables: n,A[0…(n-1)],j
Begin
spawn(P0, P1, P2,…,P└n/2┘- 1)
for all Pi ,where 0 £ i £ ën/2û -1 do
for j ¬ 0 to log én –1ù do
if i modulo 2j and 2i + 2j < n then
A[2i] ¬A[2i] + A[2i + 2j]
endif
endfor
endfor
end

Analysis of Parallel Reduction

•There is a dependency across levels of reduction. Unless reduction is ready at a lower level, a higher level cannot proceed with the operation.
• Overall time complexity of the algorithm is Θ (log n), given ën/2û processes
• The spawn routine requires élog ën/2û ù process doubling steps.
• The sequential for loop executes élog nù times.

O2. Explain Odd-Even Transposition Sort

It is designed for the processor array model in which the processing elements are organized into one-dimensional mesh.
Assume that A=(a0, a1,…,an-1) is the set of n elements to be sorted. Each of the n processing elements contains two local variables: a, unique element of array A, and t, a variable containing a value retrieved from a neighboring processing element. The algorithm performs n/2 iterations, and each iteration has two phase. In the first phase, called odd-even exchange, the value of a in every odd-numbered processor(except processor n-1) is compared with the value of a stored in the successor processor. The values are exchanged, if necessary, so that the lowered-numbered processor contain the smaller value. In the second phase, called even-odd exchange, the value of a in every even numbered processor is compared with the value of a in the successor processor. The values are exchanged, if necessary, so that lower numbered processor contain the smaller value. After n/2 iterations the value must be sorted.
Odd-Even Transposition Sort Algorithm for the one-dimensional mesh processor array model

Parameter n
Global i {Element to be sorted}
Local t {Element taken from adjacent processor}
Begin
for i ¬1 to n/2 do
for all pj, where 0 £ j £ n-1 do

if j < n-1 and odd(j) then
{Odd-even exchange}
t Ü successor(a) {Get value from successor}
successor(a)Ümax(a,t) {Give away larger value}
a¬ min(a,t) {Keep smaller value}
endif

if even(j) then
{Even-odd exchange}
t Ü successor(a) {Get value from successor}
successor(a) Ü max(a,t) {Give away larger value}
a ¬ min(a,t) {Keep smaller value}
endif
endfor
endfor
end

Example: Odd-Even Transposition Sort for eight values

Indices: 0 1 2 3 4 5 6 7
Initial Values: G H F D E C B A
After odd-even exchange: G F < H D < E B < C A
After even-odd exchange: F < G D < H B < E A < C
After odd-even exchange: F D < G B < H A < E C
After even-odd exchange: D < F B < G A < H C < E
After odd-even exchange: D B < F A < G C < H E
After even-odd exchange: B < D A < F C < G E < H
After odd-even exchange: B A < D C < F E < G H
After even-odd exchange: A < B C < D E < F G < H


Analysis of Odd-Even Transposition Sort

• The complexity of sorting n elements on a one- dimensional mesh processor array with n processors using odd-even transposition sort is Θ(n).

Q3. Explain Enumeration Sort
Assume that we are given table of n elements, denoted a0, a1,…., an-1, on which a linear order has been defined. Thus for any two elements ai and aj, exactly one of the following cases must be true: ai < ai =" aj,"> aj. The goal of sorting is to find a permutation (∏0, ∏1,…. ∏n-1) such that a∏0 £ a∏1 £ ……a ∏n-1.
An enumeration sort computes the final position of each element in the sorted list by comparing it with the other elements and counting the number of elements having smaller value. If j elements have smaller value than ai, then ∏j = i ; i.e., element ai
is the (j +1) element on the sorted list following a∏0,……, a∏j-1.

Algorithm for Enumeration Sort

ENUMERATION SORT (CRCW PRAM)
Parameter n {Number of elements}
Global a[0…(n-1)] {Elements to be sorted}
position[0…(n-1)] {Sorted positions}
sorted[0…(n-1)] {Contains sorted elements}
Begin
spawn(Pi,j, for all 0 £ i, j < n)
for all Pi,j, where 0 £ i, j< n do
position[i] ¬0
if a[i] < a[j] or (a[i] = a[j] and i < j) then
position[i] ¬ 1
endif
endfor
for all Pi,0, where 0 £ i < n do
sorted[position[i]] ¬ a[i]
endfor
end

Analysis of Enumeration Sort
•A set of n elements can be sorted in Q(log n) time with n2 processors, given a CRCW PRAM model in where simultaneous writes to the same memory location cuse the sum of the values to be assigned.
• If the time needed to spawn the processors is not counted, the algorithm executes in constant time.

Q4. Explain Parallel Quick Sort.
Quick sort is a divide and conquer algorithm that easily yields itself to parallelisation.
The array is partitioned into two parts using a pivot element, such that all elements in the left partition are smaller than the pivot element and the elements to the right are larger. The pivot element is inserted in between the two partitions and hence in the sorted position after partitioning. Apply the same algorithm recursively to each left and right parts of the array, if there is just one element in the array to stop.
To parallelise this algorithm, a set of processes(pool of workers) is created. The extents of the array to be sorted are placed on a stack to indicate the presence of work to be done. The first process to fetch the work partitions the array and places the extents of the resulting two partitions on the stack. The workers fetch the extents of a partition to be sorted from the stack and create new partitions. These are added to the stack.

Algorithm for Parallel Quick Sort

Global n {Size of array of unsorted elements}
a[0…(n-1)] {array of elements to be sorted}
sorted {Number of elements in sorted position}
min.partition {Smallest subarray that is partitioned rather than sorted directly}
Local bounds {Indices of unsorted subarray}
median {Final position in subarray of partitioning key}
Begin
sorted ¬ 0
INITIALIZE.STACK()

for all Pi, where 0 £ i < p do
while( sorted < n) do
bounds ¬ STACK.DELETE()
while(bounds.low < bounds.high) do
if(bounds.high - bounds.low < min.partition) then
INSERTION.SORT(a, bounds.low, bounds.high)
ADD.TO.SORTED(bounds.high-bounds.low +1 )
exit while
else
median¬ PARTITION(bounds.low, bounds.high)
STACK.INSERT(median + 1, bounds.high)
bounds.high ¬ median – 1

if bounds.low = bounds.high then
ADD.TO.SORTED(2)
else
ADD.TO.SORTED(1)
endif
endif
endwhile
endwhile
endfor
end

In the above algorithm, function INITIALIZE.STACK initializes the shared stack containing the indices of unsorted subarrays. When a process calls function STACK.DELETE, it receives the indices of an unsorted sub array if the stack contains indices; otherwise, there is no useful work to do at this point. Function STACK.INSERT adds the indices of an unsorted subarray to the stack. Since all these functions access the same shared data structure, their execution must be mutually exclusive. Function ADD.TO.SORTED increases the count of elements that are in their correct positions and execution of this function, too, must be mutually exclusive.
Read full story

Monday, September 17, 2007

Parallel Programming: MPI

0 comments;Click here for request info on this topic
Parallelization - Parallel Knock

We describe a simple example of message passing in a parallel code. For those you learning about parallel code, we provide the following discussion of elemental message-passing as an introduction to how communication can be performed in parallel codes.
Knock
This example, knock.c, passes a simple message from the even processors to the odd ones. Then, the odd processors send a reply in a message back to the even ones. The "conversation" is displayed on all processors, and the code ends.

This code's main routine performs the appropriate message-passing calls. In a more complex parallel code, these can be used to pass information processors to coordinate or complete useful work. Here, we focus on the basic technique to pass these messages. The C source is shown in Listing 1.
Listing 1 - knock.c, or see knock.f
#include
#include
int main(int argc, char *argv[])
{
/*
this program demonstrates a very simple MPI program illustrating
communication between two nodes. Node 0 sends a message consisting
of 3 integers to node 1. Upon receipt, node 1 replies with another
message consisting of 3 integers.
*/
/* get definition of MPI constants */
#include "mpi.h"
/* define variables
ierror = error indicator
nproc = number of processors participating
idproc = each processor's id number (0 <= idproc < nproc)
len = length of message (in words)
tag = message tag, used to distinguish between messages */
int ierror, nproc, idproc, len=3, tag=1;
/* status = status array returned by MPI_Recv */
MPI_Status status;
/* sendmsg = message being sent */
/* recvmsg = message being received */
/* replymsg = reply message being sent */
int recvmsg[3];
int sendmsg[3] = {1802399587,1798073198,1868786465};
int replymsg[3] = {2003332903,1931506792,1701995839};

printf("Knock program initializing...\n");

/* initialize the MPI execution environment */
ierror = MPI_Init(&argc, &argv);
/* stop if MPI could not be initialized */
if (ierror)
exit(1);
/* determine nproc, number of participating processors */
ierror = MPI_Comm_size(MPI_COMM_WORLD,&nproc);
/* determine idproc, the processor's id */
ierror = MPI_Comm_rank(MPI_COMM_WORLD,&idproc);
/* use only even number of processors */
nproc = 2*(nproc/2);
if (idproc < nproc) {
/* even processor sends and prints message,
then receives and prints reply */
if (idproc%2==0) {
ierror = MPI_Send(&sendmsg,len,MPI_INT,idproc+1,tag,
MPI_COMM_WORLD);
printf("proc %d sent: %.12s\n",idproc,sendmsg);
ierror = MPI_Recv(&recvmsg,len,MPI_INT,idproc+1,tag+1,
MPI_COMM_WORLD,&status);
printf("proc %d received: %.12s\n",idproc,recvmsg);
}
/* odd processor receives and prints message,
then sends reply and prints it */
else {
ierror = MPI_Recv(&recvmsg,len,MPI_INT,idproc-1,tag,
MPI_COMM_WORLD,&status);
printf("proc %d received: %.12s\n",idproc,recvmsg);
ierror = MPI_Send(&replymsg,len,MPI_INT,idproc-1,tag+1,
MPI_COMM_WORLD);
printf("proc %d sent: %.12s\n",idproc,replymsg);
}
}
/* terminate MPI execution environment */
ierror = MPI_Finalize();
if (idproc < nproc) {
printf("hit carriage return to continue\n");
getchar();
}
return 0;
}
Discussing the parallel aspects of this code:

mpi.h - the header file for the MPI library, required to access information about the parallel system and perform communication
idproc, nproc - nproc describes how many processors are currently running this job and idproc identifies the designation, labeled using integers from 0 to nproc - 1, of this processor. This information is sufficient to identify exactly which part of the problem this instance of the executable should work on.
MPI_Init - performs the actual initialization of MPI, setting up the connections between processors for any subsequent message passing. It returns an error code; zero means no error.
MPI_COMM_WORLD - MPI defines communicator worlds or communicators that define a set of processors that can communicate with each other. At initialization, one communicator, MPI_COMM_WORLD, covers all the processors in the system. Other MPI calls can define arbitrary subsets of MPI_COMM_WORLD, making it possible to confine a code to a particular processor subset just by passing it the appropriate communicator. In simple cases such as this, using MPI_COMM_WORLD is sufficient.
MPI_Comm_size - accesses the processor count of the parallel system
MPI_Comm_rank - accesses the identification number of this particular processor
if (idproc%2==0) - Up until this line, the execution of main by each processor has been the same. However, idproc distinguishes each processor. We can take advantage of this difference by guiding the processor to a different execution as a function of idproc. This Boolean statement causes those processors whose idproc is even (0, 2, 4, ...) to execute the one part of the if-block, while the odd processors (1, 3, 5, ...) execute the second part of the block. At this point, no central authority is consulted for assignments; each processor recognizes what part of the problem it should perform using only its idproc.
MPI_Recv and MPI_Send - These elemental MPI calls perform communication between processors. In this case, the even processor, executing the first block of the if statement, begins by using MPI_Send to send the contents of the sendmsg array, whose size is given as 3 long integers, specified by len and MPI_INT, to the next processor, specified by idproc+1. tag merely identifies the message for debugging purposes, and must match the tag on the corresponding receive. MPI_Send returns when the MPI library no longer requires that block of memory, allowing your code to execute, but the message may or may not have arrived at its destination yet.
Meanwhile, the odd processor uses MPI_Recv to receive the data from the even processor at idproc-1. This processor specifies that the message be written to recvmsg. The call returns information about the message in the status variable. MPI_Recv returns when it is finished with the recvmsg array.
Note that the even and odd processors might not be in lock step with one another. Such tight control is not needed to complete this task. One processor could reach one MPI call sooner than the other, but the MPI library performs the necessary coordination so that the message is properly handled while allowing each processor to continue with work as soon as possible. The processors affect one another only when the communicate.
After printing the first message, each processor switches roles. Now the odd processors call MPI_Send to send back a reply in replymsg, and the even processors call MPI_Recv to receive the odd processors' reply in their recvmsg. For debugging, the tag used is incremented on all processors. Note that this is a distributed-memory model, so the instance of recvmsg on one processor resides in the memory on one machine while the other instance of recvmsg resides on a completely different piece of hardware. Here we are using the even processors' recvmsg for the first time while the odd processors' recvmsg holds data from the first message. Finally, each processor prints this second message.
MPI_Finalize - performs the properly cleanup and close the connections between codes running on other processors and release control.
When this code, with the initialization values in sendmsg and replymsg, is run in parallel on two processors, the code produces this output:
Processor 0 output
Processor 1 output
Knock program initializing...
proc 0 sent: knock,knock!
proc 0 received: who's there?
hit carriage return to continue
Knock program initializing...
proc 1 received: knock,knock!
proc 1 sent: who's there?
hit carriage return to continue
It performs the opening lines of the classic "knock, knock" joke. This example is an appropriate introduction because that joke requires the first message to be sent one way, then the next message the other way. After each message pass, all processors print the data they have showing how the conversation is unfolding.

Conclusion
The purpose of this discussion was to highlight the basic techniques needed to perform elemental message passing between processors. It also shows how one identification number, idproc, can be used to cause different behavior by distinguishing between processors. The overhead of initialization (MPI_Init et al) and shut down (MPI_Finalize) seems large compared to the actual message passing (MPI_Send and MPI_Recv) only because the nature of the message passing is so simple. In a real-world parallel code, the message-passing calls can be much more complex and interwoven with execution code. We chose a simple example so attention could be drawn to the communications calls. The reader is welcome to explore how to integrate these calls into their own code
Read full story