Kwker

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

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

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

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:

  1. Each machine takes a sample of its data with sample.
  2. The samples are combined, and splitters picks the boundaries that divide the work evenly.
  3. Each machine splits its data at those boundaries (partition_splitters) and sends each part to its machine.
  4. 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.