From 6608b6a4675da3df92a3a6e3692a03ceb1dd4706 Mon Sep 17 00:00:00 2001 From: Tim Besard Date: Thu, 1 Oct 2026 11:07:06 +0200 Subject: [PATCH] Recycle the streams of finished tasks Every task got its own HIP stream, created on first use and only destroyed when the GC finalized it. HIP streams are expensive: creating one takes milliseconds and pins ~8 MiB of host memory. Since the GC is in no hurry to collect finished tasks, code that spawns many short GPU tasks piles up thousands of streams, making stream creation take over 100ms each and eventually hanging the GPU when pinned memory runs out. Instead, keep the streams of tasks in a pool per device and priority, and hand the stream of a task that has finished, and whose work has completed, to the next task that needs one. Tasks running at the same time never share a stream, and up to 32 idle streams are kept around. Handing a stream to another task bumps its generation, so that memory last used by a recycled stream knows that its work has finished. Such memory doesn't wait for the stream's new owner, and isn't freed on that stream either, since the new owner may be capturing it. --- docs/src/usage/multitasking.md | 2 + src/hip/stream.jl | 108 ++++++++++++++++++++++++++- src/memory.jl | 40 ++++++++-- src/tls.jl | 11 ++- test/core/tls.jl | 131 +++++++++++++++++++++++++++++++++ 5 files changed, 278 insertions(+), 14 deletions(-) diff --git a/docs/src/usage/multitasking.md b/docs/src/usage/multitasking.md index f136e8ba5..177451697 100644 --- a/docs/src/usage/multitasking.md +++ b/docs/src/usage/multitasking.md @@ -34,6 +34,8 @@ Because launches are asynchronous, synchronize before reading results back or ti AMDGPU.@sync @roc groupsize=256 gridsize=n kernel(args...) ``` +HIP streams are expensive: each one pins several MiB of host memory. To keep applications that spawn many short-lived tasks from accumulating streams, the stream of a task is recycled once the task has finished and all work on it has completed. Tasks that are running at the same time never share a stream, but a newly started task may get the stream of an earlier one. If you need a stream that outlives the task that uses it, e.g., to pass it on to other tasks, create one explicitly with `AMDGPU.HIPStream()` and activate it with `AMDGPU.stream!`. + Streams also carry a priority (`:normal`, `:low`, `:high`) to bias scheduling. See [Streams](@ref) for stream priorities, synchronization details, and the blocking-vs-nonblocking preference. ## Using multiple GPUs diff --git a/src/hip/stream.jl b/src/hip/stream.jl index b2c1dbc9b..5ba8200a9 100644 --- a/src/hip/stream.jl +++ b/src/hip/stream.jl @@ -8,6 +8,10 @@ mutable struct HIPStream ctx::HIPContext Base.@atomic valid::Bool + + # bumped when the stream is handed to another task (see `task_stream`), which only + # happens when it is idle, so work submitted during earlier generations has finished. + Base.@atomic generation::Int end """ @@ -26,7 +30,7 @@ function HIPStream(priority::Symbol = :normal) stream_ref = Ref{hipStream_t}() hipStreamCreateWithPriority(stream_ref, 0, priority_int) d = device() - stream = HIPStream(stream_ref[], priority, d, HIPContext(d), true) + stream = HIPStream(stream_ref[], priority, d, HIPContext(d), true, 0) return finalizer(stream) do s Base.@atomic s.valid = false AMDGPU.context!(s.ctx) do @@ -35,9 +39,97 @@ function HIPStream(priority::Symbol = :normal) end end +# Every task gets its own default stream, but HIP streams are expensive: creating one +# takes milliseconds and pins ~8 MiB of host memory. Since the GC is in no hurry to +# collect finished tasks (and with them, their streams), code that spawns many +# short-lived tasks would pile up thousands of streams. Instead, recycle the streams of +# tasks that have finished, keeping up to `STREAM_POOL_IDLE` unused ones per device and +# priority. +const STREAM_POOL_IDLE = 32 +struct PooledStream + stream::HIPStream + owner::WeakRef +end +const STREAM_POOLS = Dict{Tuple{Int,Symbol}, Vector{PooledStream}}() +const STREAM_POOL_LOCK = ReentrantLock() + +function task_stream(priority::Symbol = :normal) + # finalizers can't wait for the pool's lock, so give them a stream of their own + GC.in_finalizer() && return HIPStream(priority) + + key = (device_id(device()), priority) + task = current_task() + stream = Base.@lock STREAM_POOL_LOCK begin + claim_stream!(get!(Vector{PooledStream}, STREAM_POOLS, key), task) + end + stream === nothing || return stream + + # creating a stream can be slow, so don't make other tasks wait for it + stream = HIPStream(priority) + Base.@lock STREAM_POOL_LOCK begin + push!(STREAM_POOLS[key], PooledStream(stream, WeakRef(task))) + end + return stream +end + +function claim_stream!(pool::Vector{PooledStream}, task::Task) + candidate = nothing + idle = 0 + i = 1 + while i <= length(pool) + entry = pool[i] + owner = entry.owner.value + keep = if owner === task + # a task that switches back and forth between priorities keeps its streams + isvalid(entry.stream) && return entry.stream + false + elseif owner !== nothing && !istaskdone(owner::Task) + true + elseif !isvalid(entry.stream) + false + else + status = query(entry.stream) + if status == hipErrorNotReady + # don't make a new task wait for work that the previous owner left behind + true + elseif status != hipSuccess + # the stream is in an error state + false + elseif candidate === nothing + candidate = entry + true + else + (idle += 1) <= STREAM_POOL_IDLE + end + end + keep ? (i += 1) : deleteat!(pool, i) + end + candidate === nothing && return nothing + + candidate.owner.value = task + generation = Base.@atomic :monotonic candidate.stream.generation + Base.@atomic :release candidate.stream.generation = generation + 1 + return candidate.stream +end + +function query(s::HIPStream) + # querying a stream is prohibited while another one is being captured in global + # mode, even though it doesn't interfere with the capture, so temporarily relax that + mode = Ref(hipStreamCaptureModeRelaxed) + hipThreadExchangeStreamCaptureMode(mode) + try + return unchecked_hipStreamQuery(s) + finally + hipThreadExchangeStreamCaptureMode(mode) + end +end + +# only bumped while holding `STREAM_POOL_LOCK`, but read without it +generation(s::HIPStream) = Base.@atomic :acquire s.generation + isvalid(s::HIPStream) = s.valid -default_stream() = HIPStream(C_NULL, :normal, device(), HIPContext(), true) +default_stream() = HIPStream(C_NULL, :normal, device(), HIPContext(), true, 0) """ HIPStream(stream::hipStream_t) @@ -46,8 +138,18 @@ Create HIPStream from `hipStream_t` handle. Device is the default device that's currently in use. """ function HIPStream(stream::hipStream_t) + # the streams of tasks get recycled, which only the pool's objects keep track of + if !GC.in_finalizer() + Base.@lock STREAM_POOL_LOCK begin + for pool in values(STREAM_POOLS), entry in pool + s = entry.stream + s.stream == stream && isvalid(s) && return s + end + end + end + d = device() - HIPStream(stream, priority(stream), d, HIPContext(d), true) + HIPStream(stream, priority(stream), d, HIPContext(d), true, 0) end function isdone(stream::HIPStream) diff --git a/src/memory.jl b/src/memory.jl index e2b0277a6..1d5d8ff00 100644 --- a/src/memory.jl +++ b/src/memory.jl @@ -409,19 +409,31 @@ end mutable struct Managed{M} const mem::M const lock::ReentrantLock + # which stream is currently using the memory, and the generation of that stream stream::HIPStream + generation::Int dirty::Bool captured::Bool function Managed(mem; stream=AMDGPU.stream(), dirty=true, captured=false) - new{typeof(mem)}(mem, ReentrantLock(), stream, dirty, captured) + new{typeof(mem)}(mem, ReentrantLock(), stream, HIP.generation(stream), + dirty, captured) end end +# if the stream has been handed to another task since it last used the memory, that work +# has finished, and waiting for the stream would only wait for the new task's work. +recycled(m::Managed) = m.generation != HIP.generation(m.stream) + function synchronize(m::Managed) Base.@lock m.lock begin m.dirty || return - synchronize(m.stream) + if recycled(m) + # the work has finished, but may have raised an exception + throw_if_exception(m.stream.device) + else + synchronize(m.stream) + end m.dirty = false return end @@ -441,6 +453,7 @@ function take_ownership!(managed::Managed; stream::HIPStream=AMDGPU.stream()) synchronize(managed) managed.stream = stream end + managed.generation = HIP.generation(managed.stream) managed.dirty = true return managed @@ -448,7 +461,8 @@ end # Fast-path ownership transfer for the kernel-launch path @inline function take_ownership_fast!(managed::Managed, stream::HIPStream) - (managed.stream === stream && managed.dirty) && return + (managed.stream === stream && managed.dirty && + managed.generation == HIP.generation(stream)) && return Base.@lock managed.lock take_ownership!(managed; stream) return end @@ -522,7 +536,7 @@ function pool_free(managed::Managed{M}) where M try time = Base.@elapsed Base.@lock managed.lock begin - _pool_free(managed.mem, managed.stream) + _pool_free(managed) end Base.@atomic alloc_stats.free_count += 1 Base.@atomic alloc_stats.free_bytes += sz @@ -536,9 +550,19 @@ function pool_free(managed::Managed{M}) where M return end -function _pool_free(buf, stream::HIPStream) - if !HIP.isvalid(stream) - stream = AMDGPU.default_stream() +function _pool_free(managed::Managed) + buf = managed.mem + AMDGPU.context!(() -> Mem.free(buf; stream=free_stream(managed)), buf.ctx) +end + +function free_stream(managed::Managed) + if !HIP.isvalid(managed.stream) + return AMDGPU.default_stream() + elseif recycled(managed) && !GC.in_finalizer() + # the stream now belongs to another task, which may be capturing it. + # finalizers can keep using it, because capturing disables the GC. + return AMDGPU.stream() + else + return managed.stream end - AMDGPU.context!(() -> Mem.free(buf; stream), buf.ctx) end diff --git a/src/tls.jl b/src/tls.jl index 6d2823abd..d03752037 100644 --- a/src/tls.jl +++ b/src/tls.jl @@ -96,7 +96,7 @@ end function stream(state::TaskLocalState)::HIPStream i = device_id(state.device) if state.streams[i] ≡ nothing - state.streams[i] = HIPStream(:normal) + state.streams[i] = HIP.task_stream(:normal) else state.streams[i] end @@ -107,6 +107,11 @@ end Get the HIP stream that should be used as the default one for the currently executing task. + +Each task gets its own stream, which may be handed to another task once the task has +finished and all work on the stream has completed. If you need a stream that outlives the +task, or that is never shared, create one with `HIPStream()` and activate it using +[`AMDGPU.stream!`](@ref). """ stream()::HIPStream = stream(task_local_state!()) @@ -180,7 +185,7 @@ function priority!(p::Symbol) state = task_local_state!() state.stream.priority == p && return p - state.streams[device_id(state.device)] = HIPStream(p) + state.streams[device_id(state.device)] = HIP.task_stream(p) return p end @@ -201,7 +206,7 @@ function priority!(f::Function, p::Symbol) old_s = state.stream swap = p != old_s.priority - swap && (state.streams[idx] = HIPStream(p);) + swap && (state.streams[idx] = HIP.task_stream(p);) return try f() diff --git a/test/core/tls.jl b/test/core/tls.jl index bcf3af932..c40f4edac 100644 --- a/test/core/tls.jl +++ b/test/core/tls.jl @@ -88,6 +88,137 @@ end @test s1 ≡ s2 end + @testset "Recycling" begin + idle_limit = AMDGPU.HIP.STREAM_POOL_IDLE + pool(priority=:normal) = + AMDGPU.HIP.STREAM_POOLS[(AMDGPU.HIP.device_id(AMDGPU.device()), priority)] + function finished(entry) + owner = entry.owner.value + return owner === nothing || istaskdone(owner) + end + + # call `f` with the streams of `n` tasks that are alive at the same time + function with_concurrent_streams(f, n) + ready = Channel{HIPStream}(Inf) + release = Base.Event() + tasks = [Threads.@spawn begin + try + put!(ready, AMDGPU.stream()) + catch err + # don't leave the caller waiting for our stream + close(ready, err) + rethrow() + end + wait(release) + end for _ in 1:n] + try + f([take!(ready) for _ in tasks]) + finally + notify(release) + foreach(wait, tasks) + end + end + + # s_sleep instead of a clock: gfx11+ lacks s_memrealtime + function nap(n) + for _ in 1:n + AMDGPU.Device.device_sleep(Int32(127)) + end + return + end + keep_busy() = @roc nap(300_000) # for about a second + set42!(a) = (a[1] = 42; nothing) + + # finished tasks hand their stream to new ones, without having to wait for the GC + streams = [fetch(Threads.@spawn AMDGPU.stream()) for _ in 1:2idle_limit] + @test length(unique(streams)) <= idle_limit + @test all(s -> any(entry -> entry.stream === s, pool()), streams) + + # tasks running at the same time never share a stream, even beyond the pool's size, + # but only a limited number of idle streams is kept around afterwards + streams = with_concurrent_streams(identity, idle_limit + 8) + @test allunique(streams) + foreach(AMDGPU.synchronize, streams) + @test fetch(Threads.@spawn AMDGPU.stream()) in streams + @test count(finished, pool()) <= idle_limit + 1 + + # tasks that keep their stream don't prevent others from being recycled + with_concurrent_streams(idle_limit) do _ + streams = [fetch(Threads.@spawn AMDGPU.stream()) for _ in 1:8] + @test length(unique(streams)) <= 2 + end + + # switching priorities doesn't make a task take more and more streams + streams = fetch(Threads.@spawn begin + [AMDGPU.priority!(AMDGPU.stream, :high) for _ in 1:2idle_limit] + end) + @test allequal(streams) + @test fetch(Threads.@spawn begin + s1 = AMDGPU.stream() + AMDGPU.priority!(:high) + s2 = AMDGPU.stream() + AMDGPU.priority!(:normal) + s3 = AMDGPU.stream() + AMDGPU.priority!(:high) + s4 = AMDGPU.stream() + s1 === s3 && s2 === s4 && s1 !== s2 + end) + + # a stream that still has work queued isn't handed to another task + busy = fetch(Threads.@spawn begin + keep_busy() + AMDGPU.stream() + end) + with_concurrent_streams(idle_limit) do streams + if !AMDGPU.HIP.isdone(busy) + @test !(busy in streams) + end + end + AMDGPU.synchronize(busy) + + # streams that can't be used anymore are removed from the pool + s = fetch(Threads.@spawn AMDGPU.stream()) + finalize(s) + @test fetch(Threads.@spawn AMDGPU.stream()) !== s + @test !any(entry -> entry.stream === s, pool()) + + # memory knows that the work of a stream's previous owner has finished, so it + # doesn't wait for the new owner, nor gets freed on its stream (which the new owner + # may be capturing) + a = fetch(Threads.@spawn begin + a = ROCArray([42]) + AMDGPU.synchronize() + a + end) + with_concurrent_streams(idle_limit) do streams + @test a.buf[].stream in streams + @test AMDGPU.recycled(a.buf[]) === true + @test Array(a) == [42] + @test AMDGPU.free_stream(a.buf[]) === AMDGPU.stream() + end + + # wrapping the handle of a task's stream gives the pool's object, which keeps track + # of recycling, and using memory through another object for the same stream doesn't + # mistake it for a recycled stream + a = fetch(Threads.@spawn begin + s = AMDGPU.stream() + @test HIPStream(s.stream) === s + a = ROCArray([0]) + AMDGPU.synchronize() + @test AMDGPU.HIP.generation(s) > 0 + AMDGPU.stream!(HIPStream(s.stream, s.priority, s.device, s.ctx, true, 0)) + @roc set42!(a) + a + end) + @test AMDGPU.recycled(a.buf[]) === false + @test Array(a) == [42] + + # looking for a stream to recycle doesn't break graph capture + graph = AMDGPU.capture() do + @test fetch(Threads.@spawn AMDGPU.stream()) !== AMDGPU.stream() + end + end + @testset "Validity" begin s = HIPStream() @test AMDGPU.HIP.isvalid(s)