【仓颉语言入门 · 第28课】
【仓颉语言入门 · 第28课】并发实战:多线程任务处理
并发模块最后一课。前两课我们造齐了零件:第 26 课的
spawn/Future/get负责"开线程、等结果",第 27 课的Mutex/Condition/ 阻塞队列负责"线程之间不打架、能传话"。零件散着看不出威力,今天把它们组装成一台真正的机器——多线程任务处理机:一批任务进来,固定数量的工作线程(worker)并发处理,结果按编号回收,个别任务失败不拖垮整批,生产太快时队列自动削峰。本文所有代码与输出均在仓颉 SDK 1.2.0 下逐行实测编译运行。
目录(系列导航)
整套路线共 7 个模块、30 课:
| 模块 | 课次 | 内容 |
|---|---|---|
| 一、环境与入门 | 01~05 | 环境搭建与 Hello World、变量与基本类型、运算符与输入输出、分支、循环 |
| 二、常用类型与数据组织 | 06~10 | 字符串、数组与区间、ArrayList/HashMap/HashSet、可空类型、错误处理 |
| 三、函数与函数式 | 11~14 | 函数、Lambda 与高阶函数、闭包、迭代器与惰性序列 |
| 四、面向对象与类型系统 | 15~20 | struct/class、构造与属性、接口、枚举与 match 模式匹配、泛型、扩展 |
| 五、工程化与标准库 | 21~25 | cjpm 包管理与多文件、文件 IO、JSON 处理、网络编程、单元测试 |
| 六、并发编程 | 26~28 | 线程的创建与等待、线程同步、并发实战(本文) |
| 七、项目实战 | 29~30 | 命令行小工具、GeoJSON 数据处理实战 |
- 环境搭建与第一个仓颉程序
- 变量、常量与基本数据类型
- 运算符与标准输入输出
- 分支结构与 match 表达式
- 循环结构:while / for / Range
- 字符串详解与字符串插值
- 数组 Array 与区间 Range
- 集合框架:ArrayList、HashMap、HashSet
- 可空类型
?与 Option - 错误处理:异常机制与 Result
- 函数定义、参数与返回值
- Lambda 与高阶函数
- 闭包、作用域与函数类型
- 迭代器 Iterator 与 Sequence
- 结构体 struct 与类 class
- 构造函数、属性与方法
- 接口 interface 与实现
- 枚举 enum、代数数据类型与 match 模式匹配
- 泛型编程
- 扩展、类型别名与可见性控制
- cjpm 包管理与多文件项目组织
- 文件与目录 IO
- JSON 处理
- 网络编程入门
- 单元测试
- 并发基础:线程的创建与等待
- 线程同步:互斥锁、原子类型与条件变量
- 并发实战:多线程任务处理(本文)
- 实战一:带文件持久化的命令行小工具
- 实战二:GeoJSON 数据处理程序
一、三课并发知识盘点:我们手里有什么
组装机器之前,先清点工具箱:
| 来自 | 武器 | 干什么用 |
|---|---|---|
| 第 26 课 | spawn { ... } | 启动一个新线程,立刻返回 Future<T> |
| 第 26 课 | future.get() | 等线程结束、拿返回值(子线程抛异常时在这里重新抛出) |
| 第 26 课 | sleep(Duration.millisecond * n) | 让线程歇一会,模拟慢速 IO |
| 第 27 课 | Mutex + synchronized | 保护共享数据,防止"读—改—写"交错 |
| 第 27 课 | Condition | 队列空了等、队列满了等,不空转烧 CPU |
| 第 27 课 | 阻塞队列 BlockingQueue<T> | 线程之间传任务、传结果的传送带 |
| 第 18 课 | 枚举 + match | 表达"任务/停止""成功/失败"这类两种结果 |
| 第 8 课 | ArrayList | 攒任务、攒结果 |
今天只引入一个新 API:怎么测一段代码跑了多久。然后全部用旧零件搭系统。
二、先算一笔时间账:并行到底快在哪
多线程不是银弹,先搞清楚它在什么场景下才划算。现实中"慢"的任务分两类:
- IO 密集型:等网络回复、等磁盘读写、等数据库——CPU 大部分时间在干等;
- CPU 密集型:拼命算——CPU 一直满载。
sleep 模拟的就是第一类(线程在睡觉,不占 CPU)。测时间用 DateTime.now()(在 std.time 包里,第 22、24 课提过它),两个时刻相减得到一个 Duration:
import std.collection.ArrayList
import std.time.*
func work(): Unit {
sleep(Duration.millisecond * 100) // 模拟一个耗时 100ms 的 IO 任务
}
main(): Int64 {
// 串行:8 个任务一个个做
let s1 = DateTime.now()
for (_ in 1..=8) {
work()
}
let serial = DateTime.now() - s1
// 并行:8 个线程同时做
let p1 = DateTime.now()
let futures = ArrayList<Future<Unit>>()
for (_ in 1..=8) {
futures.add(spawn { work() })
}
for (f in futures) {
f.get()
}
let parallel = DateTime.now() - p1
println("并行比串行快:${parallel < serial}")
println("并行不到串行的一半:${parallel * 2 < serial}")
return 0
}
并行比串行快:true
并行不到串行的一半:true
串行要睡 8 次,约 800ms;8 个线程一起睡,约 100ms 出头(多出来的零头是线程创建和调度的开销)。输出里我们不打印具体毫秒数——那是个每次都变的数字,只比较两个 Duration 的大小,结论就是稳定的。
🔸 两点提醒:① 两个
DateTime直接用-相减,得到Duration,它支持<、>比较,也支持乘以整数;② 若是纯 CPU 计算任务,线程数超过 CPU 核心数后不仅不快,还会因切换调度变慢。本课所有例子都是 IO 密集型。
三、模式一:一任务一线程,Future 列表按序回收
最简单的并行模型:来 N 个任务,就开 N 个线程,每个返回一个 Future,把它们按提交顺序收进 ArrayList。
关键问题:任务有快有慢,完成顺序是乱的,结果还能按提交顺序拿到吗? 能。因为 Future 列表本身就是按提交顺序排好的,get() 只是"等",不改变顺序。下面让编号小的任务故意睡更久:
import std.collection.ArrayList
func job(n: Int64): Int64 {
sleep(Duration.millisecond * (10 - n) * 20) // n 越小越慢:1 号睡最久,8 号最先做完
return n * n
}
main(): Int64 {
let futures = ArrayList<Future<Int64>>()
for (n in 1..=8) {
futures.add(spawn { job(n) })
}
// 按列表顺序 get:第一次 get 会一直等到最慢的 1 号完成
for (f in futures) {
print("${f.get()} ")
}
println("")
return 0
}
1 4 9 16 25 36 49 64
8 号明明最先算完,输出仍然从 1 开始。原理:futures[0].get() 死等 1 号,等完了再取 2 号……每个 get 都在自己的位置上等,结果自然归位。第 26 课的分段求和用的就是这个手法,这里再加深一次印象。
这个模式的缺点:线程数量和任务数量一样多。8 个任务无所谓,要是来了 8000 个呢?创建 8000 个线程,内存和调度都扛不住。于是有了模式二。
四、模式二:固定数量 worker + 任务队列
工业界的标准做法是线程池的雏形:
任务队列(BlockingQueue)
┌──┬──┬──┬──┐
│J1│J2│J3│J4│ ←─ 主线程不断往里放任务
└──┴──┴──┴──┘
▲ ▲ ▲
│ │ │ take()
worker worker worker ← 固定 3 个线程,干完一个再取下一个
│ │ │
▼ ▼ ▼
结果队列(BlockingQueue)→ 主线程统一回收
线程数固定(比如 3 个),任务再多也在队列里排队,不会把系统压垮。这里用到第 27 课写的阻塞队列,原样搬过来即可。
4.1 完整代码
import std.sync.*
import std.collection.ArrayList
import std.sort.*
import std.time.*
// ===== 第 27 课的阻塞队列,原样复用 =====
class BlockingQueue<T> {
let buf = ArrayList<T>()
let capacity: Int64
let lock = Mutex()
let notEmpty: Condition
let notFull: Condition
init(capacity: Int64) {
this.capacity = capacity
lock.lock()
this.notEmpty = lock.condition()
this.notFull = lock.condition()
lock.unlock()
}
func put(item: T) {
synchronized (lock) {
while (buf.size >= capacity) {
notFull.wait()
}
buf.add(item)
notEmpty.notifyAll()
}
}
func take(): T {
synchronized (lock) {
while (buf.size == 0) {
notEmpty.wait()
}
let item = buf[0]
buf.remove(0..1)
notFull.notifyAll()
return item
}
}
}
// ===== 队列里放两种消息:干活 / 收工(毒丸) =====
enum Command {
| Job(Int64)
| Stop
}
class Result {
let id: Int64
let value: Int64
init(id: Int64, value: Int64) {
this.id = id
this.value = value
}
}
func process(id: Int64): Int64 {
sleep(Duration.millisecond * 50) // 模拟处理耗时
return id * id
}
main(): Int64 {
let jobQ = BlockingQueue<Command>(4)
let resultQ = BlockingQueue<Result>(16)
// 1) 启动 3 个常驻 worker
let workers = ArrayList<Future<Unit>>()
for (_ in 1..=3) {
workers.add(spawn {
while (true) {
match (jobQ.take()) {
case Job(id) =>
resultQ.put(Result(id, process(id)))
case Stop =>
break
}
}
})
}
// 2) 提交 9 个任务,再提交 3 个"毒丸"
let t1 = DateTime.now()
for (id in 1..=9) {
jobQ.put(Job(id))
}
for (_ in 1..=3) {
jobQ.put(Stop)
}
// 3) 收回 9 个结果
let results = ArrayList<Result>()
for (_ in 1..=9) {
results.add(resultQ.take())
}
let elapsed = DateTime.now() - t1
// 4) 等所有 worker 正常退出
for (w in workers) {
w.get()
}
// 5) 结果到达顺序是乱的,按任务编号排序后再展示
let arr = results.toArray()
sort(arr, by: { a: Result, b: Result => a.id.compare(b.id) })
for (r in arr) {
println("任务 ${r.id} 的结果:${r.value}")
}
println("耗时不到串行的一半:${elapsed < Duration.millisecond * 225}")
return 0
}
任务 1 的结果:1
任务 2 的结果:4
任务 3 的结果:9
任务 4 的结果:16
任务 5 的结果:25
任务 6 的结果:36
任务 7 的结果:49
任务 8 的结果:64
任务 9 的结果:81
耗时不到串行的一半:true
9 个任务每个 50ms,串行要 450ms;3 个 worker 分三批干完,约 150ms。最后一行的判据 225ms 正是串行的一半,本机连跑多次均为 true(如果你的机器上正赶上高负载出现 false,把任务数量加大再看数量级即可)。
4.2 三个关键设计
① 毒丸(Poison Pill)为什么是 3 个? worker 是 while(true) 循环,必须有办法让它停下。我们往队列里放一种特殊消息 Stop,worker 取到它就 break。队列是先进先出的,Stop 又在所有任务之后才放,所以 worker 取到毒丸时,前面的任务保证都处理完了。3 个 worker 各会取走一个 Stop,因此毒丸数量必须恰好等于 worker 数量:少了有 worker 永远卡在 take(),多了则永远没人取(本例主线程不再取任务队列,倒也不报错,但属于逻辑错误)。
② 结果为什么也要走队列? worker 不能直接 println 结果——多线程打印会交错,顺序也乱。把结果投到 resultQ,由主线程单点回收、排序、输出,所有展示顺序都是确定的。
③ 排序为什么要先 toArray()? 第 18 课学过的全局函数 sort(数组, by: { ... }) 只接收 Array,而我们攒结果用的是 ArrayList,所以先 toArray()。比较器返回 Ordering,直接用整数自带的 compare 即可(a.id.compare(b.id))。
4.3 别忘了 get 每一个 worker
workers 里 3 个 Future 要挨个 get()。这一步有两个作用:确认 worker 是正常退出而不是带着异常死掉;让主线程在所有线程收尾后再结束程序。
五、个别任务失败怎么办:异常隔离
真实批处理中,总会有几个任务出问题:文件损坏、网络超时、数据非法。如果 worker 不处理异常,异常会沿着 Future 传出去——可 worker 是常驻循环,我们从不 get 它处理单个任务的过程,异常会让这个 worker 线程静默死亡,剩下的任务全部堆积。
正确做法:在 worker 内部把每个任务的异常就地接住,用枚举把"成功/失败"都变成正常结果投回主线程。这正是第 10 课异常机制和第 18 课枚举的组合拳:
import std.sync.*
import std.collection.ArrayList
import std.sort.*
class BlockingQueue<T> {
let buf = ArrayList<T>()
let capacity: Int64
let lock = Mutex()
let notEmpty: Condition
let notFull: Condition
init(capacity: Int64) {
this.capacity = capacity
lock.lock()
this.notEmpty = lock.condition()
this.notFull = lock.condition()
lock.unlock()
}
func put(item: T) {
synchronized (lock) {
while (buf.size >= capacity) {
notFull.wait()
}
buf.add(item)
notEmpty.notifyAll()
}
}
func take(): T {
synchronized (lock) {
while (buf.size == 0) {
notEmpty.wait()
}
let item = buf[0]
buf.remove(0..1)
notFull.notifyAll()
return item
}
}
}
enum Command {
| Job(Int64)
| Stop
}
enum Outcome {
| Ok(Int64, Int64) // 成功:任务编号、结果
| Fail(Int64, String) // 失败:任务编号、错误信息
}
class Failure {
let id: Int64
let message: String
init(id: Int64, message: String) {
this.id = id
this.message = message
}
}
func process(id: Int64): Int64 {
if (id == 3 || id == 7) {
throw IllegalArgumentException("任务 ${id} 数据损坏")
}
sleep(Duration.millisecond * 20)
return id * id
}
main(): Int64 {
let jobQ = BlockingQueue<Command>(4)
let resultQ = BlockingQueue<Outcome>(16)
let workers = ArrayList<Future<Unit>>()
for (_ in 1..=3) {
workers.add(spawn {
while (true) {
match (jobQ.take()) {
case Job(id) =>
try {
let v = process(id)
resultQ.put(Ok(id, v))
} catch (e: IllegalArgumentException) {
resultQ.put(Fail(id, e.message))
}
case Stop =>
break
}
}
})
}
for (id in 1..=9) {
jobQ.put(Job(id))
}
for (_ in 1..=3) {
jobQ.put(Stop)
}
var okCount = 0
let failures = ArrayList<Failure>()
for (_ in 1..=9) {
match (resultQ.take()) {
case Ok(_, _) =>
okCount += 1
case Fail(id, msg) =>
failures.add(Failure(id, msg))
}
}
for (w in workers) {
w.get()
}
// 失败信息到达顺序不定,收集后按编号排序,输出才稳定
let arr = failures.toArray()
sort(arr, by: { a: Failure, b: Failure => a.id.compare(b.id) })
for (f in arr) {
println("任务 ${f.id} 失败:${f.message}")
}
println("成功 ${okCount} 个,失败 ${arr.size} 个")
return 0
}
任务 3 失败:任务 3 数据损坏
任务 7 失败:任务 7 数据损坏
成功 7 个,失败 2 个
注意三件事:
try-catch的范围只包一个任务。任务 3 炸了,worker 没死,转脸继续处理任务 4,这就是"异常隔离"。- 失败也是一种结果。
Outcome枚举逼主线程必须match两种情况,漏处理哪个分支编译期就会提醒你。 - 展示前再排序。失败结果到达主线程的时刻取决于调度,先收集到
failures里,排序后再打印,连跑 10 次输出都逐字相同。
六、生产快、消费慢:阻塞队列如何削峰填谷
上节预告里承诺的最后一个场景:假设生产任务很快(比如突发涌来的请求),处理却很慢。如果来多少开多少线程,系统瞬间被打爆。有界阻塞队列天然解决这个问题——队列满了,put() 就让生产者在条件变量上等着,生产速度被自动压低到消费速度附近,这就是"削峰填谷"。
下面做个可测量的实验:给每个任务盖上"出生时间戳";3 个 worker 每个任务处理 200ms(很慢),而主线程每 2ms 就生产一个(很快)。任务被取出时,用"当前时间 − 出生时间"算出它在队列里等了多久:
import std.sync.*
import std.collection.ArrayList
import std.sort.*
import std.time.*
class BlockingQueue<T> {
let buf = ArrayList<T>()
let capacity: Int64
let lock = Mutex()
let notEmpty: Condition
let notFull: Condition
init(capacity: Int64) {
this.capacity = capacity
lock.lock()
this.notEmpty = lock.condition()
this.notFull = lock.condition()
lock.unlock()
}
func put(item: T) {
synchronized (lock) {
while (buf.size >= capacity) {
notFull.wait()
}
buf.add(item)
notEmpty.notifyAll()
}
}
func take(): T {
synchronized (lock) {
while (buf.size == 0) {
notEmpty.wait()
}
let item = buf[0]
buf.remove(0..1)
notFull.notifyAll()
return item
}
}
}
class Job {
let id: Int64
let bornAt: DateTime
init(id: Int64, bornAt: DateTime) {
this.id = id
this.bornAt = bornAt
}
}
class Report {
let id: Int64
let waited: Duration
init(id: Int64, waited: Duration) {
this.id = id
this.waited = waited
}
}
main(): Int64 {
let jobQ = BlockingQueue<Job>(10)
let reportQ = BlockingQueue<Report>(10)
// 3 个慢吞吞的消费者
for (_ in 1..=3) {
spawn {
while (true) {
let job = jobQ.take()
let waited = DateTime.now() - job.bornAt
reportQ.put(Report(job.id, waited))
sleep(Duration.millisecond * 200) // 消费很慢
}
}
}
// 生产者很快:每 2ms 一个,6 个任务一眨眼生产完
for (id in 1..=6) {
jobQ.put(Job(id, DateTime.now()))
sleep(Duration.millisecond * 2)
}
let reports = ArrayList<Report>()
for (_ in 1..=6) {
reports.add(reportQ.take())
}
let arr = reports.toArray()
sort(arr, by: { a: Report, b: Report => a.id.compare(b.id) })
for (r in arr) {
println("任务 ${r.id} 排队超过 100ms:${r.waited > Duration.millisecond * 100}")
}
return 0
}
任务 1 排队超过 100ms:false
任务 2 排队超过 100ms:false
任务 3 排队超过 100ms:false
任务 4 排队超过 100ms:true
任务 5 排队超过 100ms:true
任务 6 排队超过 100ms:true
规律连跑 10 次都一样:前 3 个任务到了立刻被 worker 取走,几乎不排队;后 3 个必须等第一轮 200ms 处理完,老老实实在队列里等了近 200ms。任务在队列里"把峰抹平",系统始终只有 3 个线程在忙,内存里最多只有那 6 个任务。
🔸 本例的 worker 是演示用的无限循环,主线程收完报告直接返回,进程退出时线程随之结束。正式程序请照第四节配毒丸,保证 worker 优雅退出。
🔸 队列容量也是调优手段:容量大,抗突发能力强但积压占内存;容量小,生产者更早被"顶住"、任务排队时间更长。要根据内存预算和可接受延迟权衡。
七、综合实战:批量缩略图处理机
把前面所有零件装进一个完整程序:图片站后台批量生成缩略图。12 个图片任务、3 个 worker,其中 2 张图片"格式损坏"必然失败。要求:失败的图片不能影响其他图片;最终结果按任务编号排列;汇总成功/失败数和总耗时。
import std.sync.*
import std.collection.ArrayList
import std.sort.*
import std.time.*
// ---------- 阻塞队列(第 27 课的成果) ----------
class BlockingQueue<T> {
let buf = ArrayList<T>()
let capacity: Int64
let lock = Mutex()
let notEmpty: Condition
let notFull: Condition
init(capacity: Int64) {
this.capacity = capacity
lock.lock()
this.notEmpty = lock.condition()
this.notFull = lock.condition()
lock.unlock()
}
func put(item: T) {
synchronized (lock) {
while (buf.size >= capacity) {
notFull.wait()
}
buf.add(item)
notEmpty.notifyAll()
}
}
func take(): T {
synchronized (lock) {
while (buf.size == 0) {
notEmpty.wait()
}
let item = buf[0]
buf.remove(0..1)
notFull.notifyAll()
return item
}
}
}
// ---------- 任务与结果 ----------
class ImageJob {
let id: Int64
let fileName: String
init(id: Int64, fileName: String) {
this.id = id
this.fileName = fileName
}
}
enum Command {
| Run(ImageJob)
| Stop
}
enum Outcome {
| Done(Int64, String) // 编号、生成的缩略图文件名
| Failed(Int64, String) // 编号、失败原因
}
// ---------- 模拟缩略图处理 ----------
func makeThumbnail(job: ImageJob): String {
if (job.id == 4 || job.id == 9) {
throw IllegalArgumentException("格式损坏")
}
sleep(Duration.millisecond * 50)
return "thumb-" + job.fileName
}
func idOf(o: Outcome): Int64 {
match (o) {
case Done(id, _) => return id
case Failed(id, _) => return id
}
}
// ---------- 主流程 ----------
main(): Int64 {
let workerCount: Int64 = 3
let total: Int64 = 12
println("==== 批处理开始 ====")
println("收到 ${total} 个图片任务,worker 数:${workerCount}")
let jobQ = BlockingQueue<Command>(4)
let resultQ = BlockingQueue<Outcome>(16)
for (_ in 1..=workerCount) {
spawn {
while (true) {
match (jobQ.take()) {
case Run(job) =>
try {
let thumb = makeThumbnail(job)
resultQ.put(Done(job.id, thumb))
} catch (e: IllegalArgumentException) {
resultQ.put(Failed(job.id, e.message))
}
case Stop =>
break
}
}
}
}
let t1 = DateTime.now()
for (id in 1..=total) {
let name = if (id < 10) {
"img-0${id}.png"
} else {
"img-${id}.png"
}
jobQ.put(Run(ImageJob(id, name)))
}
for (_ in 1..=workerCount) {
jobQ.put(Stop)
}
let outcomes = ArrayList<Outcome>()
for (_ in 1..=total) {
outcomes.add(resultQ.take())
}
let elapsed = DateTime.now() - t1
let arr = outcomes.toArray()
sort(arr, by: { a: Outcome, b: Outcome => idOf(a).compare(idOf(b)) })
var okCount = 0
var failCount = 0
println("==== 结果(按任务编号) ====")
for (o in arr) {
match (o) {
case Done(id, thumb) =>
okCount += 1
println("任务 ${id}:${thumb}")
case Failed(id, msg) =>
failCount += 1
println("任务 ${id}:失败(${msg})")
}
}
println("==== 汇总 ====")
println("成功 ${okCount} 个,失败 ${failCount} 个")
println("耗时约 ${elapsed}(数字因机器而异)")
return 0
}
==== 批处理开始 ====
收到 12 个图片任务,worker 数:3
==== 结果(按任务编号) ====
任务 1:thumb-img-01.png
任务 2:thumb-img-02.png
任务 3:thumb-img-03.png
任务 4:失败(格式损坏)
任务 5:thumb-img-05.png
任务 6:thumb-img-06.png
任务 7:thumb-img-07.png
任务 8:thumb-img-08.png
任务 9:失败(格式损坏)
任务 10:thumb-img-10.png
任务 11:thumb-img-11.png
任务 12:thumb-img-12.png
==== 汇总 ====
成功 10 个,失败 2 个
耗时约 239ms829us700ns(数字因机器而异)
除最后一行外,每次运行逐字相同。最后一行是本次实测值,格式为 XmsYusZns,你的机器上数字一定不同——10 个成功任务每个 50ms,3 个 worker 分四批,理论约 200ms,实测加上调度开销约 240ms。
这个程序的骨架可以直接迁移到任何批处理场景:把 ImageJob 换成 URL 就是多线程抓取器,换成文件路径就是批量转换器,换成 SQL 就是并发导入器——并发结构不用动,只换 makeThumbnail 这一个函数。
八、并发设计军规
把三节课的血泪经验浓缩成七条:
- 共享可变数据必须装箱加锁。
var不能被spawn捕获;计数器、多字段状态装进class,字段访问走synchronized(或用原子类型)。 - 能用原子类型就不用锁。单个计数器用
AtomicInt64;多个字段必须"一起变"时才用Mutex(第 27 课练习 1 的坐标点)。 - 等待条件用
Condition,不要 while 自旋。队列空/满时wait(),被通知后在while里复查。 - worker 数固定,任务走队列。别一任务一线程;worker 数参考任务类型(IO 密集可多于 CPU 核心数,CPU 密集约等于核心数)。
- 毒丸数 = worker 数,且最后才放。
- 异常在 worker 内部就地接住,用枚举把失败也变成结果,否则失败的 worker 会悄悄退场,任务越积越多。
- 输出和汇总只在一个线程做。worker 只投递结果,主线程单点回收、排序、打印——这是输出结果可复现的根本保证。
九、CIDE 实操:改两个参数,感受线程池
在 CIDE 里 cjpm init --name imgbatch,把第七节完整代码贴进 src/main.cj,连跑三次:前 17 行三次逐字一致,只有耗时行在 200ms 档波动。
然后做两个对照实验(记得改回来):
- 把
workerCount改成 1:退化成串行,耗时从约 240ms 涨到约 500ms(10 个成功任务 × 50ms),但结果一行不少、顺序不变——这说明 worker 数只影响速度,不影响正确性。 - 把
workerCount改成 6:两批就干完,耗时降到约 100ms 档;再改成 12 试试,收益已经不明显——线程多了也要排队抢 CPU,加线程不是越多越快。
十、常用 API 速查
| 功能 | 写法 | 备注 |
|---|---|---|
| 当前时间 | DateTime.now() | 需 import std.time.* |
| 测耗时 | DateTime.now() - t1 | 两时刻相减得 Duration |
| 耗时比较 | elapsed < Duration.millisecond * 225 | Duration 支持 < > 与整数乘法 |
| 提交一批任务 | 循环 futures.add(spawn { ... }) | 一任务一线程模式,任务少才用 |
| 按序回收 | 按 Future 列表顺序 get() | 完成顺序乱,回收顺序不乱 |
| 常驻 worker | spawn { while (true) { match (q.take()) {...} } } | 配合毒丸退出 |
| 优雅停止 | 每 worker 发一个 Stop | 毒丸数 = worker 数,最后放 |
| 结果回传 | 结果队列 resultQ.put(...) | 主线程单点回收 |
| 列表转数组 | arrayList.toArray() | sort 只收 Array |
| 按字段排序 | sort(arr, by: { a, b => a.id.compare(b.id) }) | 比较器返回 Ordering,见第 18 课 |
| 异常隔离 | worker 内 try { ... } catch (e) { resultQ.put(Fail(...)) } | 失败也变成结果 |
十一、常见问题 FAQ
Q1:任务完成顺序是乱的,怎么保证结果按提交顺序展示?
两种手法,看架构选:一任务一线程时,按 Future 列表顺序 get()(第三节);worker 池时,给每个任务和结果带编号,回收后 toArray() + sort(by: id)(第四节)。核心思想一样:并行执行,串行汇总。
Q2:毒丸只发一个会怎样?
只有一个 worker 会取到 Stop 并退出,另外两个永远阻塞在 take() 等下一条消息。本例主线程不 join 它们时程序也能结束(进程退出带走线程),但这属于"没关干净"。记住:几个 worker 就发几个毒丸。
Q3:worker 数设多少合适?
IO 密集型任务(网络、磁盘、sleep 模拟的都是)可以多于 CPU 核心数,因为线程大部分时间在等;CPU 密集型任务设成和 CPU 核心数接近最划算,设多了只会增加调度开销。拿不准就从小(比如 2、4、8)开始压测,用第九节的办法观察耗时变化。
Q4:worker 里不写 try-catch 会怎样?
任务抛出的异常会让该次 spawn 对应的 Future 承载异常。但常驻 worker 的 Future 只有最后退出时才 get(),中途任务的异常无人接收,worker 线程会静默终止,表现为"处理速度莫名其妙变慢、任务堆积"。务必像第五节那样把 try-catch 放在循环内部、每个任务一包。
Q5:为什么 worker 里不能直接 var count = 0 统计完成数?
第 26、27 课的老规矩:spawn 闭包不能捕获可变变量(报 error: 'spawn' expressions cannot capture mutable variables; consider using 'let' or boxing)。把计数装进 class 用锁保护,或用 AtomicInt64.fetchAdd(1)。
Q6:一任务一线程和 worker 池怎么选?
任务数量少、生命周期明确(比如并行算几个分片、同时请求几个固定接口)用一任务一线程,get() 回收最简单;任务数量大、来源持续(突发请求、批量文件、消息流)用固定 worker + 有界队列,靠队列削峰、靠毒丸收尾。
十二、课后练习
- (必做)用第三节"一任务一线程"模式写一个并行平方批处理:6 个任务,每个先
sleep(Duration.millisecond * 100)再返回编号的平方。把 6 个Future<Int64>收进 ArrayList,按提交顺序get()并逐行打印。期望严格输出下面 6 行(不要排序,体会"按列表顺序 get 天然有序"):
任务 1 => 1
任务 2 => 4
任务 3 => 9
任务 4 => 16
任务 5 => 25
任务 6 => 36
-
(必做)把第四节的 worker 池动手补全并改参数运行:8 个任务、2 个 worker,任务内容改成"编号翻倍"(
return id * 2),每个任务sleep(Duration.millisecond * 30)。需要你自己写的三处:① 创建resultQ;② worker 主循环里match (jobQ.take())的两个分支;③ 毒丸提交(想清楚要发几个)。期望输出任务 1 的结果:2到任务 8 的结果:16共 8 行。 -
(必做)在第 2 题基础上加异常隔离:约定奇数编号的任务抛
IllegalArgumentException("奇数任务失败"),偶数任务正常返回编号的 10 倍。用Outcome枚举回收结果,最后按编号升序打印每个失败任务的编号,并输出统计。6 个任务、2 个 worker 时期望输出:
任务 1 失败:奇数任务失败
任务 3 失败:奇数任务失败
任务 5 失败:奇数任务失败
成功 3 个,失败 3 个
- (选做)给
BlockingQueue增加一个限时取任务方法,让 worker 空闲超时自动收工,从而连毒丸都不用发:
// 提示:wait(timeout:) 被通知返回 true、超时返回 false(第 27 课 6.3 节)
func poll(timeout: Duration): Option<T> {
synchronized (lock) {
while (buf.size == 0) {
if (!notEmpty.wait(timeout: timeout)) {
return None // 等满了还没任务
}
}
let item = buf[0]
buf.remove(0..1)
notFull.notifyAll()
return Some(item)
}
}
3 个 worker 用 poll(Duration.millisecond * 200) 取任务(队列里直接放 ImageJob,不再需要 Command 枚举),取到就处理、返回类型改成 spawn<Int64> 统计本 worker 处理了几个任务,取到 None 就 break 并返回计数。主线程只放 8 个任务、不收结果队列,最后把三个 worker 的返回值相加。期望输出:共处理 8 个任务,所有 worker 已收工。
下节预告
并发模块到此结业。从下一课开始进入最后的项目实战模块:第 29 课 实战一:带文件持久化的命令行小工具。我们会把前面学的类型系统、集合、文件 IO、异常处理组织成一个看得见、摸得着、关掉再开数据还在的完整小工具,并第一次体验"多文件项目 + 菜单主循环 + 数据落盘"的工程结构。
系列说明:本系列基于 Windows 平台 + CIDE + 仓颉 SDK(1.2.0)编写,所有代码均已实际编译运行通过。如遇 SDK 版本差异导致的细节出入,以你本地版本为准,欢迎评论区交流。
💬 遇到问题?扫码联系作者
跟着课程练习时,如果在 SDK 安装、环境变量配置、编译报错或调试上卡住,欢迎扫码加作者企业微信直接咨询(请备注"仓颉课程"):

离线环境下图片可能加载不出来,也可以在 CIDE 菜单 Help ▸ 联系作者 / Contact 中查看同一张二维码(应用内置兜底图,无需联网)。
📥 工具下载
本系列全程使用的仓颉 IDE —— CIDE(免费开源、社区版):
- GitCode 仓库 / 安装包下载:https://gitcode.com/wp_upala/cide
- 打开页面后进入 发行版(Releases),两种包任选其一:
- 安装版:下载
CIDE-<版本>-x64-Setup.exe,双击安装,适合日常长期使用; - 免安装版(Portable):下载
CIDE-<版本>-x64-Portable.zip,解压到任意目录即用,不写注册表、不留安装痕迹,拷到 U 盘也能在别的电脑直接运行(包内附《使用说明.txt》)。适合先试用、或在受限电脑上学习本系列课程。
- 安装版:下载
- 仓颉 SDK 请前往仓颉编程语言官网下载:https://cangjie-lang.cn
更多推荐




所有评论(0)