Monix Task调度原理疑问:单线程下sleep/shift为何实现并发
Understanding Monix Task's Cooperative Scheduling in Your Examples
Great question! This is a perfect illustration of how Monix's cooperative scheduling works, especially when paired with the SynchronousExecution model. Let's break down exactly what's happening in both scenarios.
First Scenario: Only Task 1 Runs
Let's start with your initial code:
val s1: Scheduler = Scheduler( ExecutionContext.fromExecutor(Executors.newSingleThreadExecutor()), ExecutionModel.SynchronousExecution) def repeat(id: Int): Task[Unit] = Task(println(s"$id ${Thread.currentThread().getName}")) >> repeat(id) val prog: Task[(Unit, Unit)] = (repeat(1), repeat(2)).parTupled prog.runToFuture(s1)
This only outputs 1 repeatedly because:
- You're using
SynchronousExecution, which tells Monix to execute tasks synchronously on the current thread whenever possible, to minimize thread-switching overhead. - The
repeat(1)task is an infinite loop of synchronous, immediate operations:Task(println(...))runs right away on the scheduler's single thread, then immediately recurses torepeat(1)again. - There are no points where the task voluntarily yields control back to the scheduler. Since the thread is stuck in this infinite synchronous loop, the scheduler never gets a chance to pull the second task (
repeat(2)) from its queue and execute it.
Second Scenario: Alternating Execution with Sleep/Shift
When you add Task.sleep(1.millis) or Task.shift to the repeat method:
def repeat(id: Int): Task[Unit] = Task(println(s"$id ${Thread.currentThread().getName}")) >> Task.sleep(1.millis) >> repeat(id)
Suddenly, tasks 1 and 2 alternate execution. Here's why:
What Task.sleep Does
Task.sleepis an asynchronous operation: it doesn't block the thread. Instead, it schedules the continuation (the nextrepeat(id)call) to run after the specified delay, then yields control back to the scheduler.- With the thread now free, the scheduler can pick up the next task in its queue (
repeat(2)), which runs itsprintlnand then also yields control viasleep. - When the delay for task 1 expires, it's added back to the scheduler's queue, and the cycle repeats—resulting in alternating output.
What Task.shift Does
Task.shiftis explicitly designed to voluntarily yield control to the scheduler. It takes the remaining part of the task (the nextrepeat(id)call) and submits it back to the scheduler's queue, then releases the current thread.- Even with a single-threaded scheduler, this gives the scheduler a chance to run the next task in the queue. After task 1 calls
shift, task 2 gets to run itsprintlnand then also callsshift, yielding back. This creates the alternating execution pattern.
Core Principle: Cooperative Scheduling
Monix uses cooperative multitasking for Task execution, unlike Java's threads which use preemptive scheduling. This means:
- Tasks don't get interrupted by the scheduler; they have to actively choose to yield control at specific points.
- Common yield points include asynchronous operations (like
sleep, IO calls),Task.shift,Task.async, and operations that cross asynchronous boundaries. - Without these yield points, a long-running synchronous task can hog the thread indefinitely, preventing other tasks from running—exactly what happened in your first example.
Summary
- In your first case, the infinite synchronous loop never yields control, so task 2 never runs.
- Adding
sleeporshiftintroduces yield points, allowing the scheduler to switch between tasks even on a single thread, simulating concurrency via cooperative task switching.
内容的提问来源于stack exchange,提问作者Kamil Kloch
相关产品推荐
相关产品推荐

