多方向マージソートを使用する

多方向マージソート (multiway merge sort / k-way merge sort) は、通常のマージソートが 2 本の整列済み列をマージするのに対し、一度に k 本(本稿では k = 4)をまとめてマージする分割統治法である。

外部整列ではテープやファイル本数に応じて k を選び、初期ランを k 本ずつまとめていく用途でも同じ部品が使われる。内部メモリ向けでも、マージ段数はおよそ \(\log_k n\) に落ち、各段のコストは線形なので全体は \(\Theta(n \log n)\) を保つ。

  1. 分割: 区間を最大 k 個の連続部分にほぼ等分する。要素が 1 つ以下ならそのままソート済みとみなす。
  2. 再帰: 各部分に対して同じ手順を繰り返す。
  3. k 方向マージ: 各部分は昇順になっている前提で、各ランの「先頭」を比較し、最小(同値ならより左のラン)を出力へ確定する。選んだランの先頭を 1 つ進め、全ランが尽きるまで繰り返す。
  4. 書き戻し: マージ結果を元の区間へ写す。

本稿のデモとベンチマークは、先頭比較を長さ k の線形走査で行う単純実装である。k が大きい外部マージでは、同じ選択を敗者木やヒープで \(O(\log k)\) に落とすことが多い。

procedure merge_k_way(runs[0..k))
  heads[i] = 0 for each run i
  while some run still has unread elements
    pick run i with smallest heads[i] value
      (ties: smallest i, for stability)
    append runs[i][heads[i]] to output
    heads[i] = heads[i] + 1
  return output

procedure multiway_merge_sort(A)
  n = length(A)
  if n <= 1 then
    return
  split A into up to k contiguous parts of nearly equal length
  for each part P
    multiway_merge_sort(P)
  merged = merge_k_way(the k sorted parts)
  copy merged back into A

分割の深さが \(O(\log_k n)\)、各層のマージが \(O(k n)\)(固定 k なら \(O(n)\))なので最悪計算量は \(O(n \log n)\) である。作業用バッファに \(O(n)\) の追加領域が要る。同値を左ラン優先で取れば安定ソートになる。

類似アルゴリズムとの相違点

マージソートは常に 2 方向マージであり、本稿はその k 一般化である。

ポリフェーズマージソートカスケードマージソートは、テープ本数が限られた外部整列向けにラン分布やパス構成を工夫する系統である。

クアッドソートは 4 本をボトムアップでマージする実装の一例だが、適応的な枝刈りや交換ネットワークなど、本稿の素直な再帰 k 方向マージとは別の工夫を含む。

敗者木ソートは要素全体をトーナメントにする整列本体であり、多方向マージの「k 本の先頭から最小を取る」部品としても使われる。

計算時間量および空間計算量を計測する

Size Average time Maximum time Average memory Maximum memory
256 0.000018 0.000247 4 4
512 0.000047 0.000263 8 8
1024 0.000077 0.000176 16 16
2048 0.000194 0.000713 32 32
4096 0.000335 0.000847 64 64
8192 0.000855 0.001521 128 128
16384 0.001479 0.002574 256 256
32768 0.004319 0.010915 512 512
65536 0.009095 0.023376 1024 1024
131072 0.019601 0.083974 2048 2048
262144 0.034262 0.070343 4096 4096
計測に使用したコードを表示する

set -euo pipefail

WORKDIR="$(mktemp -d)"
trap 'rm -rf "$WORKDIR"' EXIT

cat > "$WORKDIR/Dockerfile" <<'EOF'
FROM rust:1.95.0

WORKDIR /app

RUN mkdir -p src

RUN cat > Cargo.toml <<'CARGO'
[package]
name = "rust-benchmark"
version = "0.1.0"
edition = "2021"

[profile.release]
lto = true
codegen-units = 1
panic = "abort"
CARGO

RUN cat > src/main.rs <<'RUST'
use std::{
    alloc::{GlobalAlloc, Layout, System},
    env,
    process::Command,
    sync::atomic::{AtomicUsize, Ordering as AtomicOrdering},
    time::{Duration, Instant},
};

