From a50eb6f7cd385dfcc0680bc12c507eec6f2593d5 Mon Sep 17 00:00:00 2001 From: Thomas Waldmann Date: Mon, 27 Jul 2026 22:49:25 +0200 Subject: [PATCH] use multithreaded zstd compression for big chunks, #9961 libzstd can compress a single chunk with a thread pool, but borg never asked it to: zstd.compress() was called without options, so nb_workers stayed 0. Just setting nb_workers would not have helped either. zstd splits the input into jobs sized from the window log, and that default job size swallows any chunk borg produces whole, so the workers never engage - measured at 0.97x, i.e. slightly slower than single-threaded. The job size has to be set explicitly. libzstd refuses to use a job smaller than 512KiB, so we ask for exactly that: it gives the most jobs and thus the most parallelism. Small chunks are excluded, but not because they cannot be split - a 768KiB chunk is happily cut into 512KiB + 256KiB. The problem is that such a split is very uneven: every thread waits for the one full-size job, which takes about as long as compressing the whole chunk single-threaded, and the thread overhead comes on top. Measured on a 12-core machine over levels 1..10, that is a real loss: 0.80x..0.91x at 544KiB, 0.90x..1.07x at 608KiB, break-even around 640KiB, then 1.21x..1.38x at 768KiB. The cutoff is therefore 768KiB (= 1.5 jobs), which keeps some margin over break-even because the thread overhead is machine dependent. Measured on the same machine, zstd,3, 10GiB of compressible data in one file: borg create 55.0s -> 33.8s (1.63x). Compressor alone: 2.61x for a 2MiB chunk, 4.16x for an 8MiB one. This trades a little compression ratio for speed, because a job only sees its own data: +0.05% archive size for 1MiB chunks at zstd,3, +0.64% for 8MiB ones, and somewhat more at higher levels which rely on longer match history. It also uses more cpu time in total (+22% for the 10GiB create) to reduce wallclock time. BORG_ZSTD_MT_WORKERS=1 turns it off for those who prefer the smaller archive or have to share the cpu with other work. Note that the compressed bytes now depend on the worker count, so compressing the same chunk twice on different machines can give different (equally valid) output. Chunk ids are unaffected: they are computed on the plaintext before compression. Co-Authored-By: Claude Opus 5 --- docs/usage/general/environment.rst.inc | 11 +++++ src/borg/compress.pyi | 4 ++ src/borg/compress.pyx | 64 +++++++++++++++++++++++++- src/borg/testsuite/compress_test.py | 19 ++++++++ 4 files changed, 96 insertions(+), 2 deletions(-) diff --git a/docs/usage/general/environment.rst.inc b/docs/usage/general/environment.rst.inc index 082395dcb9..f11d3ea46f 100644 --- a/docs/usage/general/environment.rst.inc +++ b/docs/usage/general/environment.rst.inc @@ -101,6 +101,17 @@ General: set it, and can optionally show a chart of the measurements in your browser (``--html --open``). 0 means "always multi-threaded", a very large value effectively disables multi-threading. + BORG_ZSTD_MT_WORKERS + When set to a numeric value, use that many threads to zstd-compress a single chunk + (default: the cpu count). 0 or 1 means single-threaded compression. + Only relevant when compressing with ``zstd``. + Chunks below 768KiB are always compressed single-threaded: libzstd will not use a + compression job smaller than 512KiB, so a small chunk gets split very unevenly and + multi-threading it would be slower than not doing it at all. + Multi-threading trades a little compression ratio for speed (measured at ``zstd,3``: + +0.05% archive size for 1MiB chunks, +0.64% for 8MiB ones, more at higher levels), and + it uses more cpu time in total to reduce the wallclock time. Set it to 1 if you would + rather have the smaller archive, or if borg has to share the cpu with other work. BORG_SHOW_SYSINFO When set to no (default: yes), system information (like OS, Python version, ...) in exceptions is not shown. diff --git a/src/borg/compress.pyi b/src/borg/compress.pyi index 9fbdfed8c5..d11b7bff0f 100644 --- a/src/borg/compress.pyi +++ b/src/borg/compress.pyi @@ -1,6 +1,10 @@ from typing import Any +ZSTD_JOB_SIZE_MIN: int +ZSTD_MT_MIN_SIZE: int + def get_compressor(name: str, **kwargs) -> Any: ... +def get_zstd_mt_workers() -> int: ... class Compressor: def __init__(self, name: Any = ..., **kwargs) -> None: ... diff --git a/src/borg/compress.pyx b/src/borg/compress.pyx index faf3031273..739c558cbf 100644 --- a/src/borg/compress.pyx +++ b/src/borg/compress.pyx @@ -16,6 +16,7 @@ decompressor. """ import math +import os import random from struct import Struct import sys @@ -27,7 +28,7 @@ except ImportError: lzma = None from .constants import MAX_DATA_SIZE, ROBJ_FILE_STREAM -from .helpers import Buffer, DecompressionError +from .helpers import Buffer, DecompressionError, Error from .helpers.argparsing import ArgumentTypeError if sys.version_info >= (3, 14): @@ -35,6 +36,51 @@ if sys.version_info >= (3, 14): else: from backports import zstd +# libzstd will not use a compression job smaller than this, no matter what job_size says. +ZSTD_JOB_SIZE_MIN = 512 * 1024 + +# Below this chunk size we compress single-threaded. +# +# zstd cuts a chunk into ceil(size / job_size) jobs, the last one being short, so a chunk only +# a little bigger than one job is split very unevenly (e.g. 576KiB -> 512KiB + 64KiB). All +# threads then wait for that one full-size job, which takes about as long as compressing the +# whole chunk single-threaded would - but with the thread overhead added on top. Measured on a +# 12-core machine, that is a real loss: 0.80x .. 0.91x at 544KiB (levels 1 .. 10), still +# 0.90x .. 1.07x at 608KiB, break-even around 640KiB, then 1.21x .. 1.38x at 768KiB. +# 768KiB (= 1.5 jobs) keeps some margin over break-even, as the overhead is machine dependent. +ZSTD_MT_MIN_SIZE = 768 * 1024 + +_zstd_mt_workers = None + + +def get_zstd_mt_workers(): + """How many threads libzstd may use to compress a single chunk. + + Defaults to the cpu count, BORG_ZSTD_MT_WORKERS overrides it. 0 or 1 means + single-threaded compression, which also avoids the small loss of compression ratio + that splitting a chunk into jobs causes (measured at zstd,3: +0.05% for a 1MiB chunk, + +0.64% for an 8MiB one; higher levels lose a bit more as they rely on longer match + history). + + This is called for every chunk, so the result is cached. The env var is evaluated on + first use rather than at import time, so that an invalid value is reported via borg's + normal error handling rather than as a traceback while importing. + """ + global _zstd_mt_workers + if _zstd_mt_workers is None: + value = os.environ.get("BORG_ZSTD_MT_WORKERS") + if value is None: + workers = os.cpu_count() or 1 + else: + try: + workers = int(value) + except ValueError: + raise Error(f"BORG_ZSTD_MT_WORKERS must be an integer, but is: {value!r}") from None + if workers < 0: + raise Error(f"BORG_ZSTD_MT_WORKERS must not be negative, but is: {workers}") + _zstd_mt_workers = workers + return _zstd_mt_workers + cdef extern from "lz4.h": int LZ4_compress_default(const char* source, char* dest, int inputSize, int maxOutputSize) nogil int LZ4_decompress_safe(const char* source, char* dest, int inputSize, int maxOutputSize) nogil @@ -303,7 +349,21 @@ class ZSTD(DecidingCompressor): self.level = level def _decide(self, meta, data): - zstd_data = zstd.compress(data, self.level) + # 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) if len(zstd_data) < len(data): return self, (meta, zstd_data) else: diff --git a/src/borg/testsuite/compress_test.py b/src/borg/testsuite/compress_test.py index 62ef59f51b..6348246f09 100644 --- a/src/borg/testsuite/compress_test.py +++ b/src/borg/testsuite/compress_test.py @@ -4,6 +4,8 @@ import pytest from ..compress import get_compressor, Compressor, CNONE, ZLIB, LZ4, LZMA, ZSTD, Auto +from ..compress import get_zstd_mt_workers +from ..helpers import Error from ..helpers import CompressionSpec from ..constants import ROBJ_FILE_STREAM, ROBJ_ARCHIVE_META from ..helpers.argparsing import ArgumentTypeError @@ -261,3 +263,20 @@ def test_robj_specific_obfuscation(data_length, expected_padding, robj_type): assert ( len(compressed) == expected_padded_size ), f"For {data_length}, expected {expected_padded_size}, got {len(compressed)}" + + +def test_zstd_mt_workers_from_env(monkeypatch): + def workers_for(env_value): + monkeypatch.setattr("borg.compress._zstd_mt_workers", None) # drop the cache + if env_value is None: + monkeypatch.delenv("BORG_ZSTD_MT_WORKERS", raising=False) + else: + monkeypatch.setenv("BORG_ZSTD_MT_WORKERS", env_value) + return get_zstd_mt_workers() + + assert workers_for(None) == (os.cpu_count() or 1) + for value, expected in [("0", 0), ("1", 1), ("4", 4)]: + assert workers_for(value) == expected + for invalid in ["", "yes", "4x", "1.5", "-1"]: + with pytest.raises(Error): + workers_for(invalid)