
【仓颉语言入门 · 第28课】并发实战多线程任务处理并发模块最后一课。前两课我们造齐了零件第 26 课的spawn/Future/get负责开线程、等结果第 27 课的Mutex/Condition/ 阻塞队列负责线程之间不打架、能传话。零件散着看不出威力今天把它们组装成一台真正的机器——多线程任务处理机一批任务进来固定数量的工作线程worker并发处理结果按编号回收个别任务失败不拖垮整批生产太快时队列自动削峰。本文所有代码与输出均在仓颉 SDK 1.2.0 下逐行实测编译运行。目录系列导航整套路线共7 个模块、30 课模块课次内容一、环境与入门0105环境搭建与 Hello World、变量与基本类型、运算符与输入输出、分支、循环二、常用类型与数据组织0610字符串、数组与区间、ArrayList/HashMap/HashSet、可空类型、错误处理三、函数与函数式1114函数、Lambda 与高阶函数、闭包、迭代器与惰性序列四、面向对象与类型系统1520struct/class、构造与属性、接口、枚举与 match 模式匹配、泛型、扩展五、工程化与标准库2125cjpm 包管理与多文件、文件 IO、JSON 处理、网络编程、单元测试六、并发编程2628线程的创建与等待、线程同步、并发实战本文七、项目实战2930命令行小工具、GeoJSON 数据处理实战环境搭建与第一个仓颉程序变量、常量与基本数据类型运算符与标准输入输出分支结构与 match 表达式循环结构while / for / Range字符串详解与字符串插值数组 Array 与区间 Range集合框架ArrayList、HashMap、HashSet可空类型?与 Option错误处理异常机制与 Result函数定义、参数与返回值Lambda 与高阶函数闭包、作用域与函数类型迭代器 Iterator 与 Sequence结构体 struct 与类 class构造函数、属性与方法接口 interface 与实现枚举 enum、代数数据类型与 match 模式匹配泛型编程扩展、类型别名与可见性控制cjpm 包管理与多文件项目组织文件与目录 IOJSON 处理网络编程入门单元测试并发基础线程的创建与等待线程同步互斥锁、原子类型与条件变量并发实战多线程任务处理本文实战一带文件持久化的命令行小工具实战二GeoJSON 数据处理程序一、三课并发知识盘点我们手里有什么组装机器之前先清点工具箱来自武器干什么用第 26 课spawn { ... }启动一个新线程立刻返回FutureT第 26 课future.get()等线程结束、拿返回值子线程抛异常时在这里重新抛出第 26 课sleep(Duration.millisecond * n)让线程歇一会模拟慢速 IO第 27 课Mutexsynchronized保护共享数据防止读—改—写交错第 27 课Condition队列空了等、队列满了等不空转烧 CPU第 27 课阻塞队列BlockingQueueT线程之间传任务、传结果的传送带第 18 课枚举 match表达任务/停止成功/失败这类两种结果第 8 课ArrayList攒任务、攒结果今天只引入一个新 API怎么测一段代码跑了多久。然后全部用旧零件搭系统。二、先算一笔时间账并行到底快在哪多线程不是银弹先搞清楚它在什么场景下才划算。现实中慢的任务分两类IO 密集型等网络回复、等磁盘读写、等数据库——CPU 大部分时间在干等CPU 密集型拼命算——CPU 一直满载。sleep模拟的就是第一类线程在睡觉不占 CPU。测时间用DateTime.now()在std.time包里第 22、24 课提过它两个时刻相减得到一个Durationimport 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 ArrayListFutureUnit() 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 次约 800ms8 个线程一起睡约 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 ArrayListFutureInt64() 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 648 号明明最先算完输出仍然从 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 BlockingQueueT { let buf ArrayListT() 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 BlockingQueueCommand(4) let resultQ BlockingQueueResult(16) // 1) 启动 3 个常驻 worker let workers ArrayListFutureUnit() 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 ArrayListResult() 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 耗时不到串行的一半true9 个任务每个 50ms串行要 450ms3 个 worker 分三批干完约 150ms。最后一行的判据 225ms 正是串行的一半本机连跑多次均为true如果你的机器上正赶上高负载出现false把任务数量加大再看数量级即可。4.2 三个关键设计① 毒丸Poison Pill为什么是 3 个worker 是while(true)循环必须有办法让它停下。我们往队列里放一种特殊消息Stopworker 取到它就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 每一个 workerworkers里 3 个Future要挨个get()。这一步有两个作用确认 worker 是正常退出而不是带着异常死掉让主线程在所有线程收尾后再结束程序。五、个别任务失败怎么办异常隔离真实批处理中总会有几个任务出问题文件损坏、网络超时、数据非法。如果 worker 不处理异常异常会沿着Future传出去——可 worker 是常驻循环我们从不 get 它处理单个任务的过程异常会让这个 worker 线程静默死亡剩下的任务全部堆积。正确做法在 worker 内部把每个任务的异常就地接住用枚举把成功/失败都变成正常结果投回主线程。这正是第 10 课异常机制和第 18 课枚举的组合拳import std.sync.* import std.collection.ArrayList import std.sort.* class BlockingQueueT { let buf ArrayListT() 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 BlockingQueueCommand(4) let resultQ BlockingQueueOutcome(16) let workers ArrayListFutureUnit() 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 ArrayListFailure() 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 BlockingQueueT { let buf ArrayListT() 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 BlockingQueueJob(10) let reportQ BlockingQueueReport(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 ArrayListReport() 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 排队超过 100msfalse 任务 2 排队超过 100msfalse 任务 3 排队超过 100msfalse 任务 4 排队超过 100mstrue 任务 5 排队超过 100mstrue 任务 6 排队超过 100mstrue规律连跑 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 BlockingQueueT { let buf ArrayListT() 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 BlockingQueueCommand(4) let resultQ BlockingQueueOutcome(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 ArrayListOutcome() 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 结果按任务编号 任务 1thumb-img-01.png 任务 2thumb-img-02.png 任务 3thumb-img-03.png 任务 4失败格式损坏 任务 5thumb-img-05.png 任务 6thumb-img-06.png 任务 7thumb-img-07.png 任务 8thumb-img-08.png 任务 9失败格式损坏 任务 10thumb-img-10.png 任务 11thumb-img-11.png 任务 12thumb-img-12.png 汇总 成功 10 个失败 2 个 耗时约 239ms829us700ns数字因机器而异除最后一行外每次运行逐字相同。最后一行是本次实测值格式为XmsYusZns你的机器上数字一定不同——10 个成功任务每个 50ms3 个 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 涨到约 500ms10 个成功任务 × 50ms但结果一行不少、顺序不变——这说明 worker 数只影响速度不影响正确性。把workerCount改成 6两批就干完耗时降到约 100ms 档再改成 12 试试收益已经不明显——线程多了也要排队抢 CPU加线程不是越多越快。十、常用 API 速查功能写法备注当前时间DateTime.now()需import std.time.*测耗时DateTime.now() - t1两时刻相减得Duration耗时比较elapsed Duration.millisecond * 225Duration 支持与整数乘法提交一批任务循环futures.add(spawn { ... })一任务一线程模式任务少才用按序回收按 Future 列表顺序get()完成顺序乱回收顺序不乱常驻 workerspawn { 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(...)) }失败也变成结果十一、常见问题 FAQQ1任务完成顺序是乱的怎么保证结果按提交顺序展示两种手法看架构选一任务一线程时按Future列表顺序get()第三节worker 池时给每个任务和结果带编号回收后toArray()sort(by: id)第四节。核心思想一样并行执行串行汇总。Q2毒丸只发一个会怎样只有一个 worker 会取到Stop并退出另外两个永远阻塞在take()等下一条消息。本例主线程不 join 它们时程序也能结束进程退出带走线程但这属于没关干净。记住几个 worker 就发几个毒丸。Q3worker 数设多少合适IO 密集型任务网络、磁盘、sleep 模拟的都是可以多于 CPU 核心数因为线程大部分时间在等CPU 密集型任务设成和 CPU 核心数接近最划算设多了只会增加调度开销。拿不准就从小比如 2、4、8开始压测用第九节的办法观察耗时变化。Q4worker 里不写 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 个FutureInt64收进 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): OptionT { 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枚举取到就处理、返回类型改成spawnInt64统计本 worker 处理了几个任务取到None就break并返回计数。主线程只放 8 个任务、不收结果队列最后把三个 worker 的返回值相加。期望输出共处理 8 个任务所有 worker 已收工。下节预告并发模块到此结业。从下一课开始进入最后的项目实战模块第 29 课实战一带文件持久化的命令行小工具。我们会把前面学的类型系统、集合、文件 IO、异常处理组织成一个看得见、摸得着、关掉再开数据还在的完整小工具并第一次体验多文件项目 菜单主循环 数据落盘的工程结构。系列说明本系列基于 Windows 平台 CIDE 仓颉 SDK1.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