文章目录
-
- Task
- @async
- Channel
- 多协程调度
- 取消任务
- 总结
Task
【Task】在Julia中表示协程,由Julia运行时来调度,不经过操作系统来实现异步工作。实际使用中,可直接在非异步的函数外面套一层Task,并通过【schedule】来调度,通过【wait】来等待协程结束,从而阻塞运行时,以实现同步。示例如下,千万注意,在wait之前必须schedule,否则由于协程尚未运行,那么wait将永远也等不到它结束,导致卡死。
t = Task(() -> println("task"))
schedule(t) # 手动加入调度队列
wait(t) # 等待完成
协程的返回值存储在t.result中,示例如下
t = Task(() -> "task")
schedule(t)
wait(t) # 等待完成
t.result # task
【fetch】可将上述代码的最后两行合二为一
t = Task(() -> "task")
schedule(t)
fetch(t) # 直接返回task
@async
【@async】是一个宏,可将表达式打包成Task,并马上执行schedule。例如,在命令行中运行下面的程序,命令行会先输出启动,再输出完成。
function test()
sleep(2)
println("Task完成")
end
t = @async test()
println("Task启动")
如果并不追求test的复用,那么可以直接将代码段写在@async后面
t = @async begin
sleep(2)
print("Task完成")
end
千万注意一点,@async后面的代码如果超过一行,必须用begin…end括起来,否则会只执行离它最近的那条。
Channel
【Channel】是多协程之间的同步通道,有了这个,就可以实现多个Task之间的通信,示例如下
ch = Channel{Int}(32)
@async for x in ch # 若 ch 空,自动挂起
println("Got: ", x)
end
put!(ch, 3) # 命令行立即打印 Got: 3
@async for i in 1:10
put!(ch, i) # 若 ch 满,自动挂起当前 Task
end
# 命令行打印Got: 1到Got: 10
close(ch) # 通知消费者结束
此代码中,先创建了一个包含32个Int的通道,然后创建接收通道内容的协程。接下来,先在系统线程载入一个值,然后新建一个协程载入10个值。最后关闭通道。
其中,put!很容易理解,就是为通道新装进去一个值。
【for x in ch】并不是普通的for循环,而是由Channel驱动的、可挂起/恢复的异步迭代协议。在普通的for循环中,若ch为空,则直接退出循环,而在Channel的for循环中,若ch为空,则当前协程自动挂起,直到close(ch),才算循环结束。
这种看上去像是for循环的表达式,其实际运行逻辑类似下面的过程,其中take!是与put!互相对偶的函数,表示从通道中取出一个值(先入先出)。
while (ch not closed) #伪代码
x = take!(ch)
end
多协程调度
【@sync】也是一个宏,可以统一wait块内所有@async代码执行结束,如下例所示,程序输出了ACB,并阻塞了2秒钟。
@sync begin
@async begin sleep(1); println("A") end
@async begin sleep(2); println("B") end
@async begin sleep(1); println("C") end
end
Julia只提供了这一种调度模式,即等待所有任务完成。在实际工作中,另一种等待模式也十分常见,即任一任务完成(或失败)就立即返回,此功能可通过Channel来实现,示例如下。
function wait_any(tasks)
ch = Channel{Any}(length(tasks))
for t in tasks
@async begin
try
put!(ch, fetch(t))
catch
end
end
end
res = take!(ch) # 等待第一个成功结果
close(ch)
return res
end
tasks = [
@async begin sleep(4); "A"; end
@async begin sleep(3); "B"; end
@async begin sleep(3); "C"; end
]
wait_any(tasks) # 返回一个 B
取消任务
Julia中并未提供取消任务的函数,但取消任务的机制,可通过Channel轻而易举的实现。下面从两种应用场景出发,来实现任务的取消功能。
其一是延时取消,即给定一个最长延时,如果超过这个时间仍然没有个结果,就取消这个任务。其实现方法非常简单,只要把当前任务和一个延时任务丢尽前面实现的wait_any里就可以,当然也可以单独封装成一个函数,示例如下
function with_timeout(f::Function, ts::Float64)
ch = Channel{Any}(1)
@async begin
try
put!(ch, :success => f())
catch e
put!(ch, :error => e)
end
end
@async begin
sleep(ts)
put!(ch, :timeout => "Timed out")
end
status, value = take!(ch)
return status, value
end
status, value = with_timeout(() -> (sleep(5); "Done"), 2.0)
# 返回值 (:timeout, "Timed out")
延时取消,以及前面的wait_any函数,都存在一个很严重的问题,即无法真正取消任务。它只是输出了我们想要的运行结果,但那些未能输出结果的任务,仍在继续工作。
因此,为了真正取消任务,我们必须在任务函数内部实现可被取消的机制,比如通过来自Channel的指令,来结束某个死循环。这种方案的合理性在于,实际使用的异步工作,并不存在大面积的sleep,而是由一个个实际任务组成的,这些任务都可以写成死循环的形式。简单实现如下
function cancellable_worker(stop_ch::Channel{Nothing})
while true
if isready(stop_ch)
take!(stop_ch) # 消费信号,避免堆积
println("收到取消信号,退出…")
break
end
println("Working…")
sleep(1) # 假装在工作
end
end
stop_ch = Channel{Nothing}(1) # 缓冲1,避免发送阻塞
t = @async cancellable_worker(stop_ch)
# 3秒后取消
@async begin
sleep(3)
put!(stop_ch, nothing)
end
wait(t) # 等待任务响应
#=运行结果
Working…
Working…
收到取消信号,退出…
+#
总结
本文介绍了Julia中协程(Task)的基本用法和高级特性。主要内容包括:1) 使用Task创建协程,通过schedule调度和wait等待执行;2) @async宏简化Task创建和调度;3) Channel实现协程间通信,支持阻塞式读写;4) 多协程调度技术,包括@sync同步等待和自定义wait_any函数;5) 任务取消机制,包括超时取消和通过Channel发送取消信号。文章通过大量代码示例展示了Julia协程编程的核心概念和实用技巧,特别强调了任务调度顺序和资源清理的重要性。