/// Counts live heap bytes and the high-water mark so auxiliary sort buffers
/// (swap Vecs, etc.) are measured as explicit heap growth during the sort.
struct TrackingAllocator;

static LIVE_BYTES: AtomicUsize = AtomicUsize::new(0);
static PEAK_BYTES: AtomicUsize = AtomicUsize::new(0);

fn record_alloc(size: usize) {
    let live = LIVE_BYTES.fetch_add(size, AtomicOrdering::Relaxed) + size;
    PEAK_BYTES.fetch_max(live, AtomicOrdering::Relaxed);
}

unsafe impl GlobalAlloc for TrackingAllocator {
    unsafe fn alloc(&self, layout: Layout) -> *mut u8 {
        let ptr = System.alloc(layout);
        if !ptr.is_null() {
            record_alloc(layout.size());
        }
        ptr
    }

    unsafe fn dealloc(&self, ptr: *mut u8, layout: Layout) {
        LIVE_BYTES.fetch_sub(layout.size(), AtomicOrdering::Relaxed);
        System.dealloc(ptr, layout);
    }

    unsafe fn alloc_zeroed(&self, layout: Layout) -> *mut u8 {
        let ptr = System.alloc_zeroed(layout);
        if !ptr.is_null() {
            record_alloc(layout.size());
        }
        ptr
    }

    unsafe fn realloc(&self, ptr: *mut u8, layout: Layout, new_size: usize) -> *mut u8 {
        let new_ptr = System.realloc(ptr, layout, new_size);
        if !new_ptr.is_null() {
            LIVE_BYTES.fetch_sub(layout.size(), AtomicOrdering::Relaxed);
            record_alloc(new_size);
        }
        new_ptr
    }
}

#[global_allocator]
static GLOBAL: TrackingAllocator = TrackingAllocator;
const MIN_POWER: u32 = 8;
const MAX_POWER: u32 = 18;
const RUNS: usize = 8192;


/// Branching factor for multiway (k-way) merge sort. Fixed for the pedagogical
/// benchmark so asymptotics stay Θ(n log n) with a constant-factor log_k.
const WAY: usize = 4;

/// Stable k-way merge: when heads compare equal, the leftmost run wins.
fn merge_k_way(runs: &[&[usize]]) -> Vec<usize> {
    let k = runs.len();
    let mut heads = vec![0usize; k];
    let total: usize = runs.iter().map(|r| r.len()).sum();
    let mut out = Vec::with_capacity(total);

    loop {
        let mut best: Option<(usize, usize)> = None; // (run_index, value)
        for i in 0..k {
            if heads[i] < runs[i].len() {
                let v = runs[i][heads[i]];
                match best {
                    None => best = Some((i, v)),
                    Some((bi, bv)) => {
                        if v < bv || (v == bv && i < bi) {
                            best = Some((i, v));
                        }
                    }
                }
            }
        }
        match best {
            None => break,
            Some((i, v)) => {
                out.push(v);
                heads[i] += 1;
            }
        }
    }
    out
}

fn multiway_merge_sort(a: &mut [usize]) {
    let n = a.len();
    if n <= 1 {
        return;
    }

    let mut bounds: Vec<(usize, usize)> = Vec::with_capacity(WAY);
    let base = n / WAY;
    let rem = n % WAY;
    let mut start = 0usize;
    for i in 0..WAY {
        let len = base + usize::from(i < rem);
        if len == 0 {
            continue;
        }
        let end = start + len;
        multiway_merge_sort(&mut a[start..end]);
        bounds.push((start, end));
        start = end;
    }

    if bounds.len() <= 1 {
        return;
    }

    let runs: Vec<Vec<usize>> = bounds
        .iter()
        .map(|&(lo, hi)| a[lo..hi].to_vec())
        .collect();
    let refs: Vec<&[usize]> = runs.iter().map(|r| r.as_slice()).collect();
    let merged = merge_k_way(&refs);
    a.copy_from_slice(&merged);
}


fn benchmark_sort(array: &mut [usize]) {

    multiway_merge_sort(array);

}

