You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

为何独立线程发起请求时subscribeOn会死锁?请求堆积原因求解

Why Request Backlog Causes Deadlock with Eager/Blocking create + subscribeOn

Great question—this is one of those tricky Reactor details that confuses even experienced developers. Let’s break down exactly what’s happening here, step by step.

First, Let’s Recap Default subscribeOn Behavior

When you use subscribeOn(Scheduler) without specifying the second parameter, it defaults to subscribeOn(scheduler, true). That true means the entire subscription flow—including the propagation of request signals from the subscriber back to the emitter—will be dispatched to the specified Scheduler’s thread pool.

What Makes an Eager/Blocking create Different?

The create(Consumer, FluxSink.OverflowStrategy) variant you’re referring to is "eager" because it starts emitting data immediately when the subscription is triggered—before it even receives any request signals from the subscriber. If it’s also blocking, it might:

  • Run an infinite loop emitting data, blocking the thread until the sink’s buffer is full
  • Use blocking I/O (like reading from a file synchronously) to generate data, holding the thread hostage until the operation finishes

This type of emitter doesn’t respect Reactor’s backpressure rules out of the box—it pushes data regardless of whether the subscriber has requested it.

The Deadlock Chain Explained

Here’s how the pieces line up to create a deadlock:

  1. You call subscribeOn(Schedulers.boundedElastic()) (default true), so Reactor schedules the entire subscription logic (including running your create callback) onto a thread from the boundedElastic pool.
  2. Your eager/blocking create callback starts running on that thread: it immediately begins emitting data, fills up the FluxSink’s buffer, and then blocks waiting for more capacity (or just keeps blocking indefinitely if it’s an infinite loop).
  3. Meanwhile, the subscriber tries to send a request signal to tell the emitter it can accept more data. But because subscribeOn is set to true, this request signal is also scheduled to run on the same boundedElastic thread that’s already stuck in your blocking create callback.
  4. The request signal can never execute because the thread is occupied by the blocking emitter. The emitter, in turn, can’t proceed because it’s waiting for the request signal to free up buffer space. You’ve got a classic deadlock: two tasks waiting on each other, with no way to break the cycle.

Why subscribeOn(scheduler, false) Fixes It

Setting the second parameter to false (named requestOnSeparateThread) changes one critical behavior: the propagation of request signals will no longer be dispatched to the specified Scheduler. Instead, request signals are handled on the thread that initiated the subscription (or the thread that’s processing the subscriber’s demand).

This breaks the deadlock because:

  • The create callback still runs on the Scheduler thread (so your blocking logic doesn’t block the main thread)
  • The request signal can now bypass the blocked Scheduler thread and reach the emitter directly
  • The emitter receives the green light to keep emitting (or adjust its flow based on backpressure), unblocking the cycle

Example of the Deadlock Risk

Here’s a simplified code snippet that could trigger this issue:

// Risk of deadlock!
Flux.create(sink -> {
    // Eager, blocking emission: infinite loop sending data
    while (!sink.isCancelled()) {
        sink.next("message-" + System.currentTimeMillis());
        // Simulate blocking work
        try { Thread.sleep(50); } catch (InterruptedException e) {}
    }
}, FluxSink.OverflowStrategy.BUFFER)
.subscribeOn(Schedulers.boundedElastic()) // Defaults to true
.subscribe();

In this case, the boundedElastic thread gets stuck in the while loop, and the request signal can never reach the sink to manage backpressure—leading to a deadlock as the buffer fills and the thread stays blocked.

Switching to subscribeOn(Schedulers.boundedElastic(), false) would let the request signal reach the sink without waiting for the blocked thread, allowing backpressure to work as intended.

内容的提问来源于stack exchange,提问作者anonk

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.26 10:46:19