Large data
This page covers threads, memory limits, files larger than memory, and the building blocks for sorting across machines.
Use several cores
Sorts, argsorts and key-value sorts can run on several threads: threads= in Python, the _mt calls in Rust, C and
C++. Passing 0 threads uses the default count for your machine. The result is always exactly the same as with one
thread, so you can turn threads on without changing any test.
import numpy as np
import kwker
a = np.random.default_rng(7).integers(0, 1_000_000, 4_000_000)
one = kwker.argsort(a)
four = kwker.argsort(a, threads=4)
print(np.array_equal(one, four))
True
use kwker::Order;
fn main() {
// 4 million pseudo-random keys below 1,000,000
let mut x = 7u64;
let a: Vec<i64> = (0..4_000_000)
.map(|_| {
x = x.wrapping_mul(6364136223846793005).wrapping_add(1442695040888963407);
((x >> 33) % 1_000_000) as i64
})
.collect();
let one: Vec<u64> = kwker::argsort(&a, Order::ASCENDING);
let mut four = vec![0u64; a.len()];
kwker::sort_indexed_mt(&a, Order::ASCENDING, None, &mut four, 4);
println!("{}", one == four);
}
true
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <kwker.h>
int main(void) {
size_t n = 4000000;
int64_t* a = malloc(n * sizeof(int64_t));
uint64_t* one = malloc(n * sizeof(uint64_t));
uint64_t* four = malloc(n * sizeof(uint64_t));
uint64_t x = 7;
for (size_t i = 0; i < n; i++) {
x = x * 6364136223846793005u + 1442695040888963407u;
a[i] = (int64_t)((x >> 33) % 1000000);
}
kwker_i64_argsort(a, n, KWKER_ASCENDING, one);
kwker_i64_sort_indexed_mt(a, n, KWKER_ASCENDING, NULL, four, 4);
printf("%s\n", memcmp(one, four, n * sizeof(uint64_t)) == 0 ? "true" : "false");
free(a);
free(one);
free(four);
return 0;
}
true
#include <cstdint>
#include <iostream>
#include <vector>
#include <kwker.hpp>
int main() {
std::vector<int64_t> a(4000000);
uint64_t x = 7;
for (auto& v : a) {
x = x * 6364136223846793005u + 1442695040888963407u;
v = int64_t((x >> 33) % 1000000);
}
auto one = kwker::argsort(a);
auto four = kwker::argsort_mt(a, 4);
std::cout << std::boolalpha << (one == four) << '\n';
}
true
package main
import (
"fmt"
"slices"
"kwker.io/go/kwker"
)
func main() {
// 4 million pseudo-random keys below 1,000,000
a := make([]int64, 4_000_000)
x := uint64(7)
for i := range a {
x = x*6364136223846793005 + 1442695040888963407
a[i] = int64((x >> 33) % 1_000_000)
}
one := slices.Clone(a)
kwker.Sort(one)
kwker.SortMT(a, kwker.Ascending, 4) // 4 threads
fmt.Println(slices.Equal(one, a))
}
true
import io.kwker.Kwker;
import java.util.Arrays;
public class Example {
public static void main(String[] args) {
// 4 million pseudo-random keys below 1,000,000
long[] a = new long[4_000_000];
long x = 7;
for (int i = 0; i < a.length; i++) {
x = x * 6364136223846793005L + 1442695040888963407L;
a[i] = (x >>> 33) % 1_000_000;
}
long[] one = a.clone();
Kwker.sort(one);
Kwker.sortMt(a, Kwker.ASCENDING, 4); // 4 threads
System.out.println(Arrays.equals(one, a));
}
}
true
using Kwker;
// 4 million pseudo-random keys below 1,000,000
var a = new long[4_000_000];
ulong x = 7;
for (int i = 0; i < a.Length; i++)
{
x = x * 6364136223846793005UL + 1442695040888963407UL;
a[i] = (long)((x >> 33) % 1_000_000);
}
var one = (long[])a.Clone();
Sorter.Sort(one);
Sorter.SortMT(a, Order.Ascending, threads: 4);
Console.WriteLine(one.SequenceEqual(a));
True
More threads pay off for large arrays, from millions of elements. For small arrays one core is usually fastest.
Limit the extra memory
Some fast paths use scratch memory next to your array. set_scratch_limit(bytes) caps it for the calling thread.
Results never change; some inputs get slower. With a limit of 0,
sort, select, partial_sort and sort_kv allocate nothing at all, which suits real-time code and tight memory
budgets. No limit (None in Python and Rust, SIZE_MAX in C and C++) removes it again.
import numpy as np
import kwker
a = np.random.default_rng(3).integers(0, 100, 1_000_000)
expected = np.sort(a)
kwker.set_scratch_limit(0)
kwker.sort(a)
kwker.set_scratch_limit(None)
print(np.array_equal(a, expected), kwker.scratch_limit())
True None
fn main() {
let mut x = 3u64;
let mut a: Vec<i64> = (0..1_000_000)
.map(|_| {
x = x.wrapping_mul(6364136223846793005).wrapping_add(1442695040888963407);
((x >> 33) % 100) as i64
})
.collect();
let mut expected = a.clone();
expected.sort_unstable();
kwker::set_scratch_limit(Some(0)); // no allocation from here on
kwker::sort(&mut a);
kwker::set_scratch_limit(None);
println!("{} {:?}", a == expected, kwker::scratch_limit());
}
true None
#include <stdint.h>
#include <stdio.h>
#include <stdlib.h>
#include <kwker.h>
int main(void) {
size_t n = 1000000;
int64_t* a = malloc(n * sizeof(int64_t));
uint64_t x = 3;
for (size_t i = 0; i < n; i++) {
x = x * 6364136223846793005u + 1442695040888963407u;
a[i] = (int64_t)((x >> 33) % 100);
}
kwker_set_scratch_limit(0); /* no allocation from here on */
kwker_i64_sort(a, n);
kwker_set_scratch_limit(SIZE_MAX);
int sorted = 1;
for (size_t i = 1; i < n; i++) sorted &= a[i - 1] <= a[i];
printf("%s %s\n", sorted ? "true" : "false", kwker_scratch_limit() == SIZE_MAX ? "unlimited" : "limited");
free(a);
return 0;
}
true unlimited
#include <algorithm>
#include <cstdint>
#include <iostream>
#include <vector>
#include <kwker.hpp>
int main() {
std::vector<int64_t> a(1000000);
uint64_t x = 3;
for (auto& v : a) {
x = x * 6364136223846793005u + 1442695040888963407u;
v = int64_t((x >> 33) % 100);
}
kwker::set_scratch_limit(0); // no allocation from here on
kwker::sort(a);
kwker::set_scratch_limit(SIZE_MAX);
std::cout << std::boolalpha << std::is_sorted(a.begin(), a.end()) << ' '
<< (kwker::scratch_limit() == SIZE_MAX ? "unlimited" : "limited") << '\n';
}
true unlimited
package main
import (
"fmt"
"slices"
"kwker.io/go/kwker"
)
func main() {
a := make([]int64, 1_000_000)
x := uint64(3)
for i := range a {
x = x*6364136223846793005 + 1442695040888963407
a[i] = int64((x >> 33) % 100)
}
expected := slices.Clone(a)
slices.Sort(expected)
kwker.WithScratchLimit(0, func() { kwker.Sort(a) }) // no extra memory inside
fmt.Println(slices.Equal(a, expected))
}
true
import io.kwker.Kwker;
import java.util.Arrays;
public class Example {
public static void main(String[] args) {
long[] a = new long[1_000_000];
long x = 3;
for (int i = 0; i < a.length; i++) {
x = x * 6364136223846793005L + 1442695040888963407L;
a[i] = (x >>> 33) % 100;
}
long[] expected = a.clone();
Arrays.sort(expected);
Kwker.setScratchLimit(0); // this thread: no extra memory
Kwker.sort(a);
Kwker.setScratchLimit(-1); // no limit again
System.out.println(Arrays.equals(a, expected) + " " + Kwker.scratchLimit());
}
}
true -1
using Kwker;
var a = new long[1_000_000];
ulong x = 3;
for (int i = 0; i < a.Length; i++)
{
x = x * 6364136223846793005UL + 1442695040888963407UL;
a[i] = (long)((x >> 33) % 100);
}
var expected = (long[])a.Clone();
Array.Sort(expected);
Sorter.ScratchLimit = 0; // this thread: no extra memory
Sorter.Sort(a);
Sorter.ScratchLimit = null; // no limit again
Console.WriteLine($"{a.SequenceEqual(expected)} {Sorter.ScratchLimit?.ToString() ?? "none"}");
True none
Sort a file larger than memory
sort_file(input, output) sorts a binary file of numbers and writes the result to output. The memory setting says
how much RAM it may use; the rest goes through temporary files. The input and output can be the same file.
import os
import tempfile
import numpy as np
import kwker
path = os.path.join(tempfile.mkdtemp(), "keys.u64")
np.random.default_rng(4).integers(0, 2**63, 2_000_000, dtype=np.uint64).tofile(path)
kwker.sort_file(path, path, np.uint64, memory=4 << 20) # a 16 MB file with 4 MB of memory
k = np.fromfile(path, dtype=np.uint64)
print(len(k), bool(np.all(k[:-1] <= k[1:])))
2000000 True
use kwker::ExternalSort;
use std::path::Path;
fn main() {
// a 16 MB file of 2 million u64 keys
let mut x = 4u64;
let bytes: Vec<u8> = (0..2_000_000)
.flat_map(|_| {
x = x.wrapping_mul(6364136223846793005).wrapping_add(1442695040888963407);
(x >> 1).to_le_bytes()
})
.collect();
std::fs::write("keys.u64", &bytes).unwrap();
let path = Path::new("keys.u64");
let opts = ExternalSort { memory: 4 << 20, ..Default::default() }; // 4 MB of memory
kwker::sort_file::<u64>(path, path, &opts).unwrap();
let data = std::fs::read(path).unwrap();
let k: Vec<u64> = data.chunks(8).map(|c| u64::from_le_bytes(c.try_into().unwrap())).collect();
println!("{} {}", k.len(), k.windows(2).all(|w| w[0] <= w[1]));
std::fs::remove_file(path).unwrap();
}
2000000 true
#include <stdio.h>
#include <stdlib.h>
#include <kwker.h>
int main(void) {
size_t n = 2000000; /* a 16 MB file of uint64_t keys */
uint64_t* k = malloc(n * sizeof(uint64_t));
uint64_t x = 4;
for (size_t i = 0; i < n; i++) {
x = x * 6364136223846793005u + 1442695040888963407u;
k[i] = x >> 1;
}
FILE* f = fopen("keys.u64", "wb");
fwrite(k, sizeof(uint64_t), n, f);
fclose(f);
kwker_u64_sort_file("keys.u64", "keys.u64", 4 << 20, NULL, KWKER_ASCENDING, 1); /* 4 MB of memory */
f = fopen("keys.u64", "rb");
size_t got = fread(k, sizeof(uint64_t), n, f);
fclose(f);
int sorted = 1;
for (size_t i = 1; i < got; i++) sorted &= k[i - 1] <= k[i];
printf("%zu %s\n", got, sorted ? "true" : "false");
remove("keys.u64");
free(k);
return 0;
}
2000000 true
#include <algorithm>
#include <cstdint>
#include <cstdio>
#include <fstream>
#include <iostream>
#include <vector>
#include <kwker.hpp>
int main() {
std::vector<uint64_t> k(2000000); // a 16 MB file of uint64_t keys
uint64_t x = 4;
for (auto& v : k) {
x = x * 6364136223846793005u + 1442695040888963407u;
v = x >> 1;
}
std::ofstream("keys.u64", std::ios::binary).write(reinterpret_cast<const char*>(k.data()), k.size() * 8);
kwker::sort_file<uint64_t>("keys.u64", "keys.u64", kwker::Order::ascending, 4 << 20); // 4 MB of memory
std::ifstream("keys.u64", std::ios::binary).read(reinterpret_cast<char*>(k.data()), k.size() * 8);
std::cout << k.size() << ' ' << std::boolalpha << std::is_sorted(k.begin(), k.end()) << '\n';
std::remove("keys.u64");
}
2000000 true
package main
import (
"encoding/binary"
"fmt"
"os"
"path/filepath"
"slices"
"kwker.io/go/kwker"
)
func main() {
path := filepath.Join(os.TempDir(), "keys.u64")
k := make([]uint64, 2_000_000) // a 16 MB file of uint64 keys
x := uint64(4)
buf := make([]byte, 0, len(k)*8)
for range k {
x = x*6364136223846793005 + 1442695040888963407
buf = binary.NativeEndian.AppendUint64(buf, x>>1)
}
os.WriteFile(path, buf, 0o644)
err := kwker.SortFile[uint64](path, path, 4<<20, kwker.Ascending, 1) // 4 MB of memory
data, _ := os.ReadFile(path)
for i := range k {
k[i] = binary.NativeEndian.Uint64(data[8*i:])
}
fmt.Println(err, len(k), slices.IsSorted(k))
os.Remove(path)
}
<nil> 2000000 true
import io.kwker.Kwker;
import java.nio.ByteBuffer;
import java.nio.ByteOrder;
import java.nio.file.Files;
import java.nio.file.Path;
public class Example {
public static void main(String[] args) throws Exception {
Path path = Path.of(System.getProperty("java.io.tmpdir"), "kwker-keys.i64");
long[] k = new long[2_000_000]; // a 16 MB file of long keys
ByteBuffer buf = ByteBuffer.allocate(k.length * 8).order(ByteOrder.nativeOrder());
long x = 4;
for (int i = 0; i < k.length; i++) {
x = x * 6364136223846793005L + 1442695040888963407L;
buf.putLong(x >>> 1);
}
Files.write(path, buf.array());
Kwker.sortFile(long.class, path.toString(), path.toString(), 4 << 20, Kwker.ASCENDING, 1); // 4 MB of memory
ByteBuffer.wrap(Files.readAllBytes(path)).order(ByteOrder.nativeOrder()).asLongBuffer().get(k);
boolean sorted = true;
for (int i = 1; i < k.length; i++) sorted &= k[i - 1] <= k[i];
System.out.println(k.length + " " + sorted);
Files.delete(path);
}
}
2000000 true
using System.Runtime.InteropServices;
using Kwker;
var path = Path.Combine(Path.GetTempPath(), "kwker-keys.u64");
var k = new ulong[2_000_000]; // a 16 MB file of ulong keys
ulong x = 4;
for (int i = 0; i < k.Length; i++)
{
x = x * 6364136223846793005UL + 1442695040888963407UL;
k[i] = x >> 1;
}
File.WriteAllBytes(path, MemoryMarshal.AsBytes(k.AsSpan()).ToArray());
Sorter.SortFile<ulong>(path, path, memory: 4 << 20, threads: 1); // 4 MB of memory
var back = MemoryMarshal.Cast<byte, ulong>(File.ReadAllBytes(path));
bool sorted = true;
for (int i = 1; i < back.Length; i++) sorted &= back[i - 1] <= back[i];
Console.WriteLine($"{back.Length} {sorted}");
File.Delete(path);
2000000 True
sort_file also sorts fixed-size records by one field (key="field" for NumPy structured types, sort_file_records
in Rust, C and C++).
Sorting across machines
Kwker does not move data between machines, but it provides the pieces a distributed sort needs:
- Each machine takes a sample of its data with
sample. - The samples are combined, and
splitterspicks the boundaries that divide the work evenly. - Each machine splits its data at those boundaries (
partition_splitters) and sends each part to its machine. - Each machine merges what it receives with
kway_merge.
Equal values at a boundary go to one side by a fixed, documented rule (splitters_exact), so every machine agrees on
where each value belongs.
Related
- Runtime controls: engines, threads and memory in detail.
- API reference:
sort_file,set_scratch_limit