fn is_non_decreasing(a: &[usize]) -> bool {
    a.windows(2).all(|w| w[0] <= w[1])
}

fn same_multiset(a: &[usize], b: &[usize]) -> bool {
    if a.len() != b.len() {
        return false;
    }

    let mut left = a.to_vec();
    let mut right = b.to_vec();
    left.sort_unstable();
    right.sort_unstable();
    left == right
}

fn check_correctness_case(label: &str, mut input: Vec<usize>) {
    let original = input.clone();

    benchmark_sort(&mut input);

    if !is_non_decreasing(&input) {
        panic!("correctness case {}: output is not sorted", label);
    }

    if !same_multiset(&input, &original) {
        panic!("correctness case {}: elements were lost or added", label);
    }
}

// Skip cases larger than the algorithm's measured size cap (MAX_POWER). That
// cap exists because larger inputs are impractically slow; forcing them here
// would stall the published measurement script before any table rows print.
fn check_correctness_case_within_limit(label: &str, input: Vec<usize>) {
    if input.len() > (1usize << MAX_POWER) {
        return;
    }
    check_correctness_case(label, input);
}

fn few_unique_values(size: usize, unique: usize, seed: u64) -> Vec<usize> {
    let mut state = seed;

    (0..size)
        .map(|_| {
            state ^= state << 13;
            state ^= state >> 7;
            state ^= state << 17;
            (state as usize % unique) + 1
        })
        .collect()
}

fn run_correctness_checks() {
    check_correctness_case("empty", vec![]);
    check_correctness_case("single", vec![42]);
    check_correctness_case("duplicates", vec![3, 1, 3, 2, 1, 2]);
    check_correctness_case("sorted", vec![1, 2, 3, 4, 5]);
    check_correctness_case("reverse", vec![5, 4, 3, 2, 1]);
    check_correctness_case("all_equal", vec![7, 7, 7, 7]);
    check_correctness_case("skewed_range", vec![1_000_000, 2, 1_000_001, 1, 999_999]);
    // Static-buffer Grail skips the in-buffer build when key collection is sparse
    // (ideal_buffer = false). Exercising that path catches regressions in buffer gating.
    check_correctness_case(
        "few_keys_len16",
        vec![2, 2, 2, 2, 2, 2, 2, 2, 4, 3, 1, 2, 3, 4, 1, 4],
    );
    // Seed 0 is a fixed point of the xorshift below, so it would degenerate into
    // yet another all-equal case instead of a 4-value mix. Start at 1.
    for seed in 1..=32 {
        check_correctness_case(
            &format!("few_keys_len32_seed_{seed}"),
            few_unique_values(32, 4, seed),
        );
    }
    // Small-input cutoffs (insertion sort below 32 elements, etc.) hide duplicate-key
    // bugs in the recursive path, so repeat the duplicate cases at the smallest
    // benchmark size, which every algorithm must handle within reasonable time.
    check_correctness_case("all_equal_len256", vec![7; 256]);
    for seed in 1..=4 {
        check_correctness_case(
            &format!("few_keys_len256_seed_{seed}"),
            few_unique_values(256, 4, seed),
        );
    }
    // Blit's equal-key second sweep used to copy the whole range into a fixed
    // 512-element swap; lengths above that must still sort without panicking.
    // Respect MAX_POWER so algorithms with a low measured-size cap (slow,
    // sleep) do not hang here for minutes or months.
    check_correctness_case_within_limit("all_equal_len600", vec![7; 600]);
    for seed in 1..=4 {
        check_correctness_case_within_limit(
            &format!("few_keys_len2048_seed_{seed}"),
            few_unique_values(2048, 4, seed),
        );
    }
}


fn shuffled(size: usize, seed: u64) -> Vec<usize> {
    let mut v: Vec<usize> = (1..=size).collect();

    let mut state = seed;

    for i in (1..size).rev() {
        state ^= state << 13;
        state ^= state >> 7;
        state ^= state << 17;

        let j = (state as usize) % (i + 1);

        v.swap(i, j);
    }

    v
}

fn micros(d: Duration) -> u128 {
    d.as_micros()
}

