Parallel Encoding and Decoding in Nimble: Faster Flushes and Reads
Introduction
Nimble can reduce wall-clock time when a stripe or batch contains enough independent work. It parallelizes encoding across complete streams and decoding across eligible sibling fields. The writer bounds the number of encode tasks using the total stream count, maximum parallelism, and the minimum-streams-per-task setting. Each encode task may process several streams during a flush, encoding one stream at a time. A decode task handles a contiguous group of eligible sibling fields under a row or FlatMap-as-struct reader. The nested schema defines parent-child dependencies: Nimble reads a row's null information before its child values, and array lengths before element values because they provide the element offsets.
How the writer and reader split work
The writer chooses the encode-task count and sorts buffered streams before dispatch. The reader runs decode tasks at eligible row and FlatMap-as-struct readers in the nested schema. The decode-task budget is the configured decode parallelism, shared across the schema. With a limit of four, eligible readers share up to four-way parallelism across the schema. The planner starts at the outermost eligible readers, shares capacity across readers at that level, then moves inward if capacity remains.
| Path | What one task handles | Scheduling | Coordination |
|---|---|---|---|
| Encoding | One stream at a time; a task may process several streams during a flush | Task count and stream order are fixed up front; tasks claim streams as they run | The writer waits for all encode tasks before continuing the flush |
| Decoding | A contiguous group of eligible child fields | The planner shares the decode-parallelism budget across eligible readers | Each parent waits for its child decode tasks |
Parallel encoding
Field writers create stream buffers on the calling thread. Nimble then runs a separate stream-encoding phase with bounded parallelism. The task count is capped by the number of streams and the configured maximum. The minimum-streams-per-task threshold keeps task granularity coarse enough to amortize scheduling overhead across multiple streams.
Nimble sorts streams largest first by buffered size, an estimate of encoding work. Each task encodes one stream at a time. The task count and stream order are fixed up front, but tasks claim streams from that order at runtime, so they may process different numbers. With buffer caching enabled, each encode task gets its own buffer pools to avoid sharing non-thread-safe buffer pools.
The Nimble writer waits for all encode tasks to finish before handing the encoded stripe to the tablet writer. Parallel encoding is configurable: callers provide an encoding executor, a maximum parallelism that caps concurrent encode tasks, and a minimum number of streams per encode task.
Parallel decoding
The reader planner marks row and FlatMap-as-struct readers with multiple readable children as candidates for parallel decoding. For each candidate, it estimates work from the scalar streams beneath its children, chooses a task count to avoid splitting that work too finely, and caps the count at the number of readable children. A single decode-parallelism budget is shared across the reader tree. The planner allocates it level by level, starting at the outermost eligible readers and moving inward.
At batch-read time, Nimble divides a reader's readable child fields into contiguous groups that balance the number of direct children. For example, five fields across two tasks become groups of three and two. Each decode task is a Folly coroutine on the supplied executor. Encode tasks are plain functions on the encoding executor, and the writer waits for all of them to finish.
Within a decode task, child fields are read sequentially while separate decode tasks process sibling groups concurrently. The parent awaits all child decode tasks before returning the batch, and each child reader preserves its dependency order.
Parallel decoding is configurable: callers provide a decode executor, a maximum parallelism that caps decode tasks across the whole reader tree, and a minimum number of scalar streams per decode task.
How parallelism affects wall-clock time and CPU
Parallelism reduces wall-clock time. The same work runs on more threads, and scheduling adds a small overhead, so total CPU time may stay similar or increase. Wide schemas, large batches, and several substantial streams or fields expose more parallel work. A narrow schema, a small batch, or a single dominant stream limits the benefit.
Several factors also limit scaling:
- Memory bandwidth.
- Uneven field costs.
- I/O.
- Serial phases.
Tune encoding and decoding separately on production schemas and batch sizes, comparing wall-clock and CPU time against a serial run with the same input and executor. Each path has two control parameters:
- Maximum parallelism: Keep it at or below the executor's thread count, and use the smallest value that meets the latency target.
- Minimum streams per task: Raise it when small tasks add scheduling overhead.


