Skip to content

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 p=comm_szp = \text{comm\_sz} processes and nn total keys, each process starts and finishes with a local block of n/pn / p keys (assuming pp evenly divides nn).

The algorithm terminates when the following conditions are satisfied:

  1. Local Order: The n/pn / p keys residing on each process are sorted in non-decreasing order.
  2. Global Order: If 0≤q<r<p0 \le q < r < p, then every key assigned to process qq is less than or equal to every key assigned to process rr:

∀x∈Keys(q),  ∀y∈Keys(r):x≤y\forall x \in \text{Keys}(q), \; \forall y \in \text{Keys}(r): \quad x \le y

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”

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 a=[9,5,7]a = [9, 5, 7], comparing 9 and 5 first yields [5,9,7][5, 9, 7]. Attempting to compare pairs concurrently or out of order produces incorrect orderings.


Odd-even transposition sort decouples compare-swaps into alternating independent phases:

  • Even Phase: Compare-swaps are performed on even-indexed pairs: (a[0],a[1]),(a[2],a[3]),(a[4],a[5]),…(a[0], a[1]), \quad (a[2], a[3]), \quad (a[4], a[5]), \quad \dots
  • Odd Phase: Compare-swaps are performed on odd-indexed pairs: (a[1],a[2]),(a[3],a[4]),(a[5],a[6]),…(a[1], a[2]), \quad (a[3], a[4]), \quad (a[5], a[6]), \quad \dots
Sorting [5, 9, 4, 3]:
Start: 5, 9, 4, 3
Even Phase: (5, 9) and (4, 3) -> 5, 9, 3, 4
Odd Phase: (9, 3) -> 5, 3, 9, 4
Even Phase: (5, 3) and (9, 4) -> 3, 5, 4, 9
Odd 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
  end

Handling Blocks of Keys (n/p>1n / p > 1)

Section titled “Handling Blocks of Keys (n/p>1n / p > 1n/p>1)”

Assigning a single key per process (n=pn = p) is inefficient due to network latency. Instead, each process manages n/pn / p keys:

  1. Local Sort: Each process begins by sorting its local n/pn / p keys using an efficient sequential algorithm (e.g., C standard library qsort).
  2. Phase Exchanges: In each phase, active process pairs exchange their entire blocks of n/pn / p keys.
  3. Merge-Split:
    • The process with the smaller rank retains the n/pn / p smallest elements from the combined 2n/p2n / p keys.
    • The process with the larger rank retains the n/pn / p largest elements.

Table 3.8: Parallel Odd-Even Transposition Sort Trace (p=4,n=16p = 4, n = 16)

Section titled “Table 3.8: Parallel Odd-Even Transposition Sort Trace (p=4,n=16p = 4, n = 16p=4,n=16)”
Phase / StepProcess 0Process 1Process 2Process 3
Start15, 11, 9, 163, 14, 8, 74, 6, 12, 105, 2, 13, 1
After Local Sort9, 11, 15, 163, 7, 8, 144, 6, 10, 121, 2, 5, 13
After Phase 03, 7, 8, 911, 14, 15, 161, 2, 4, 56, 10, 12, 13
After Phase 13, 7, 8, 91, 2, 4, 511, 14, 15, 166, 10, 12, 13
After Phase 21, 2, 3, 45, 7, 8, 96, 10, 11, 1213, 14, 15, 16
After Phase 31, 2, 3, 45, 6, 7, 89, 10, 11, 1213, 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 with my_rank + 1.
  • In an odd phase, odd ranks pair with my_rank + 1, and even ranks pair with my_rank - 1.

Boundary processes (e.g., rank 0 or rank p−1p - 1) may compute an invalid partner rank of −1-1 or pp. 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);
  • 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_Send waiting for the other to call MPI_Recv, causing a deadlock.
  • Safety Verification: You can verify whether an MPI program is safe by replacing all calls to MPI_Send with MPI_Ssend (Synchronous Send). MPI_Ssend is guaranteed to block until the matching receive has initiated. If the program completes without hanging under MPI_Ssend, the communication structure is provably safe.

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
  end
if (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);
}

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.


When two processes exchange their n/pn / p sorted keys, they do not need to run a full sorting algorithm on all 2n/p2n / p elements. Because both individual arrays are already sorted, they can be merged in linear time O(n/p)\mathcal{O}(n / p):

/* 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 from local_n - 1.
  • Further performance gains can be achieved by swapping pointers between my_keys and temp_keys rather than copying elements back.

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 (pp)200k Keys400k Keys800k Keys1.6M Keys3.2M Keys
1881903908301800
24391190410860
4224696200430
8122451110220
167.5142960130

Across all problem sizes, parallel odd-even sort achieves near-linear speedup, scaling efficiently as the dataset grows into millions of keys.


  • 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.
  • 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_Sendrecv guarantees deadlock-free bidirectional communications.
  • Performance: Benchmarking requires wall-clock timing (MPI_Wtime), synchronization (MPI_Barrier), and measuring speedup (S=Tserial/TparallelS = T_{\text{serial}} / T_{\text{parallel}}) and efficiency (E=S/pE = S / p).