Make synchronization cooperative, using GPUToolbox - #654
Merged
Merged
Conversation
`synchronize` blocked the calling thread in the driver until the work had completed, so no other task could run on it in the meantime, and no other thread could run the GC either. KernelInterface requires `synchronize` to be cooperative. Use GPUToolbox's `cooperative_wait`, as CUDA.jl does: poll a non-blocking query first, which keeps the latency of short operations low, and then block in the driver (GC-safe) on a separate thread, while the calling task waits without blocking the scheduler. This applies to `synchronize()` and `synchronize(::oneStream)`, i.e., to user code, KernelAbstractions and the synchronizing copies; `synchronize(; blocking=true)` restores the old behavior. The command list and queue methods, which also run from finalizers, keep blocking. With `ONEAPI_SYNC_EACH_SUBMISSION` set, as on Aurora's LTS stack, the wait after every launch is cooperative too; otherwise it would block the thread, and leave `synchronize` nothing to wait for.
Contributor
|
Your PR requires formatting changes to meet the project's style guidelines. Click here to view the suggested changes.diff --git a/lib/level-zero/synchronization.jl b/lib/level-zero/synchronization.jl
index e3a3a4e..f6c2aff 100644
--- a/lib/level-zero/synchronization.jl
+++ b/lib/level-zero/synchronization.jl
@@ -30,10 +30,12 @@ Base.isdone(queue::ZeCommandQueue) =
# the blocking synchronization, marked GC-safe so that it doesn't keep the GC from running
gcsafe_synchronize(list::ZeImmediateCommandList) =
@gcsafe_ccall libze_loader.zeCommandListHostSynchronize(
- list::ze_command_list_handle_t, typemax(UInt64)::UInt64)::ze_result_t
+ list::ze_command_list_handle_t, typemax(UInt64)::UInt64
+)::ze_result_t
gcsafe_synchronize(queue::ZeCommandQueue) =
@gcsafe_ccall libze_loader.zeCommandQueueSynchronize(
- queue::ze_command_queue_handle_t, typemax(UInt64)::UInt64)::ze_result_t
+ queue::ze_command_queue_handle_t, typemax(UInt64)::UInt64
+)::ze_result_t
"""
nonblocking_synchronize(list_or_queue)
@@ -44,8 +46,10 @@ while this one waits.
"""
function nonblocking_synchronize(obj::SyncObject)
# when polling found the work to be done, synchronize again to check for errors
- res = @something(cooperative_wait(gcsafe_synchronize, obj; isdone=Base.isdone),
- gcsafe_synchronize(obj))
+ res = @something(
+ cooperative_wait(gcsafe_synchronize, obj; isdone = Base.isdone),
+ gcsafe_synchronize(obj)
+ )
res == RESULT_SUCCESS || throw_api_error(res)
return
end
diff --git a/src/context.jl b/src/context.jl
index 6c9e188..36f0cd0 100644
--- a/src/context.jl
+++ b/src/context.jl
@@ -401,7 +401,7 @@ println("GPU work completed")
See also: [`global_stream`](@ref), [`context`](@ref), [`device`](@ref)
"""
-function oneL0.synchronize(s::oneStream; blocking::Bool=false)
+function oneL0.synchronize(s::oneStream; blocking::Bool = false)
sync = blocking ? oneL0.synchronize : oneL0.nonblocking_synchronize
sync(s.list)
q = s.queue
@@ -412,8 +412,8 @@ function oneL0.synchronize(s::oneStream; blocking::Bool=false)
return
end
-function oneL0.synchronize(; blocking::Bool=false)
- oneL0.synchronize(global_stream(context(), device()); blocking)
+function oneL0.synchronize(; blocking::Bool = false)
+ return oneL0.synchronize(global_stream(context(), device()); blocking)
end
# Julia → MKL ordering: everything Julia appended to the task's immediate list must be
diff --git a/test/execution.jl b/test/execution.jl
index c32cb17..b0f232b 100644
--- a/test/execution.jl
+++ b/test/execution.jl
@@ -753,8 +753,10 @@ end
# it, which makes these tests independent of timing: synchronizing can only return after
# that task got to run. if it cannot (e.g., because synchronizing blocks the thread), a
# watchdog on another thread signals the event after a while instead, and records that.
- pool = oneL0.ZeEventPool(context(), 1, device();
- flags=oneL0.ZE_EVENT_POOL_FLAG_HOST_VISIBLE)
+ pool = oneL0.ZeEventPool(
+ context(), 1, device();
+ flags = oneL0.ZE_EVENT_POOL_FLAG_HOST_VISIBLE
+ )
gate = pool[1]
timeout = UInt64(60_000_000_000) # ns
@@ -777,10 +779,11 @@ end
try
# the watchdog waits on one of `cooperative_wait`'s worker threads, which keep
# running when this thread is blocked
- watchdog = @async oneL0.cooperative_wait(gate; spin=false) do gate
+ watchdog = @async oneL0.cooperative_wait(gate; spin = false) do gate
watching[] = true
res = oneL0.@gcsafe_ccall oneL0.libze_loader.zeEventHostSynchronize(
- gate::oneL0.ze_event_handle_t, timeout::UInt64)::oneL0.ze_result_t
+ gate::oneL0.ze_event_handle_t, timeout::UInt64
+ )::oneL0.ze_result_t
res == oneL0.RESULT_NOT_READY || return false
oneL0.signal(gate)
return true
@@ -812,13 +815,13 @@ end
# compile everything beforehand, as that could wait for the GPU to become idle
a = oneArray{Int32}(undef, 64)
- @oneapi items=64 fill_kernel(a, Int32(0))
+ @oneapi items = 64 fill_kernel(a, Int32(0))
synchronize()
@test gated(synchronize)
@test gated() do
oneL0.sync_each_submission(false) do
- @oneapi items=64 fill_kernel(a, Int32(1))
+ @oneapi items = 64 fill_kernel(a, Int32(1))
end
synchronize()
end
@@ -826,7 +829,7 @@ end
@test gated() do
oneL0.sync_each_submission(false) do
- @oneapi items=64 fill_kernel(a, Int32(2))
+ @oneapi items = 64 fill_kernel(a, Int32(2))
end
synchronize(oneAPI.global_stream(context(), device()))
end
@@ -835,14 +838,14 @@ end
# synchronizing after every launch, as on the LTS stack
@test gated() do
oneL0.sync_each_submission(true) do
- @oneapi items=64 fill_kernel(a, Int32(3))
+ @oneapi items = 64 fill_kernel(a, Int32(3))
end
end
@test Array(a) == fill(Int32(3), 64)
# blocking synchronization is still available
- @oneapi items=64 fill_kernel(a, Int32(4))
- @test synchronize(; blocking=true) === nothing
+ @oneapi items = 64 fill_kernel(a, Int32(4))
+ @test synchronize(; blocking = true) === nothing
@test Array(a) == fill(Int32(4), 64)
end
|
Codecov Report❌ Patch coverage is
Additional details and impacted files@@ Coverage Diff @@
## main #654 +/- ##
==========================================
+ Coverage 80.88% 80.90% +0.02%
==========================================
Files 56 57 +1
Lines 4090 4105 +15
==========================================
+ Hits 3308 3321 +13
- Misses 782 784 +2 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
3.3.1 keeps cooperative_wait's worker threads from running finalizers. On the LTS stack, freeing memory from a finalizer drains every stream, so a worker could block in there before waking the task it waited for.
The test counted how often a timer task ran while a calibrated kernel kept the GPU busy, which isn't reliable on a loaded machine. Instead, keep the work on the stream waiting for a host-signalled event, which another task on the same thread only signals after yielding many times, and check that synchronizing returned after it did. A watchdog on one of cooperative_wait's worker threads signals the event after a timeout and records that, so blocking synchronization fails the test instead of hanging it. Also cover synchronizing a stream, and the wait after every launch with ONEAPI_SYNC_EACH_SUBMISSION.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
synchronize()currently blocks the calling thread inside the Level Zero driver until the GPU is done. While it waits, no other Julia task can run on that thread, so work like launching kernels from another task, doing I/O, or answering a heartbeat stalls behind a long-running kernel. KernelInterface also requires back-ends to synchronize cooperatively.This PR makes it cooperative using
cooperative_waitfrom GPUToolbox 3.3.1 (JuliaGPU/GPUToolbox.jl#23), the same implementation CUDA.jl, OpenCL.jl and KernelAbstractions' POCL back-end are switching to:zeCommandListHostSynchronizewith a zero timeout), so they don't pay for involving another thread.For example, a task that ticks every millisecond keeps ticking while another task waits for a long kernel:
What changes:
synchronize()andsynchronize(::oneStream), which also back KernelAbstractions'synchronizeand the synchronizing copies, now wait cooperatively.synchronize(; blocking=true)gives the old behavior.ONEAPI_SYNC_EACH_SUBMISSION(the Aurora LTS workaround), the wait after every launch is cooperative too. Otherwise that wait would block the thread and leave nothing forsynchronizeto wait for.oneL0.synchronize(list)) and the stream drain before freeing memory on the LTS stack keep blocking, because they run from finalizers or with finalizers disabled, where a task can't yield. The internal barriers that order Julia and oneMKL work are also left as they are.The tests don't depend on timing. They put a wait on a host-signalled Level Zero event in front of the work on the task's stream, so it can't complete until the host signals that event. Then they call
synchronize()(orsynchronize(stream), or launch a kernel with sync-each-submission enabled) while another task on the same thread signals the event after yielding 10,000 times, well past the polling phase. Each test checks that the call returned only after that task had run, and that the kernel's result is there. If synchronization blocked the thread, the signalling task couldn't run and the call would hang. To avoid that, a watchdog on one ofcooperative_wait's worker threads waits for the event with a 60 s timeout, signals it itself if it times out, and records that, so the test fails instead. I checked this by swapping blocking synchronization back in: every gated test fails after the timeout, on 1.10, 1.12 and 1.13, with 1 or 4 threads, and with and withoutONEAPI_SYNC_EACH_SUBMISSION.Cost, measured on an Iris Xe (sagittarius), alternating blocking and cooperative synchronization after a launch:
Short operations are unaffected. Longer ones pay a few tens of µs, about 1%, to wake up the waiting task. That machine was busy running other jobs, which makes the wake-ups slower than they'd normally be.
The LTS check failed once on
reductions of host-accessible arrays(sum(a)returned 0 on Julia 1.13). That failure isn't caused by this PR: the same test failed the same way on Julia 1.13 before this PR existed, on #649's own branch (5c33db5) and on theunsafe_wrapbranch (24ba399). Synchronization was still blocking then. Across the LTS runs that include the test, it failed in 3 of 6 runs on Julia 1.13 and never on 1.10.