From 6c2f23cfa3e5ac1829ac21c2262bad96f9e61524 Mon Sep 17 00:00:00 2001 From: Thomas Waldmann Date: Thu, 3 Sep 2026 20:49:32 +0200 Subject: [PATCH] compress: reuse the zstd compressor per thread libzstd starts a worker pool when a compression context with nb_workers > 1 is set up and stops it again when that context is freed, so the one-shot zstd.compress() started and stopped the workers for every single chunk. Cache the compressor per thread instead (keyed by level and worker count), so its pool is built once and stays alive between chunks. Measured on the 2MiB chunks of the Silesia corpus with 4 workers on an M3 Pro: at zstd,-4 +20% throughput (2345 -> 2815 MB/s), at zstd,3 +9%. The compressed data does not change: set_pledged_input_size() plus FLUSH_FRAME produce exactly the bytes the one-shot API produced. Single-threaded compression goes through the same cache. That gains next to nothing (0.98 .. 1.07x), but keeps compression to a single code path. The price of caching is that a thread holds on to the contexts it used: +1.6MB peak RSS for a single-threaded level 3 context, +8MB for a 4 worker one. Compressors are stateful, so they are per thread, like the scratch buffer. In the child of a fork they are dropped, but deliberately not freed: the worker threads do not exist there, and freeing such a context would make libzstd wait for them forever (verified: without this, a child hangs as soon as it compresses). Co-Authored-By: Claude Opus 5 --- src/borg/compress.pyi | 2 + src/borg/compress.pyx | 96 ++++++++++++++++++++++++----- src/borg/testsuite/compress_test.py | 91 ++++++++++++++++++++++++++- 3 files changed, 174 insertions(+), 15 deletions(-) diff --git a/src/borg/compress.pyi b/src/borg/compress.pyi index e3575840ea..723e3cee13 100644 --- a/src/borg/compress.pyi +++ b/src/borg/compress.pyi @@ -5,6 +5,8 @@ ZSTD_MT_MIN_SIZE: int def get_compressor(name: str, **kwargs) -> Any: ... def get_zstd_mt_workers(stream: bool = ...) -> int: ... +def get_zstd_compressor(level: int, workers: int) -> Any: ... +def forget_zstd_compressors(keep_alive: bool = ...) -> None: ... class Compressor: def __init__(self, name: Any = ..., **kwargs) -> None: ... diff --git a/src/borg/compress.pyx b/src/borg/compress.pyx index f26b2d6ebd..dc2ea31760 100644 --- a/src/borg/compress.pyx +++ b/src/borg/compress.pyx @@ -130,6 +130,72 @@ def get_buffer(): buffer = _thread_local.buffer = Buffer(bytearray, size=0) return buffer + +_zstd_compressors_kept_after_fork = [] + + +def get_zstd_compressor(level, workers): + """Return this thread's zstd compressor for *level* / *workers*, creating it on first use. + + Setting up a compression context is not free, and for a multithreaded one (workers > 1) it + is expensive: libzstd starts a worker pool then and stops it again when the context is + freed, so the one-shot zstd.compress() would start and stop *workers* threads for every + single chunk. Keeping the compressor - and with it its pool - alive between chunks avoids + that. Measured on the 2MiB chunks of the Silesia corpus with 4 workers on an M3 Pro: at + zstd,-4 +20% throughput (2345 -> 2815 MB/s), at zstd,3 +9% - the slower the level, the + less the pool setup weighs against the compression itself. Single-threaded compression goes through the same cache, + which gains next to nothing (0.98 .. 1.07x, measured for 16KiB .. 2MiB chunks at zstd,3 and + zstd,-4), but keeps this to one code path. + + The price is that each thread holds on to the contexts it used, rather than freeing them + after every chunk: measured +1.6MB peak RSS for a single-threaded level 3 context, +8MB for + a 4 worker one (which also keeps its job buffers). + + Compressors are stateful, so - like the scratch buffer, see get_buffer() - they must not be + shared between threads. The cache is keyed by level and worker count because one thread may + use more than one of them (e.g. a "borg recreate" recompressing to a different level). + """ + workers = max(workers, 1) # BORG_ZSTD_MT_WORKERS=0 and =1 both mean single-threaded + try: + compressors = _thread_local.zstd_compressors + except AttributeError: + compressors = _thread_local.zstd_compressors = {} + compressor = compressors.get((level, workers)) + if compressor is None: + params = zstd.CompressionParameter + options = {params.compression_level: level} + if workers > 1: + # the smallest job libzstd accepts gives the most parallelism, see ZSTD._decide. + # note: nb_workers == 1 would mean one *asynchronous* worker, not single-threaded. + options[params.nb_workers] = workers + options[params.job_size] = ZSTD_JOB_SIZE_MIN + compressor = zstd.ZstdCompressor(options=options) + compressors[(level, workers)] = compressor + return compressor + + +def forget_zstd_compressors(keep_alive=False): + """Drop this thread's cached zstd compressors, so that the next call builds fresh ones. + + With *keep_alive*, the compressors are parked in a module-level list instead of being freed. + That is for the child of a fork: only the forking thread exists there, so the worker threads + of a cached multithreaded compressor are gone, while its libzstd context still believes it + has them - and freeing it would make libzstd wait for threads that do not exist. Leaking a + few contexts in a child (which usually execs or is short-lived anyway) is the cheaper + problem. + """ + compressors = getattr(_thread_local, "zstd_compressors", None) + if not compressors: + return + if keep_alive: + _zstd_compressors_kept_after_fork.append(compressors) + _thread_local.zstd_compressors = {} + + +if hasattr(os, "register_at_fork"): # posix only + os.register_at_fork(after_in_child=lambda: forget_zstd_compressors(keep_alive=True)) + + cdef class CompressorBase: """ base class for all (de)compression classes, @@ -442,24 +508,26 @@ class ZSTD(DecidingCompressor): def _decide(self, meta, data): # libzstd can compress a single chunk with a thread pool, but only if it is told to # split it into jobs: the default job size is derived from the window log and swallows - # any chunk borg produces whole, so the workers would never engage. We ask for the - # smallest job libzstd accepts, which gives the most parallelism (and, as a trade-off, - # the largest ratio loss). See ZSTD_MT_MIN_SIZE for why small chunks are excluded. - workers = get_zstd_mt_workers() - if workers > 1 and len(data) >= ZSTD_MT_MIN_SIZE: - params = zstd.CompressionParameter - zstd_data = zstd.compress(data, options={ - params.compression_level: self.level, - params.nb_workers: workers, - params.job_size: ZSTD_JOB_SIZE_MIN, - }) - else: - zstd_data = zstd.compress(data, self.level) + # any chunk borg produces whole, so the workers would never engage. get_zstd_compressor + # asks for the smallest job libzstd accepts, which gives the most parallelism (and, as a + # trade-off, the largest ratio loss). See ZSTD_MT_MIN_SIZE for why small chunks are + # compressed single-threaded. The compressor is this thread's, reused between chunks. + workers = get_zstd_mt_workers() if len(data) >= ZSTD_MT_MIN_SIZE else 1 + compressor = get_zstd_compressor(self.level, workers) + try: + # we know the size, so the frame header gets it just like the one-shot API's does + compressor.set_pledged_input_size(len(data)) + zstd_data = compressor.compress(data, mode=zstd.ZstdCompressor.FLUSH_FRAME) + except Exception: + # a failed call can leave the compressor inside a frame, and whatever it produced + # next would continue that frame - so this thread starts over with fresh ones. + forget_zstd_compressors() + raise if len(zstd_data) < len(data): return self, (meta, zstd_data) else: return NONE_COMPRESSOR, (meta, None) - + def decompress(self, meta, data): meta, data = super().decompress(meta, data) try: diff --git a/src/borg/testsuite/compress_test.py b/src/borg/testsuite/compress_test.py index 3a86af5708..e42f58e784 100644 --- a/src/borg/testsuite/compress_test.py +++ b/src/borg/testsuite/compress_test.py @@ -6,7 +6,8 @@ import pytest from ..compress import get_compressor, Compressor, CNONE, ZLIB, LZ4, LZMA, ZSTD, Auto -from ..compress import get_zstd_mt_workers +from ..compress import get_zstd_mt_workers, get_zstd_compressor, forget_zstd_compressors +from ..compress import ZSTD_JOB_SIZE_MIN, ZSTD_MT_MIN_SIZE from ..helpers import Error from ..helpers import CompressionSpec from ..constants import ROBJ_FILE_STREAM, ROBJ_ARCHIVE_META @@ -359,3 +360,91 @@ def workers_for(env_value, stream=False): for invalid in ["yes", "4x", "1.5", "-1"]: with pytest.raises(Error): workers_for(invalid) + + +@pytest.fixture +def zstd_mt(monkeypatch): + """Force the multithreaded zstd path on, whatever the machine's cpu count is.""" + monkeypatch.setattr("borg.compress._zstd_mt_workers", (4, 4)) + forget_zstd_compressors() # do not inherit compressors cached by another test + yield 4 + forget_zstd_compressors() + + +def mt_data(size=ZSTD_MT_MIN_SIZE): + """Compressible data, big enough for the MT path to engage.""" + rnd = random.Random(0) + words = [b"borg", b"backup", b"chunk", b"compress", b"zstd"] + out = bytearray() + while len(out) < size: + out += rnd.choice(words) + return bytes(out[:size]) + + +def test_zstd_compressor_is_cached_per_thread_and_level(zstd_mt): + workers = zstd_mt + compressor = get_zstd_compressor(3, workers) + assert get_zstd_compressor(3, workers) is compressor # reused, that is the point + assert get_zstd_compressor(-4, workers) is not compressor # but not across levels + assert get_zstd_compressor(3, workers + 1) is not compressor # nor across worker counts + assert get_zstd_compressor(3, 0) is get_zstd_compressor(3, 1) # 0 workers == 1 == no pool + + forget_zstd_compressors() + assert get_zstd_compressor(3, workers) is not compressor # dropped, so rebuilt + + in_thread = ThreadPoolExecutor(max_workers=1).submit(get_zstd_compressor, 3, workers).result() + assert in_thread is not get_zstd_compressor(3, workers) # each thread has its own + + +def test_zstd_mt_reused_compressor_output(zstd_mt): + """Reusing the compressor must give the same bytes as a fresh one-shot compression.""" + from borg.compress import zstd # the stdlib module or its backport, whichever is in use + + data = mt_data() + params = zstd.CompressionParameter + expected = zstd.compress( + data, options={params.compression_level: 3, params.nb_workers: zstd_mt, params.job_size: ZSTD_JOB_SIZE_MIN} + ) + compressor = get_compressor(name="zstd", level=3) + meta, first = compressor.compress({}, data) + meta, second = compressor.compress({}, data) # same compressor object, second frame + assert first == expected + assert second == expected # no state carried over from the previous chunk + + +def test_zstd_mt_roundtrip_concurrent(zstd_mt): + """Several threads compressing at once must not share a compressor (or corrupt each other).""" + chunks = [mt_data() + bytes([i]) * 1024 for i in range(8)] + compressor = get_compressor(name="zstd", level=3) # one instance, shared by all threads + + def roundtrip(data): + meta, cdata = compressor.compress({}, data) + meta, plain = compressor.decompress(dict(meta), cdata) + return plain + + with ThreadPoolExecutor(max_workers=4) as pool: + assert list(pool.map(roundtrip, chunks)) == chunks + + +def test_zstd_below_threshold_is_single_threaded(zstd_mt): + """Small chunks are compressed single-threaded, by a compressor cached for one worker.""" + from borg.compress import zstd, _thread_local + + data = mt_data(ZSTD_MT_MIN_SIZE - 1) + compressor = get_compressor(name="zstd", level=3) + meta, cdata = compressor.compress({}, data) + assert cdata == zstd.compress(data, 3) # same bytes as the one-shot single-threaded API + cached = _thread_local.zstd_compressors + assert list(cached) == [(3, 1)] # one worker: no thread pool was set up for this chunk + assert cached[(3, 1)] is get_zstd_compressor(3, 1) + + +def test_zstd_single_threaded_reused_compressor_output(zstd_mt): + """Reuse must not change the single-threaded output either, chunk after chunk.""" + from borg.compress import zstd + + data = mt_data(ZSTD_MT_MIN_SIZE - 1) + expected = zstd.compress(data, 3) + compressor = get_compressor(name="zstd", level=3) + assert compressor.compress({}, data)[1] == expected + assert compressor.compress({}, data)[1] == expected