Skip to content

Recycle the streams of finished tasks - #1120

Merged
maleadt merged 1 commit into
mainfrom
tb/stream-pool
Oct 6, 2026
Merged

maleadt merged 1 commit into
mainfrom
tb/stream-pool

Conversation

@maleadt

@maleadt maleadt commented Oct 1, 2026 •

Copy link
Copy Markdown
Member

AMDGPU.jl gives every Julia task its own HIP stream. A task creates its stream the first time it touches the GPU, and the stream is destroyed when the GC finalizes it. CUDA.jl uses the same model, but it works out badly with HIP, because HIP streams are expensive. Creating one takes milliseconds and pins about 8 MiB of host memory: a kernel-argument pool plus another 4 MiB buffer. The GC doesn't know about that memory, so it sees no reason to collect the streams of finished tasks. Code that spawns many short tasks that use the GPU, like Dagger.jl, piles up thousands of streams.

To measure this, I ran a producer/consumer loop in which each iteration spawns one task that fills an array and another that reduces and frees it, with four such loops running concurrently. On a Ryzen 9950X iGPU with ROCm 7.2.4:

iterations  time per 250 iterations   RSS
       250          0.30s              5.2 GB
       750          5.71s             12.9 GB
      1250          9.53s             21.0 GB
      2000          6.42s             33.0 GB

By then, creating a single stream took more than 100 ms. Much of that time was spent in the kernel, compacting memory to pin the new pages for the GPU. In a test that only creates streams, the process died somewhere between 2,000 and 4,000 live streams with:

HW Exception by GPU node-1 (Agent handle: 0x36bcd6e0) reason :GPU Hang

Dagger ran into this earlier and worked around it on its side in JuliaParallel/Dagger.jl@4ebcb22, describing how stream creation "stalls inside the driver, after which unrelated HIP calls start reporting illegal addresses". With the Dagger version from just before that commit, a 4096×4096 stencil sweep reproduces it here: on a single chunk it takes 282 ms per iteration, and with 2×2 chunks it hangs, with every thread stuck in hipStreamCreateWithPriority.

With this PR, task streams come from a pool per device and priority. A task that needs a stream gets the stream of a task that has finished and whose work on that stream has completed. Only when no such stream exists does it create a new one. Tasks that run at the same time never share a stream, however many there are. Once they have finished, up to 32 idle streams per device and priority are kept for reuse, and the rest are left to the GC. A task that switches priorities with priority! gets its own stream for that priority back each time. Results on the same machine:

  • The producer/consumer benchmark runs at about 0.03s per 250 iterations once warmed up, with RSS around 950 MB. The pool never holds more than 11 streams.
  • A loop that spawns 20,000 tasks one after another used to get through about 9,000 of them in two minutes, getting slower all the time. It now finishes all 20,000 in under 4 seconds, including compilation.
  • The old Dagger's stencil sweep takes 37, 42 and 52–71 ms per iteration for 1×1, 2×2 and 4×4 chunks (the 4×4 time varies between runs). Current Dagger with its own workaround takes 42, 58 and 75 ms. Results match a CPU reference.

