Parallel Sorting, Safety, and Deadlock Prevention
Sorting is a fundamental problem in computing. In a distributed-memory system, sorting takes on a distinct character because memory is physically partitioned among separate processes.
In this chapter, we develop a parallel odd-even transposition sort, examine the critical concept of MPI communication safety, and resolve deadlocks using MPI_Sendrecv.
3.16 Problem Definition: Distributed Memory Sorting
Section titled “3.16 Problem Definition: Distributed Memory Sorting”In a distributed-memory environment with processes and total keys, each process starts and finishes with a local block of keys (assuming evenly divides ).
The algorithm terminates when the following conditions are satisfied:
- Local Order: The keys residing on each process are sorted in non-decreasing order.
- Global Order: If , then every key assigned to process is less than or equal to every key assigned to process :
If we align the keys across processes in ascending rank order (Process 0 followed by Process 1, etc.), the overall collection is sorted globally.
3.17 Serial Sorting: From Bubble Sort to Odd-Even Transposition
Section titled “3.17 Serial Sorting: From Bubble Sort to Odd-Even Transposition”Why Bubble Sort Fails to Parallelize
Section titled “Why Bubble Sort Fails to Parallelize”Bubble sort iterates through a list, comparing adjacent pairs and swapping them if out of order:
/* Program 3.14: Serial Bubble Sort */void Bubble_sort(int a[], int n) { int list_length, i, temp; for (list_length = n; list_length >= 2; list_length--) { for (i = 0; i < list_length - 1; i++) { if (a[i] > a[i + 1]) { temp = a[i]; a[i] = a[i + 1]; a[i + 1] = temp; } } }}Bubble sort cannot be parallelized effectively because of strict sequential data dependencies. If , comparing 9 and 5 first yields . Attempting to compare pairs concurrently or out of order produces incorrect orderings.
Serial Odd-Even Transposition Sort
Section titled “Serial Odd-Even Transposition Sort”Odd-even transposition sort decouples compare-swaps into alternating independent phases:
- Even Phase: Compare-swaps are performed on even-indexed pairs:
- Odd Phase: Compare-swaps are performed on odd-indexed pairs:
Sorting [5, 9, 4, 3]:Start: 5, 9, 4, 3Even Phase: (5, 9) and (4, 3) -> 5, 9, 3, 4Odd Phase: (9, 3) -> 5, 3, 9, 4Even Phase: (5, 3) and (9, 4) -> 3, 5, 4, 9Odd Phase: (5, 4) -> 3, 4, 5, 9 (Sorted!)/* Program 3.15: Serial Odd-Even Transposition Sort */void Odd_even_sort(int a[], int n) { int phase, i, temp; for (phase = 0; phase < n; phase++) { if (phase % 2 == 0) { /* Even phase */ for (i = 1; i < n; i += 2) { if (a[i - 1] > a[i]) { temp = a[i]; a[i] = a[i - 1]; a[i - 1] = temp; } } } else { /* Odd phase */ for (i = 1; i < n - 1; i += 2) { if (a[i] > a[i + 1]) { temp = a[i]; a[i] = a[i + 1]; a[i + 1] = temp; } } } }}3.18 Parallelizing Odd-Even Transposition Sort
Section titled “3.18 Parallelizing Odd-Even Transposition Sort”Because all compare-swaps within a given phase are completely independent, they can occur concurrently across multiple processes.
flowchart TD
subgraph Tasks["Figure 3.12: Task Communication in Odd-Even Sort"]
direction TB
subgraph PhaseJ["Phase j"]
Aj_minus["Task a[i-1]"]
Aj["Task a[i]"]
Aj_plus["Task a[i+1]"]
end
subgraph PhaseJ1["Phase j + 1"]
Aj1_minus["Task a[i-1]"]
Aj1["Task a[i]"]
Aj1_plus["Task a[i+1]"]
end
Aj <--> Aj_plus
Aj <--> Aj_minus
Aj --> Aj1
Aj_minus --> Aj1_minus
Aj_plus --> Aj1_plus
endHandling Blocks of Keys ()
Section titled “Handling Blocks of Keys (n/p>1n / p > 1n/p>1)”Assigning a single key per process () is inefficient due to network latency. Instead, each process manages keys:
- Local Sort: Each process begins by sorting its local keys using an efficient sequential algorithm (e.g., C standard library
qsort). - Phase Exchanges: In each phase, active process pairs exchange their entire blocks of keys.
- Merge-Split:
- The process with the smaller rank retains the smallest elements from the combined keys.
- The process with the larger rank retains the largest elements.
Table 3.8: Parallel Odd-Even Transposition Sort Trace ()
Section titled “Table 3.8: Parallel Odd-Even Transposition Sort Trace (p=4,n=16p = 4, n = 16p=4,n=16)”| Phase / Step | Process 0 | Process 1 | Process 2 | Process 3 |
|---|---|---|---|---|
| Start | 15, 11, 9, 16 | 3, 14, 8, 7 | 4, 6, 12, 10 | 5, 2, 13, 1 |
| After Local Sort | 9, 11, 15, 16 | 3, 7, 8, 14 | 4, 6, 10, 12 | 1, 2, 5, 13 |
| After Phase 0 | 3, 7, 8, 9 | 11, 14, 15, 16 | 1, 2, 4, 5 | 6, 10, 12, 13 |
| After Phase 1 | 3, 7, 8, 9 | 1, 2, 4, 5 | 11, 14, 15, 16 | 6, 10, 12, 13 |
| After Phase 2 | 1, 2, 3, 4 | 5, 7, 8, 9 | 6, 10, 11, 12 | 13, 14, 15, 16 |
| After Phase 3 | 1, 2, 3, 4 | 5, 6, 7, 8 | 9, 10, 11, 12 | 13, 14, 15, 16 |
Computing Partner Ranks and Handling Idle Processes
Section titled “Computing Partner Ranks and Handling Idle Processes”During each phase, processes pair up to exchange data:
- In an even phase, odd ranks pair with
my_rank - 1, and even ranks pair withmy_rank + 1. - In an odd phase, odd ranks pair with
my_rank + 1, and even ranks pair withmy_rank - 1.
Boundary processes (e.g., rank 0 or rank ) may compute an invalid partner rank of or . When this occurs, the process is idle for that phase. In MPI, setting the destination rank to MPI_PROC_NULL ensures the communication call returns immediately without performing any action:
if (phase % 2 == 0) { /* Even phase */ if (my_rank % 2 != 0) partner = my_rank - 1; else partner = my_rank + 1;} else { /* Odd phase */ if (my_rank % 2 != 0) partner = my_rank + 1; else partner = my_rank - 1;}
if (partner == -1 || partner == comm_sz) { partner = MPI_PROC_NULL;}3.19 Safety in MPI Programs and Deadlock Prevention
Section titled “3.19 Safety in MPI Programs and Deadlock Prevention”If two partner processes execute communications naively:
/* DANGEROUS: May cause deadlock depending on buffer size */MPI_Send(my_keys, n / comm_sz, MPI_INT, partner, 0, comm);MPI_Recv(temp_keys, n / comm_sz, MPI_INT, partner, 0, comm, MPI_STATUS_IGNORE);Safe vs. Unsafe Programs
Section titled “Safe vs. Unsafe Programs”- Unsafe Program: A program whose correctness depends on the MPI implementation providing internal message buffering. If the message exceeds the system threshold, both processes block inside
MPI_Sendwaiting for the other to callMPI_Recv, causing a deadlock. - Safety Verification: You can verify whether an MPI program is safe by replacing all calls to
MPI_SendwithMPI_Ssend(Synchronous Send).MPI_Ssendis guaranteed to block until the matching receive has initiated. If the program completes without hanging underMPI_Ssend, the communication structure is provably safe.
Restructuring Point-to-Point Exchanges
Section titled “Restructuring Point-to-Point Exchanges”To make pairwise exchange safe without relying on buffers, communication must be scheduled so that some processes receive before sending:
flowchart TD
subgraph Alternating["Alternating Schedule (comm_sz even)"]
direction LR
subgraph Step1["Step 1"]
E1["Even Ranks: Call MPI_Send"] --> O1["Odd Ranks: Call MPI_Recv"]
end
subgraph Step2["Step 2"]
O2["Odd Ranks: Call MPI_Send"] --> E2["Even Ranks: Call MPI_Recv"]
end
Step1 --> Step2
endif (my_rank % 2 == 0) { MPI_Send(msg, size, MPI_INT, partner, 0, comm); MPI_Recv(new_msg, size, MPI_INT, partner, 0, comm, MPI_STATUS_IGNORE);} else { MPI_Recv(new_msg, size, MPI_INT, partner, 0, comm, MPI_STATUS_IGNORE); MPI_Send(msg, size, MPI_INT, partner, 0, comm);}Atomic Send-Receive: MPI_Sendrecv
Section titled “Atomic Send-Receive: MPI_Sendrecv”Rather than manually scheduling branches for odd and even ranks, MPI provides MPI_Sendrecv, which executes a blocking send and receive atomically:
int MPI_Sendrecv( void* send_buf_p /* in: send buffer */, int send_buf_size /* in: count to send */, MPI_Datatype send_buf_type /* in: send datatype */, int dest /* in: destination rank */, int send_tag /* in: send tag */, void* recv_buf_p /* out: receive buffer */, int recv_buf_size /* in: max count to receive */, MPI_Datatype recv_buf_type /* in: receive datatype */, int source /* in: source rank */, int recv_tag /* in: receive tag */, MPI_Comm communicator /* in: communicator */, MPI_Status* status_p /* out: status object */);The underlying MPI implementation automatically schedules network transfers to prevent deadlocks.
MPI_Sendrecv(my_keys, local_n, MPI_INT, partner, 0, recv_keys, local_n, MPI_INT, partner, 0, comm, MPI_STATUS_IGNORE);If the send and receive buffers occupy the same memory location, MPI_Sendrecv_replace can be used.
3.20 Merge-Split Optimization
Section titled “3.20 Merge-Split Optimization”When two processes exchange their sorted keys, they do not need to run a full sorting algorithm on all elements. Because both individual arrays are already sorted, they can be merged in linear time :
/* Program 3.16: Merge_low function */void Merge_low( int my_keys[], /* in/out: local keys */ int recv_keys[], /* in: keys received from partner */ int temp_keys[], /* scratch buffer of size local_n */ int local_n /* number of keys per process (n/p) */) { int m_i = 0, r_i = 0, t_i = 0;
/* Merge smallest local_n elements into temp_keys */ while (t_i < local_n) { if (my_keys[m_i] <= recv_keys[r_i]) { temp_keys[t_i++] = my_keys[m_i++]; } else { temp_keys[t_i++] = recv_keys[r_i++]; } }
/* Copy merged elements back into my_keys */ for (m_i = 0; m_i < local_n; m_i++) { my_keys[m_i] = temp_keys[m_i]; }}- For
Merge_high(retaining the largest keys), the merge simply iterates backwards starting fromlocal_n - 1. - Further performance gains can be achieved by swapping pointers between
my_keysandtemp_keysrather than copying elements back.
Empirical Benchmarks
Section titled “Empirical Benchmarks”Table 3.9: Run-times of Parallel Odd-Even Sort (ms)
Section titled “Table 3.9: Run-times of Parallel Odd-Even Sort (ms)”| Processes () | 200k Keys | 400k Keys | 800k Keys | 1.6M Keys | 3.2M Keys |
|---|---|---|---|---|---|
| 1 | 88 | 190 | 390 | 830 | 1800 |
| 2 | 43 | 91 | 190 | 410 | 860 |
| 4 | 22 | 46 | 96 | 200 | 430 |
| 8 | 12 | 24 | 51 | 110 | 220 |
| 16 | 7.5 | 14 | 29 | 60 | 130 |
Across all problem sizes, parallel odd-even sort achieves near-linear speedup, scaling efficiently as the dataset grows into millions of keys.
3.21 Chapter Summary
Section titled “3.21 Chapter Summary”- Message Passing Model: Independent processes on core-memory pairs communicate across an interconnect by calling library routines in MPI.
- SPMD Architecture: A single executable binary runs on all cores, branching on
my_rank. - Point-to-Point vs. Collective:
- Point-to-point (
MPI_Send,MPI_Recv) connects pairs of processes using tags and communicators. - Collective communication (
MPI_Bcast,MPI_Reduce,MPI_Allreduce,MPI_Scatter,MPI_Gather,MPI_Allgather) coordinates all processes within a communicator without tags.
- Point-to-point (
- Derived Datatypes: Eliminate latency overhead by packing non-contiguous, heterogeneous memory elements into a single message using
MPI_Type_create_struct. - Safety & Deadlock: Unsafe programs depend on internal buffering. Using
MPI_Sendrecvguarantees deadlock-free bidirectional communications. - Performance: Benchmarking requires wall-clock timing (
MPI_Wtime), synchronization (MPI_Barrier), and measuring speedup () and efficiency ().