How to Create and Use a Fiber Pool for High-Performance Task Execution in libfiber
The fiber_pool class in iqiyi/libfiber implements a dynamic worker pool where each fiber maintains a lock-free task queue, enabling high-throughput scheduling with automatic scaling and minimal context-switch overhead.
The fiber_pool component serves as the core task scheduler in the libfiber coroutine library, distributing work across OS-level fibers to maximize CPU utilization while maintaining predictable latency. By combining dynamic fiber creation with lock-free task boxes, it provides a robust solution for demanding server-side workloads. This guide explains how to instantiate, configure, and operate the pool using the implementation found in cpp/include/fiber/fiber_pool.hpp and cpp/src/fiber_pool.cpp.
Architecture Overview
The fiber pool architecture centers on three primary components that coordinate task distribution and lifecycle management.
Core Components
-
fiber_pool: Manages the dynamic collection of worker fibers, handles task distribution across lock-free queues, and implements automatic scaling based on load. The class definition resides incpp/include/fiber/fiber_pool.hppwith implementation incpp/src/fiber_pool.cpp. -
task_box: Wraps afiber_sbox2<task_fn>lock-free queue and stores the owning fiber along with bookkeeping indices. Each worker fiber owns one box, ensuring queue contention never occurs between different fibers. -
wait_group: Provides a reference-counted synchronization barrier used to block until all submitted tasks complete. Defined incpp/include/fiber/wait_group.hpp, it is the standard mechanism for coordinating shutdown or result aggregation. -
go_*helpers: Found incpp/include/fiber/go_fiber.hpp, these utilities spawn fibers using either shared-stack (go_share) or private-stack (go_stack) modes, which the pool leverages when creating workers.
Creating a Fiber Pool
Instantiate the pool by including the public header and configuring the constructor parameters that control resource limits and behavior.
#include <fiber/fiber_pool.hpp>
// Configure pool parameters
size_t min_fibers = 10; // Minimum workers to maintain
size_t max_fibers = 100; // Maximum workers allowed
int idle_timeout = -1; // Idle milliseconds before shrink (-1 = never)
size_t box_buf = 500; // Queue depth before back-pressure yields
size_t stack_size = 64000; // Bytes per fiber stack
bool share_stack = false; // true enables shared-stack mode
auto pool = std::make_shared<acl::fiber_pool>(
min_fibers, max_fibers, idle_timeout,
box_buf, stack_size, share_stack);
Parameter details:
min_fibers/max_fibers: Establish the bounds for dynamic scaling. The pool starts withmin_fibersand grows up tomax_fiberswhen all boxes report idle starvation.idle_timeout: Specifies how long a worker waits for tasks before exiting. Set to-1to prevent automatic shrinking.box_buf: Acts as a back-pressure threshold; when a box reaches this size, the submitting thread yields to allow consumption.share_stack: Whentrue, the pool usesgo_shareto create fibers, reducing memory footprint by sharing stack space at the cost of slightly higher context-switch overhead.
Submitting Tasks
Tasks are submitted via the exec() method, which accepts any callable and its arguments, wrapping them in a std::function<void()> and pushing the result to an available worker's box.
// Define a task function
void process_data(int id, const std::string& payload) {
// ... computation or I/O bound work ...
}
// Submit with perfect forwarding
pool->exec(process_data, 42, std::string("high-priority job"));
Execution flow:
- Packing:
exec()constructs atask_fn(internallystd::function<void()>) from the callable and arguments. - Distribution: The pool selects an idle box or uses round-robin assignment if all boxes are active.
- Back-pressure: If the selected box reaches
box_bufcapacity, the calling fiber yields to prevent unbounded memory growth. - FIFO guarantee: Tasks within a single box execute in the order submitted, though ordering across different boxes is not guaranteed.
Synchronization and Shutdown
Coordinate completion using wait_group and terminate cleanly with stop().
auto wg = std::make_shared<acl::wait_group>();
wg->add(task_count); // Increment counter for each task
for (size_t i = 0; i < task_count; ++i) {
pool->exec([wg, i]{
// ... perform work ...
wg->done(); // Signal completion
});
}
// Block until all tasks finish
wg->wait();
// Gracefully terminate all fibers
pool->stop(); // Calls kill() on each fiber and waits via internal wait_group
The stop() method iterates through all alive fibers, invokes kill() on each, and blocks until the internal wait_group confirms complete termination.
Complete Working Example
The repository provides a comprehensive demonstration in samples/cxx/fiber_pool/main.cpp. The following excerpt illustrates the canonical pattern for high-throughput task submission and performance measurement:
// Excerpt from samples/cxx/fiber_pool/main.cpp
static void task_run(acl::wait_group* wg, std::atomic_long* res, long long i) {
(*res) += i;
wg->done();
}
static void benchmark(long long count, size_t min, size_t max, size_t buf,
int idle_ms, bool shared) {
auto pool = std::make_shared<acl::fiber_pool>(
min, max, idle_ms, buf, 64000, shared);
auto wg = std::make_shared<acl::wait_group>();
auto result = std::make_shared<std::atomic_long>(0);
wg->add(1);
go[wg, pool, count, result]{
for (long long i = 0; i < count; ++i) {
wg->add(1);
pool->exec(task_run, wg.get(), result.get(), i);
}
wg->done();
};
go[wg, result, pool, count]{
struct timeval begin, end;
gettimeofday(&begin, nullptr);
wg->wait(); // Synchronize all workers
gettimeofday(&end, nullptr);
// ... calculate queries-per-second ...
pool->stop(); // Shutdown pool
};
acl::fiber::schedule(); // Enter fiber scheduler
}
This sample demonstrates dynamic scaling under load, atomic result aggregation, and clean shutdown sequencing.
Summary
fiber_poolmanages dynamic worker fibers incpp/src/fiber_pool.cpp, automatically scaling betweenminandmaxlimits based on task queue pressure.- Lock-free task boxes eliminate contention; each worker owns a
task_boxwrappingfiber_sbox2<task_fn>for zero-lock task retrieval. exec()provides transparent task submission with perfect forwarding, whilebox_bufcontrols memory pressure through cooperative yielding.wait_groupoffers the standard synchronization primitive for blocking on task batches, whereasstop()handles graceful fiber termination viakill()and internal reference counting.- Shared-stack mode (
share_stack = true) reduces memory usage significantly for high-fiber-count scenarios, implemented throughgo_shareingo_fiber.hpp.
Frequently Asked Questions
How does the fiber pool handle dynamic scaling?
The pool monitors box_idle_ after each task execution. When all boxes report idle starvation and the current fiber count remains below max_fibers, the running() loop triggers fiber_create(1) to spawn an additional worker. Conversely, workers exit when their box remains empty for idle_timeout milliseconds and the pool size exceeds min_fibers, as implemented in fiber_pool::running() inside cpp/src/fiber_pool.cpp.
Is the fiber pool thread-safe for task submission?
Yes. The exec() method uses lock-free operations on task_box instances. Each box is owned by exactly one fiber, and cross-box distribution uses atomic indices, ensuring safe concurrent submission from multiple threads or fibers without explicit locking.
What is the difference between shared-stack and private-stack modes?
When share_stack is false, each fiber allocates a private stack of stack_size bytes via go_stack, maximizing isolation at the cost of higher memory usage. When true, fibers utilize go_share to allocate a small save area while sharing a common stack segment, drastically reducing memory consumption but requiring careful avoidance of stack-blocking operations like large automatic arrays.
How do I ensure all tasks complete before shutting down?
Use acl::wait_group to track outstanding tasks: call add(n) before submitting n tasks and invoke done() at the end of each task body. The main thread calls wait() to block until the counter reaches zero, then calls pool->stop() to trigger the shutdown sequence that kills all fibers and waits for their termination.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →