Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -34,7 +34,7 @@ metadata = damacy.ZarrMetadata(
planner = damacy.ChunkPlanner(metadata=metadata, limits=damacy.PlanLimits())
executor = damacy.CpuExecutor(
reader=chunk_reader,
limits=damacy.CpuLimits(max_memory_bytes=1 << 30, decode_workers=8),
limits=damacy.CpuLimits(max_memory_bytes=3 << 30, decode_workers=8),
)

with damacy.Pipeline(
Expand Down
12 changes: 11 additions & 1 deletion bench/main.c
Original file line number Diff line number Diff line change
Expand Up @@ -120,6 +120,7 @@ struct scenario
uint64_t max_gpu_memory_bytes;
uint64_t max_cpu_memory_bytes;
uint32_t decode_workers;
uint32_t chunks_per_input_buffer;
int cpu;
uint32_t max_chunk_uncompressed_bytes; // 0 → tuning_defaults() baseline
uint64_t max_read_op_bytes; // 0 → tuning_defaults() baseline
Expand Down Expand Up @@ -451,6 +452,10 @@ parse_scenario(struct cslice src, struct scenario* sc)
static const struct json_query p_workers[] = {
{ QUERY_KEY, .key = "pipeline" }, { QUERY_KEY, .key = "decode_workers" }
};
static const struct json_query p_buffer_chunks[] = {
{ QUERY_KEY, .key = "pipeline" },
{ QUERY_KEY, .key = "chunks_per_input_buffer" }
};
read_uint_opt(src, p_cpu, countof(p_cpu), &v, 0);
if (v > (UINT64_MAX >> 20) || (sc->cpu && !v))
return 1;
Expand All @@ -459,6 +464,10 @@ parse_scenario(struct cslice src, struct scenario* sc)
if (!v || v > UINT32_MAX)
return 1;
sc->decode_workers = (uint32_t)v;
read_uint_opt(src, p_buffer_chunks, countof(p_buffer_chunks), &v, 256);
if (!v || v > UINT32_MAX)
return 1;
sc->chunks_per_input_buffer = (uint32_t)v;
read_uint_opt(src, p_g, countof(p_g), &v, 0);
sc->max_gpu_memory_bytes = v << 20;
read_uint_opt(src, p_c, countof(p_c), &v, 0);
Expand Down Expand Up @@ -1134,7 +1143,8 @@ pipeline_create(const struct scenario* scenario,
.decode_workers = scenario->decode_workers,
.max_encoded_chunk_bytes = (uint32_t)cfg->tuning.max_read_op_bytes,
.max_decoded_chunk_bytes = cfg->tuning.max_chunk_uncompressed_bytes,
.max_memory_bytes = scenario->max_cpu_memory_bytes },
.max_memory_bytes = scenario->max_cpu_memory_bytes,
.chunks_per_input_buffer = scenario->chunks_per_input_buffer },
&pipeline->executor);
if (status != DAMACY_OK)
return status;
Expand Down
5 changes: 3 additions & 2 deletions bench/scenarios/throughput-cpu.json
Original file line number Diff line number Diff line change
Expand Up @@ -49,7 +49,8 @@
"n_shard_index_cache": 32768,
"max_shards_per_sample": 8,
"executor": "cpu",
"max_cpu_memory_mb": 4096,
"decode_workers": 8
"max_cpu_memory_mb": 5120,
"decode_workers": 8,
"chunks_per_input_buffer": 256
}
}
4 changes: 3 additions & 1 deletion dev/cpu-pipeline-validation.md
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,9 @@ eight I/O workers, and a 4 GiB executor memory limit. Each worker count was
measured once in the same 16-core CPU allocation, in the order 1, 4, 8, 16.
Filesystem caches were shared between runs. These are rates after warmup, not
cold-storage measurements or training-loop timings. The scenario
is [throughput-cpu.json](../bench/scenarios/throughput-cpu.json).
is [throughput-cpu.json](../bench/scenarios/throughput-cpu.json). It now sets
a 5 GiB limit, because later 256-chunk input buffers reserve 2 GiB with the
default 4 MiB encoded-chunk bound.

