Use a proper Queue in thread.Pool

With lots of tasks the dynamic array takes a big performance hit as its
allocating all the time on pop_front
This commit is contained in:
Waqar Ahmed
2024-11-30 22:29:47 +05:00
parent 314c41ef33
commit 8a27042d24
+7 -6
View File
@@ -9,6 +9,7 @@ package thread
import "base:intrinsics" import "base:intrinsics"
import "core:sync" import "core:sync"
import "core:mem" import "core:mem"
import "core:container/queue"
Task_Proc :: #type proc(task: Task) Task_Proc :: #type proc(task: Task)
@@ -40,7 +41,7 @@ Pool :: struct {
threads: []^Thread, threads: []^Thread,
tasks: [dynamic]Task, tasks: queue.Queue(Task),
tasks_done: [dynamic]Task, tasks_done: [dynamic]Task,
} }
@@ -75,7 +76,7 @@ pool_thread_runner :: proc(t: ^Thread) {
pool_init :: proc(pool: ^Pool, allocator: mem.Allocator, thread_count: int) { pool_init :: proc(pool: ^Pool, allocator: mem.Allocator, thread_count: int) {
context.allocator = allocator context.allocator = allocator
pool.allocator = allocator pool.allocator = allocator
pool.tasks = make([dynamic]Task) queue.init(&pool.tasks)
pool.tasks_done = make([dynamic]Task) pool.tasks_done = make([dynamic]Task)
pool.threads = make([]^Thread, max(thread_count, 1)) pool.threads = make([]^Thread, max(thread_count, 1))
@@ -92,7 +93,7 @@ pool_init :: proc(pool: ^Pool, allocator: mem.Allocator, thread_count: int) {
} }
pool_destroy :: proc(pool: ^Pool) { pool_destroy :: proc(pool: ^Pool) {
delete(pool.tasks) queue.destroy(&pool.tasks)
delete(pool.tasks_done) delete(pool.tasks_done)
for &t in pool.threads { for &t in pool.threads {
@@ -144,7 +145,7 @@ pool_join :: proc(pool: ^Pool) {
pool_add_task :: proc(pool: ^Pool, allocator: mem.Allocator, procedure: Task_Proc, data: rawptr, user_index: int = 0) { pool_add_task :: proc(pool: ^Pool, allocator: mem.Allocator, procedure: Task_Proc, data: rawptr, user_index: int = 0) {
sync.guard(&pool.mutex) sync.guard(&pool.mutex)
append(&pool.tasks, Task{ queue.push_back(&pool.tasks, Task{
procedure = procedure, procedure = procedure,
data = data, data = data,
user_index = user_index, user_index = user_index,
@@ -288,10 +289,10 @@ pool_is_empty :: #force_inline proc(pool: ^Pool) -> bool {
pool_pop_waiting :: proc(pool: ^Pool) -> (task: Task, got_task: bool) { pool_pop_waiting :: proc(pool: ^Pool) -> (task: Task, got_task: bool) {
sync.guard(&pool.mutex) sync.guard(&pool.mutex)
if len(pool.tasks) != 0 { if queue.len(pool.tasks) != 0 {
intrinsics.atomic_sub(&pool.num_waiting, 1) intrinsics.atomic_sub(&pool.num_waiting, 1)
intrinsics.atomic_add(&pool.num_in_processing, 1) intrinsics.atomic_add(&pool.num_in_processing, 1)
task = pop_front(&pool.tasks) task = queue.pop_front(&pool.tasks)
got_task = true got_task = true
} }