fn input_array(size: usize, seed: u64) -> Vec<usize> {
    shuffled(size, seed)
}

/// Peak heap growth during `benchmark_sort`, in bytes (explicit buffers such as swap).
/// Kept in bytes so the parent can average before rounding; converting to KiB here
/// would truncate sub-KiB buffers to 0 in every run and hide them from the average.
fn run_once(size: usize, seed: usize) -> (u128, usize) {
    let mut array = input_array(size, seed as u64);

    let base_bytes = LIVE_BYTES.load(AtomicOrdering::Relaxed);
    PEAK_BYTES.store(base_bytes, AtomicOrdering::Relaxed);

    let start = Instant::now();

    benchmark_sort(&mut array);

    let elapsed = start.elapsed();
    let peak_bytes = PEAK_BYTES.load(AtomicOrdering::Relaxed);
    let aux_bytes = peak_bytes.saturating_sub(base_bytes);

    let expected: Vec<usize> = (1..=size).collect();
    if array != expected {
        panic!(
            "sort failed with seed {} for size {}",
            seed,
            size
        );
    }

    (micros(elapsed), aux_bytes)
}

fn run_child(args: &[String]) {
    let size = args[2].parse::<usize>().expect("invalid size");
    let seed = args[3].parse::<usize>().expect("invalid seed");
    let (elapsed_us, mem) = run_once(size, seed);
    println!("{} {}", elapsed_us, mem);
}

fn main() {
    let args: Vec<String> = env::args().collect();
    if args.get(1).is_some_and(|arg| arg == "--run-once") {
        run_child(&args);
        return;
    }

    run_correctness_checks();

    println!(
        "| {:>10} | {:>15} | {:>15} | {:>15} | {:>15} |",
        "Size",
        "Average time",
        "Maximum time",
        "Average memory",
        "Maximum memory"
    );

    println!(
        "|{:-<11}:|{:-<16}:|{:-<16}:|{:-<16}:|{:-<16}:|",
        "",
        "",
        "",
        "",
        ""
    );

    for power in MIN_POWER..=MAX_POWER {
        let size = 1usize << power;

        let mut total_time: u128 = 0;
        let mut max_time: u128 = 0;

        let mut total_mem: usize = 0;
        let mut max_mem: usize = 0;

        for seed in 1..=RUNS {
            let output = Command::new(env::current_exe().expect("failed to find current executable"))
                .arg("--run-once")
                .arg(size.to_string())
                .arg(seed.to_string())
                .output()
                .expect("failed to run benchmark child process");

            if !output.status.success() {
                panic!(
                    "benchmark child process failed: {}",
                    String::from_utf8_lossy(&output.stderr)
                );
            }

            let stdout = String::from_utf8(output.stdout)
                .expect("child process returned non-UTF-8 output");
            let mut fields = stdout.split_whitespace();
            let elapsed_us = fields
                .next()
                .expect("missing elapsed time")
                .parse::<u128>()
                .expect("invalid elapsed time");
            let aux_mem = fields
                .next()
                .expect("missing memory usage")
                .parse::<usize>()
                .expect("invalid memory usage");

            total_time += elapsed_us;

            if elapsed_us > max_time {
                max_time = elapsed_us;
            }

            total_mem += aux_mem;

            if aux_mem > max_mem {
                max_mem = aux_mem;
            }
        }

        let avg_time = total_time / RUNS as u128;
        // Memory is summed in bytes and converted to KiB once, after averaging.
        let avg_mem_kb = total_mem / RUNS / 1024;
        let max_mem_kb = max_mem / 1024;

        println!(
            "| {:>10} | {:>15} | {:>15} | {:>15} | {:>15} |",
            size,
            format!("{}.{:06}", avg_time / 1_000_000, avg_time % 1_000_000),
            format!("{}.{:06}", max_time / 1_000_000, max_time % 1_000_000),
            avg_mem_kb,
            max_mem_kb
        );
    }
}
RUST

RUN cargo build --release

CMD ["./target/release/rust-benchmark"]
EOF

docker build -t rust-benchmark "$WORKDIR"
docker run --rm --init rust-benchmark