CUDA comparisons use five warmup batches and thirty measured batches, with a
6 GiB device-memory limit. Baseline and refactor runs alternate on one L40
Expand Down
11 changes: 7 additions & 4 deletions dev/cpu-pipeline.md
Original file line number Diff line number Diff line change
Expand Up @@ -88,10 +88,13 @@ reuse can be optimized separately.

CPU I/O workers and decode workers are configured independently. There are two
input groups and two output buffers. Each input group holds at most
`decode_workers` chunks and `decode_workers * max_encoded_chunk_bytes` encoded
bytes. Merged reads fit these limits, and their count respects the reader
capacity. Decoding uses a bounded per-worker output buffer and zstd context;
memory admission also reserves conservative Blosc workspace. The memory cap
`chunks_per_input_buffer` chunks and
`chunks_per_input_buffer * max_encoded_chunk_bytes` encoded bytes. The setting
is at least `decode_workers`, so every worker can get a chunk. Merged reads
contain at most `min(decode_workers, chunks_per_input_buffer)` chunks, and their
count respects the reader capacity. Decoding uses a bounded per-worker output
buffer and zstd context; memory admission also reserves conservative Blosc
workspace. The memory cap
covers executor buffers, active read plans, and temporary read-planning
scratch. It excludes metadata, reader queues, prepared-plan storage, thread
stacks, allocator overhead, and total process RSS.
Expand Down
20 changes: 17 additions & 3 deletions docs/pipeline.md
Original file line number Diff line number Diff line change
Expand Up @@ -29,10 +29,11 @@ planner = damacy.ChunkPlanner(
executor = damacy.CpuExecutor(
reader=chunk_reader,
limits=damacy.CpuLimits(
max_memory_bytes=1 << 30,
max_memory_bytes=3 << 30,
decode_workers=8,
max_encoded_chunk_bytes=4 << 20,
max_decoded_chunk_bytes=2 << 20,
chunks_per_input_buffer=256,
),
)
output = damacy.BatchSpec(samples=2, shape=(64, 256, 256), dtype="f32")
Expand Down Expand Up @@ -112,6 +113,7 @@ capacities. Defaults come from the Python value objects.
| `PlanLimits.max_plan_bytes` | Owned storage per prepared plan. |
| `CpuLimits.max_memory_bytes` | CPU executor's buffers, codec workspace, active read plans, and temporary read-planning scratch. |
| `CpuLimits.decode_workers` | Total decoding/assembly workers, including the calling scheduler thread. |
| `CpuLimits.chunks_per_input_buffer` | Chunks each of the two encoded-input buffers holds; from `decode_workers` to 16384. |
| `FileReader.workers` | Bulk I/O workers, separate from decoding workers. |
| `FileReader.max_inflight_reads` | Bulk read capacity; execution respects this bound and retries saturation. |
| `CudaLimits` | GPU memory and execution geometry, plus the CUDA codec-layout cache capacity. |
Expand All @@ -126,8 +128,20 @@ older, stricter cache validation.
The CPU executor merges adjacent or overlapping encoded ranges within each
shard, then interleaves the reads across shards. It decodes each unique source
chunk once per batch, including when several output samples use that chunk.
Both input groups stay within the decode-worker, encoded-byte, and reader
limits; merging does not increase the number of decode workers.
Each of its two encoded-input buffers holds `chunks_per_input_buffer` chunks
(default 256). A merged read contains at most
`min(decode_workers, chunks_per_input_buffer)` chunks, so one buffer can hold
several reads at once; with one decode worker, reads do not merge. Submission
also respects the reader limit. Decoder workspaces remain per worker.

The two input buffers reserve
`2 * chunks_per_input_buffer * max_encoded_chunk_bytes` bytes, plus per-chunk
bookkeeping. With the defaults, 256 chunks and a 4 MiB encoded-chunk bound,
that is 2 GiB before decoder workspaces and output buffers. Set the bound to
match the largest encoded chunk expected, and include this reserve in
`max_memory_bytes`. If the budget is too small, starting the pipeline reports
`BUDGET` and logs how many bytes the input buffers, decoder workspaces, and
output buffers need.

CPU memory admission includes active read plans, temporary planning scratch,
and a conservative allowance for Blosc scratch storage.
Expand Down
8 changes: 8 additions & 0 deletions python/damacy/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -826,12 +826,19 @@ class CpuLimits:
decode_workers: int = 8
max_encoded_chunk_bytes: int = 4 << 20
max_decoded_chunk_bytes: int = 2 << 20
chunks_per_input_buffer: int = 256

def __post_init__(self) -> None:
_positive_int(self.max_memory_bytes, "max_memory_bytes", (1 << 64) - 1)
_positive_int(self.decode_workers, "decode_workers", _native.MAX_IO_THREADS)
_positive_int(self.max_encoded_chunk_bytes, "max_encoded_chunk_bytes")
_positive_int(self.max_decoded_chunk_bytes, "max_decoded_chunk_bytes")
_positive_int(self.chunks_per_input_buffer, "chunks_per_input_buffer", 16384)
if self.chunks_per_input_buffer < self.decode_workers:
raise ValueError(
f"chunks_per_input_buffer ({self.chunks_per_input_buffer}) must be "
f"at least decode_workers ({self.decode_workers})"
)


@dataclass(frozen=True, slots=True)
Expand Down Expand Up @@ -962,6 +969,7 @@ def __init__(self, *, reader: FileReader, limits: CpuLimits) -> None:
limits.max_encoded_chunk_bytes,
limits.max_decoded_chunk_bytes,
limits.max_memory_bytes,
limits.chunks_per_input_buffer,
)


