【仓颉语言入门 · 第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~20struct/class、构造与属性、接口、枚举与 match 模式匹配、泛型、扩展
五、工程化与标准库21~25cjpm 包管理与多文件、文件 IO、JSON 处理、网络编程、单元测试
六、并发编程26~28线程的创建与等待、线程同步、并发实战(本文)
七、项目实战29~30命令行小工具、GeoJSON 数据处理实战
  1. 环境搭建与第一个仓颉程序
  2. 变量、常量与基本数据类型
  3. 运算符与标准输入输出
  4. 分支结构与 match 表达式
  5. 循环结构:while / for / Range
  6. 字符串详解与字符串插值
  7. 数组 Array 与区间 Range
  8. 集合框架:ArrayList、HashMap、HashSet
  9. 可空类型 ? 与 Option
  10. 错误处理:异常机制与 Result
  11. 函数定义、参数与返回值
  12. Lambda 与高阶函数
  13. 闭包、作用域与函数类型
  14. 迭代器 Iterator 与 Sequence
  15. 结构体 struct 与类 class
  16. 构造函数、属性与方法
  17. 接口 interface 与实现
  18. 枚举 enum、代数数据类型与 match 模式匹配
  19. 泛型编程
  20. 扩展、类型别名与可见性控制
  21. cjpm 包管理与多文件项目组织
  22. 文件与目录 IO
  23. JSON 处理
  24. 网络编程入门
  25. 单元测试
  26. 并发基础:线程的创建与等待
  27. 线程同步:互斥锁、原子类型与条件变量
  28. 并发实战:多线程任务处理(本文)
  29. 实战一:带文件持久化的命令行小工具
  30. 实战二: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 个

注意三件事:

  1. try-catch 的范围只包一个任务。任务 3 炸了,worker 没死,转脸继续处理任务 4,这就是"异常隔离"。
  2. 失败也是一种结果。Outcome 枚举逼主线程必须 match 两种情况,漏处理哪个分支编译期就会提醒你。
  3. 展示前再排序。失败结果到达主线程的时刻取决于调度,先收集到 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 这一个函数。


八、并发设计军规

把三节课的血泪经验浓缩成七条:

  1. 共享可变数据必须装箱加锁。var 不能被 spawn 捕获;计数器、多字段状态装进 class,字段访问走 synchronized(或用原子类型)。
  2. 能用原子类型就不用锁。单个计数器用 AtomicInt64;多个字段必须"一起变"时才用 Mutex(第 27 课练习 1 的坐标点)。
  3. 等待条件用 Condition,不要 while 自旋。队列空/满时 wait(),被通知后在 while 里复查。
  4. worker 数固定,任务走队列。别一任务一线程;worker 数参考任务类型(IO 密集可多于 CPU 核心数,CPU 密集约等于核心数)。
  5. 毒丸数 = worker 数,且最后才放。
  6. 异常在 worker 内部就地接住,用枚举把失败也变成结果,否则失败的 worker 会悄悄退场,任务越积越多。
  7. 输出和汇总只在一个线程做。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 * 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(...)) }失败也变成结果

十一、常见问题 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 + 有界队列,靠队列削峰、靠毒丸收尾。


十二、课后练习

  1. (必做)用第三节"一任务一线程"模式写一个并行平方批处理:6 个任务,每个先 sleep(Duration.millisecond * 100) 再返回编号的平方。把 6 个 Future<Int64> 收进 ArrayList,按提交顺序 get() 并逐行打印。期望严格输出下面 6 行(不要排序,体会"按列表顺序 get 天然有序"):
任务 1 => 1
任务 2 => 4
任务 3 => 9
任务 4 => 16
任务 5 => 25
任务 6 => 36
  1. (必做)把第四节的 worker 池动手补全并改参数运行:8 个任务、2 个 worker,任务内容改成"编号翻倍"(return id * 2),每个任务 sleep(Duration.millisecond * 30)。需要你自己写的三处:① 创建 resultQ;② worker 主循环里 match (jobQ.take()) 的两个分支;③ 毒丸提交(想清楚要发几个)。期望输出 任务 1 的结果:2 到 任务 8 的结果:16 共 8 行。

  2. (必做)在第 2 题基础上加异常隔离:约定奇数编号的任务抛 IllegalArgumentException("奇数任务失败"),偶数任务正常返回编号的 10 倍。用 Outcome 枚举回收结果,最后按编号升序打印每个失败任务的编号,并输出统计。6 个任务、2 个 worker 时期望输出:

任务 1 失败:奇数任务失败
任务 3 失败:奇数任务失败
任务 5 失败:奇数任务失败
成功 3 个,失败 3 个
  1. (选做)给 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
Logo

一站式 AI 云服务平台

更多推荐