perf: size the encoder pool to the CPUs the container may use #3

Merged
lyrathorpe merged 1 commits from perf/encode-concurrency into main 2026-08-24 11:35:04 +01:00
2 changed files with 49 additions and 4 deletions
+14
View File
@@ -78,6 +78,20 @@ music-mirror --source /music --mirror /music-mp3 --subdir "Artist/Album"
`--subdir` never prunes: a partial pass cannot tell an orphan from a file `--subdir` never prunes: a partial pass cannot tell an orphan from a file
outside its own scope. outside its own scope.
### Concurrency
LAME is single-threaded — ffmpeg reports `Threading capabilities: none` for
`libmp3lame` — so throughput comes entirely from running several encoders at
once, one process per file. `--jobs` defaults to the CPUs the process may
actually use, which inside a container means the `cpus:` allowance rather than
the host's core count. Each pass logs the number it settled on.
As a rough guide, a Zen 3 core encodes about 4060× realtime at V0 depending on
clock, so six cores clear roughly 250 hours of audio per hour of wall clock.
The first full pass is the expensive one; after that only new and changed files
are touched. Lower `MUSIC_MIRROR_JOBS` if you would rather the NAS stayed
responsive than finished sooner.
Requires `ffmpeg` and `ffprobe` on `PATH`. The container image provides both. Requires `ffmpeg` and `ffprobe` on `PATH`. The container image provides both.
## Running it on TrueNAS Scale ## Running it on TrueNAS Scale
+35 -4
View File
@@ -16,6 +16,7 @@ makes runs idempotent without a database to keep in step.
import argparse import argparse
import concurrent.futures import concurrent.futures
import fcntl import fcntl
import functools
import logging import logging
import os import os
import re import re
@@ -123,8 +124,13 @@ def is_current(source, mirror):
return abs(source.stat().st_mtime - mirror.stat().st_mtime) <= MTIME_TOLERANCE_SECONDS return abs(source.stat().st_mtime - mirror.stat().st_mtime) <= MTIME_TOLERANCE_SECONDS
@functools.lru_cache(maxsize=4096)
def find_cover(directory): def find_cover(directory):
"""Return an external cover image for a directory, if one is present.""" """Return an external cover image for a directory, if one is present.
Cached because an album's tracks all ask the same question, and the answer
costs one stat per candidate name.
"""
for name in COVER_NAMES: for name in COVER_NAMES:
candidate = directory / name candidate = directory / name
if candidate.is_file(): if candidate.is_file():
@@ -201,7 +207,12 @@ def encode(source, mirror, quality_args, dry_run):
return Result("encoded", mirror) return Result("encoded", mirror)
mirror.parent.mkdir(parents=True, exist_ok=True) mirror.parent.mkdir(parents=True, exist_ok=True)
cover = None if has_embedded_picture(source) else find_cover(source.parent) # Probing costs an ffprobe process per file, so only ask when the answer
# can change the command. With no cover file beside the track, `-map
# 0:v:0?` carries embedded art if there is any and shrugs if there is not.
cover = find_cover(source.parent)
if cover is not None and has_embedded_picture(source):
cover = None
# Read the source's mtime before encoding, not after. If the file is still # Read the source's mtime before encoding, not after. If the file is still
# being written -- a Lidarr import landing mid-pass -- stamping the mirror # being written -- a Lidarr import landing mid-pass -- stamping the mirror
@@ -331,6 +342,7 @@ def run_once(scan_root, source_root, mirror_root, quality_args, jobs, dry_run, d
computed against; they differ only for a partial pass over one directory. computed against; they differ only for a partial pass over one directory.
""" """
started = time.monotonic() started = time.monotonic()
logger.info("pass starting with %d concurrent encoders", jobs)
counts = {"encoded": 0, "copied": 0, "skipped": 0, "failed": 0} counts = {"encoded": 0, "copied": 0, "skipped": 0, "failed": 0}
failures = [] failures = []
@@ -376,6 +388,25 @@ def acquire_lock(mirror_root):
return handle return handle
def default_jobs():
"""Return the number of CPUs this process may actually use.
os.cpu_count() reports the host's total, which in a container with a `cpus:`
limit means starting several times more encoders than there is CPU to run
them. libmp3lame is single-threaded, so one process per available CPU is the
whole of the concurrency story.
"""
try:
quota, period = Path("/sys/fs/cgroup/cpu.max").read_text().split()
if quota != "max":
return max(1, round(int(quota) / int(period)))
except (OSError, ValueError):
pass
if hasattr(os, "process_cpu_count"): # 3.13+, respects CPU affinity
return os.process_cpu_count() or 4
return os.cpu_count() or 4
def build_parser(): def build_parser():
"""Return the argument parser. Every option also reads an env var, so the """Return the argument parser. Every option also reads an env var, so the
container can be configured without a command line.""" container can be configured without a command line."""
@@ -401,8 +432,8 @@ def build_parser():
parser.add_argument( parser.add_argument(
"--jobs", "--jobs",
type=int, type=int,
default=int(os.getenv("MUSIC_MIRROR_JOBS", "0")) or (os.cpu_count() or 4), default=int(os.getenv("MUSIC_MIRROR_JOBS", "0")) or default_jobs(),
help="concurrent encodes (env MUSIC_MIRROR_JOBS; default: CPU count)", help="concurrent encodes (env MUSIC_MIRROR_JOBS; default: available CPUs)",
) )
parser.add_argument( parser.add_argument(
"--interval", "--interval",