Expand Down
5 changes: 3 additions & 2 deletions python/damacy/_components.c
Original file line number Diff line number Diff line change
Expand Up @@ -186,12 +186,13 @@ create_cpu_executor(PyObject* self, PyObject* args)
unsigned long long bytes;
struct damacy_cpu_config config;
if (!PyArg_ParseTuple(args,
"OIIIK",
"OIIIKI",
&reader_object,
&config.decode_workers,
&config.max_encoded_chunk_bytes,
&config.max_decoded_chunk_bytes,
&bytes))
&bytes,
&config.chunks_per_input_buffer))
return NULL;
config.max_memory_bytes = bytes;
struct damacy_reader* reader = component_value(reader_object, READER);
Expand Down
8 changes: 7 additions & 1 deletion python/damacy/_native.pyi
Original file line number Diff line number Diff line change
Expand Up @@ -214,7 +214,13 @@ def create_planner(
/,
) -> object: ...
def create_cpu_executor(
reader: object, workers: int, max_encoded: int, max_decoded: int, max_memory: int, /
reader: object,
workers: int,
max_encoded: int,
max_decoded: int,
max_memory: int,
chunks_per_input_buffer: int,
/,
) -> object: ...
def create_cuda_executor(
reader: object,
Expand Down
3 changes: 3 additions & 0 deletions python/tests/test_components.py
Original file line number Diff line number Diff line change
Expand Up @@ -284,6 +284,9 @@ def test_invalid_limits():
lambda: damacy.CpuLimits(0),
lambda: damacy.PlanLimits(max_chunks=0),
lambda: damacy.CpuLimits(1 << 20, decode_workers=-1),
lambda: damacy.CpuLimits(1 << 20, chunks_per_input_buffer=0),
lambda: damacy.CpuLimits(1 << 20, decode_workers=8, chunks_per_input_buffer=7),
lambda: damacy.CpuLimits(1 << 20, chunks_per_input_buffer=16385),
]:
with pytest.raises(ValueError):
make()
Expand Down
2 changes: 1 addition & 1 deletion src/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -278,7 +278,7 @@ if(NOT DAMACY_FUZZ)
LINKS pipeline_components plan_builder prefetcher array_meta shard_index
metadata_store_async lookahead
)
add_src_lib(cpu_executor SOURCES executor/cpu_executor.c
add_src_lib(cpu_executor SOURCES executor/cpu_executor.c executor/cpu_read_plan.h
LINKS pipeline_components prepared_plan dispatch_utils threadpool damacy_stats
PkgConfig::ZSTD PkgConfig::BLOSC
)
Expand Down
1 change: 1 addition & 0 deletions src/damacy_pipeline.h
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,7 @@ extern "C"
uint32_t max_encoded_chunk_bytes;
uint32_t max_decoded_chunk_bytes;
uint64_t max_memory_bytes;
uint32_t chunks_per_input_buffer;
};

struct damacy_cuda_config
Expand Down
Loading
Loading