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)