Arrays remember which stream last used them, so that other tasks can wait for that work. Each stream now also has a generation, which is bumped whenever the stream goes to another task. If an array's stream has changed hands since the array last used it, that work must have finished. The array then doesn't wait for whatever the stream's new owner is doing, although it still checks for kernel exceptions. It also doesn't free its memory on that stream anymore. If the new owner is capturing a graph, the free becomes part of the graph, and launching that graph segfaults (see @gbaraldi's reproducer in the review). Freeing an array that a live task is using while that task captures a graph has the same problem, with or without recycling; that's #1123. HIPStream(handle) now returns the pool's object for the handle of a task's stream, so that such wrappers also notice when the stream is recycled.

The visible change is that a task's default stream may be passed on to another task after it finishes. Code that keeps using a task's stream elsewhere, for example by handing AMDGPU.stream() to another task, should create a stream explicitly with AMDGPU.HIPStream() instead. The docs now say so.

I first implemented this with a fixed pool of 32 streams shared round-robin between tasks, like PyTorch's stream pool. That also fixed the benchmarks, but concurrently running tasks could end up sharing a stream. Work captured into a graph by one task could then include another task's work, and synchronize() would wait for other tasks. CUDA.jl has the same per-task stream problem in a milder form, and gets the same fix in JuliaGPU/CUDA.jl#3317.

The full test suite passes on the gfx1036 iGPU (with AMD_OPT_FLUSH=0, see #1119). The new tests cover reusing the streams of finished tasks, distinct streams for concurrent tasks, the limit on idle streams, long-lived tasks and priority switches not using up the pool, streams that still have work queued or can't be used anymore, arrays that last used a recycled stream, wrapped handles, and graph capture.

@maleadt
maleadt requested a review from jpsamaroo October 1, 2026 08:04
@maleadt
maleadt marked this pull request as draft October 1, 2026 08:35
@maleadt maleadt changed the title Pool task-local HIP streams Recycle the streams of finished tasks Oct 1, 2026
@maleadt
maleadt marked this pull request as ready for review October 1, 2026 09:07

@jpsamaroo jpsamaroo left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM! But I'd leave final say up to @luraess

Comment thread src/hip/stream.jl Outdated
Comment thread src/hip/stream.jl Outdated

@gbaraldi gbaraldi left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Segfault: a free from another task gets recorded into a capture on the recycled stream.

using AMDGPU
a = fetch(Threads.@spawn (x = AMDGPU.ones(Float32, 1 << 20); AMDGPU.synchronize(); x))
capturing, freed = Base.Event(), Base.Event()
t = Threads.@spawn begin  # gets the finished task's stream from the pool
    b = AMDGPU.zeros(Float32, 16); AMDGPU.synchronize()
    AMDGPU.capture() do
        notify(capturing); wait(freed)
        b .+= 1f0
    end
end
wait(capturing)
AMDGPU.unsafe_free!(a)  # hipFreeAsync on a's last stream = t's capturing stream
notify(freed)
g = fetch(t)            # graph now contains a MemFree node
AMDGPU.HIP.launch(AMDGPU.HIP.instantiate(g))  # segfault in hipGraphLaunch

julia -t 8, MI300A + MI250, ROCm 7.2.4: segfaults every run. On the base commit the free errors (HIP error 900) and leaks, but doesn't crash. Suggested fix: if recycled(managed), free on the caller's stream instead of managed.stream.

@gbaraldi gbaraldi left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Other findings (MI300A + MI250):

  1. Skipped wait (take_ownership!). It reads the generation from stream, but keeps the old managed.stream when the two compare equal (same handle). With w = HIPStream(AMDGPU.stream().stream); stream!(w), a 2 s kernel, and then a read from another task: the read doesn't wait and returns 0 instead of 42. Always assign managed.stream = stream, or read the generation from managed.stream.
  2. priority! holds pool slots. priority!(p) / priority!(f, p) leave the swapped-out entry owned by a live task. 40 priority!(:high) do … end blocks on main fill the :high pool, after which every other task creates a fresh stream. More generally, the cap counts streams of live tasks, so 32 long-lived workers disable recycling. Release the entry on restore, and cap idle entries rather than the total.
  3. Stream creation under the lock. HIPStream(priority) runs while holding STREAM_POOL_LOCK. With 64 concurrent first touches, the median wait is 51 ms (MI300A) vs 1.5 ms on base, and when creation is slow every task queues behind it. Reserve the slot under the lock and create the stream outside it.
  4. Nits:
    • A stream in an error state is never isidle, so its slot is lost and re-queried on every lookup.
    • A user-finalized pooled stream could be handed out; check isvalid.
    • The test's "~1s at 100 MHz": memrealtime is 25 MHz on MI250, so tls.jl takes 37 s there.

@maleadt

maleadt commented Oct 1, 2026

Copy link
Copy Markdown
Member Author

Review addressed.

@luraess

luraess commented Oct 1, 2026

Copy link
Copy Markdown
Member

Error seems here as well related to now passing Int128 test (possibly fixed by having AMDGPU_LLVM_Backend_jll on v23.1.1+3 ?).

For the remaining, LGTM after having addressed the review.

@maleadt
maleadt added this pull request to stack #1130 October 2, 2026 12:20
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.
@maleadt

maleadt commented Oct 6, 2026

Copy link
Copy Markdown
Member Author

CI was failing because the new Recycling test kept a stream busy with a kernel that spins on AMDGPU.Device.memrealtime(). That lowers to s_memrealtime, which gfx11 and later don't have, so on the RDNA3 runner the kernel didn't compile (Cannot select: intrinsic %llvm.amdgcn.s.memrealtime). It passed locally because my gfx1036 still has the instruction. The test now loops over s_sleep instead, which every generation supports. I also rebased onto main.

@github-actions

github-actions Bot commented Oct 6, 2026

Copy link
Copy Markdown
Contributor

AMDGPU.jl Benchmarks

Details
Benchmark suite Current: 6608b6a Previous: e0a2787 Ratio
amdgpu/synchronization/context/device 565 ns 540 ns 1.05
amdgpu/synchronization/stream/blocking 227.5 ns 225 ns 1.01
amdgpu/synchronization/stream/nonblocking 310 ns 315 ns 0.98
applications/bitonic_sort 1121041.25 ns 1186679.5 ns 0.94
applications/convolution 91051.5 ns 99008.75 ns 0.92
applications/floyd_warshall 9123158.75 ns 9127855.75 ns 1.00
applications/histogram 816721.75 ns 806011.5 ns 1.01
applications/prefix_sum 188297.5 ns 188540.25 ns 1.00
array/accumulate/Float32/1d 82158.75 ns 78923.75 ns 1.04
array/accumulate/Float32/dims=1 265258.75 ns 265591.25 ns 1.00
array/accumulate/Float32/dims=1L 71246.25 ns 89393.75 ns 0.80
array/accumulate/Float32/dims=2 104691.75 ns 93303.75 ns 1.12
array/accumulate/Float32/dims=2L 3031211.5 ns 3025113.75 ns 1.00
array/accumulate/Int64/1d 87078.75 ns 81546.25 ns 1.07
array/accumulate/Int64/dims=1 244636 ns 241893.5 ns 1.01
array/accumulate/Int64/dims=1L 111751.75 ns 87441.25 ns 1.28
array/accumulate/Int64/dims=2 107069 ns 89436.25 ns 1.20
array/accumulate/Int64/dims=2L 3391724.25 ns 3391916.5 ns 1.00
array/broadcast 28528 ns 29423 ns 0.97
array/construct 2230 ns 2220 ns 1.00
array/copy 40603 ns 38595.75 ns 1.05
array/copyto!/cpu_to_gpu 88883.75 ns 89093.75 ns 1.00
array/copyto!/gpu_to_cpu 89403.75 ns 89286.25 ns 1.00
array/copyto!/gpu_to_gpu 33815.5 ns 34493 ns 0.98
array/iteration/findall/bool 121376.75 ns 111534 ns 1.09
array/iteration/findall/int 111856.5 ns 113686.75 ns 0.98
array/iteration/findfirst/bool 150942.25 ns 152607 ns 0.99
array/iteration/findfirst/int 163042.25 ns 154782.25 ns 1.05
array/iteration/findmin/1d 107749 ns 92681.25 ns 1.16
array/iteration/findmin/2d 91891.25 ns 78871 ns 1.17
array/iteration/logical 192295.25 ns 182967.5 ns 1.05
array/iteration/scalar 298049.25 ns 305309.5 ns 0.98
array/permutedims/2d 58921 ns 59268.25 ns 0.99
array/permutedims/3d 57485.75 ns 38405.5 ns 1.50
array/permutedims/4d 64021 ns 63411 ns 1.01
array/random/rand/Float32 45123 ns 45373.25 ns 0.99
array/random/rand/Int64 36733 ns 44520.75 ns 0.83
array/random/rand!/Float32 38848 ns 41420.75 ns 0.94
array/random/rand!/Int64 37688 ns 38573 ns 0.98
array/random/randn/Float32 71543.5 ns 69016 ns 1.04
array/random/randn!/Float32 53508.25 ns 57575.75 ns 0.93
array/reductions/mapreduce/Float32/1d 109424 ns 81668.75 ns 1.34
array/reductions/mapreduce/Float32/dims=1 68993.5 ns 84246.25 ns 0.82
array/reductions/mapreduce/Float32/dims=1L 833947 ns 838687 ns 0.99
array/reductions/mapreduce/Float32/dims=2 93451.5 ns 84033.75 ns 1.11
array/reductions/mapreduce/Float32/dims=2L 135119.5 ns 136332 ns 0.99
array/reductions/mapreduce/Int64/1d 109334 ns 81271 ns 1.35
array/reductions/mapreduce/Int64/dims=1 88223.75 ns 83796.25 ns 1.05
array/reductions/mapreduce/Int64/dims=1L 842914.75 ns 849044.5 ns 0.99
array/reductions/mapreduce/Int64/dims=2 92408.75 ns 84718.75 ns 1.09
array/reductions/mapreduce/Int64/dims=2L 136126.75 ns 137594.5 ns 0.99
array/reductions/reduce/Float32/1d 109636.5 ns 81216 ns 1.35
array/reductions/reduce/Float32/dims=1 87531.25 ns 81858.5 ns 1.07
array/reductions/reduce/Float32/dims=1L 841907.25 ns 833192 ns 1.01
array/reductions/reduce/Float32/dims=2 93441.5 ns 81401 ns 1.15
array/reductions/reduce/Float32/dims=2L 135052 ns 136399.5 ns 0.99
array/reductions/reduce/Int64/1d 109359 ns 81121.25 ns 1.35
array/reductions/reduce/Int64/dims=1 88203.75 ns 84481 ns 1.04
array/reductions/reduce/Int64/dims=1L 837902 ns 843867 ns 0.99
array/reductions/reduce/Int64/dims=2 92143.75 ns 82028.75 ns 1.12
array/reductions/reduce/Int64/dims=2L 136197 ns 136837 ns 1.00
array/reverse/1d 42775.5 ns 41165.5 ns 1.04
array/reverse/1dL 68498.5 ns 53138.25 ns 1.29
array/reverse/1dL_inplace 53538.25 ns 53878.5 ns 0.99
array/reverse/1d_inplace 37063 ns 35405.5 ns 1.05
array/reverse/2d 46010.5 ns 45943.25 ns 1.00
array/reverse/2dL 89101.25 ns 86261.25 ns 1.03
array/reverse/2dL_inplace 64373.5 ns 65116 ns 0.99
array/reverse/2d_inplace 36443 ns 36043 ns 1.01
array/sorting/1d 319414.75 ns 322337.25 ns 0.99
gemm/tiled 1927162.75 ns 1936378 ns 1.00
gemm/tiled_unbounded 1981478.5 ns 1968285.5 ns 1.01
integration/byval/reference 38691 ns 39721 ns 0.97
integration/byval/slices=1 40381 ns 39841 ns 1.01
integration/byval/slices=2 157442 ns 158692 ns 0.99
integration/byval/slices=3 238254 ns 237193 ns 1.00
integration/volumerhs 4927121 ns 4883600 ns 1.01
kernel/indexing 28073 ns 30012.75 ns 0.94
kernel/indexing_checked 36108 ns 29265.5 ns 1.23
kernel/launch 1220 ns 1200 ns 1.02
kernel/rand 50523.25 ns 38615.5 ns 1.31
latency/import 1526614789 ns 1516524125 ns 1.01
latency/precompile 28082127166 ns 28086448609 ns 1.00
latency/ttfp 2315622127 ns 2303896666 ns 1.01
stencil/diffusion3d 1584615.25 ns 1623923.25 ns 0.98
stencil/diffusion3d_checked 1626153.5 ns 1660986.5 ns 0.98

This comment was automatically generated by workflow using github-action-benchmark.

@maleadt
maleadt merged commit 130f06f into main Oct 6, 2026
14 of 15 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants