文件协程与管道
十一、文件操作
十二、协程和管道
基本概念
- 并发(Concurrency):并发是指程序中多个操作可以同时进行。在Go中,goroutine使得并发执行变得简单。
- 并行(Parallelism):并行是指程序中多个操作同时在多个CPU核心上执行。并发不一定意味着并行,但并行是并发的一种形式。
- goroutine:goroutine是Go语言中的并发执行单元,它比传统的线程更轻量级,开销更小。
1、协程
Goroutine 是 Go 语言中的一种轻量级线程,由 Go 运行时管理。每个 Goroutine 都由 Go 运行时的调度器(scheduler)调度执行,而不像操作系统的线程一样由操作系统直接调度。Goroutine 允许我们在程序中并发地执行多个任务。
1.1、协程的创建
package main
import (
"fmt"
"time"
)
func sayHello() {
for i := 0; i < 5; i++ {
time.Sleep(100 * time.Millisecond)
fmt.Println("Hello!")
}
}
func main() {
go sayHello() // 启动一个协程
for i := 0; i < 5; i++ {
time.Sleep(150 * time.Millisecond)
fmt.Println("Hi!")
}
}1.2、Goroutine 的特性
- 轻量性:相比于操作系统线程,Goroutine 更加轻量。每个 Goroutine 的栈空间非常小,通常只有 2KB 左右,而操作系统线程的栈通常需要几 MB 的内存。Goroutine 会在运行时动态扩展和收缩栈空间,因此可以更高效地使用内存。
- 并发性:Goroutine 通过 Go 运行时的调度器进行调度,多个 Goroutine 可能被映射到少数几个操作系统线程上执行。调度器负责将 Goroutine 分配到可用的 CPU 核心上执行,实现高效的并发。
- 非阻塞执行:当启动一个 Goroutine 后,主程序会继续执行,而不会等待该 Goroutine 执行完毕。换句话说,
go关键字会让指定的函数以并发的方式执行,不会阻塞主程序。 - 调度机制:Go 运行时的调度器使用的是 M:N 调度模型(即多个 Goroutine 运行在多个操作系统线程上)。Go 调度器通过抢占式调度和协作式调度来管理 Goroutine 的执行。这个调度模型是 Go 并发模型的核心。
示例:启动多个 Goroutine
通过 go 关键字,我们可以启动多个 Goroutine 来并发执行多个任务。以下是一个例子:
package main
import "fmt"
func task(id int) {
fmt.Printf("Task %d is running\n", id)
}
func main() {
// 启动多个 Goroutine
for i := 1; i <= 3; i++ {
go task(i)
}
// 为了让 Goroutine 执行完毕,主线程等待
fmt.Scanln() // 阻塞主线程,等待输入
}在这个例子中,main 函数启动了三个 Goroutine,每个 Goroutine 执行 task 函数。fmt.Scanln() 用来阻塞主线程,确保所有的 Goroutine 都有机会执行。
为什么 fmt.Scanln() 是必要的?
fmt.Scanln() 用来阻塞主程序,保证 Goroutine 有足够的时间执行。在实际应用中,Goroutine 通常会与其他同步机制配合使用(例如,使用 Channel 或 WaitGroup)来等待所有任务完成。
1.3、 协程的调度机制
M:N调度模型
Go运行时采用M:N调度模型,其中:
- M 表示操作系统线程。
- N 表示协程(goroutine)。
这种模型允许多个协程在少数几个操作系统线程上多路复用,从而提高资源利用率和并发性能。当某个协程阻塞时,调度器会将其挂起,并将其他就绪的协程调度到当前线程上执行,从而避免线程阻塞。
调度器的工作原理
Go的调度器负责将协程分配到可用的操作系统线程上执行。调度器会根据协程的状态(如就绪、运行、阻塞等)进行动态调度。当某个协程阻塞时(例如等待I/O操作),调度器会将其挂起,并将其他就绪的协程调度到当前线程上执行。
1.4、主死从随
在 Go 中,协程(goroutine)是非常轻量的线程,主程序中的主协程并不会因为启动了其他协程而阻塞。主协程结束后,程序就会退出,无论其他协程是否执行完毕。因此,通常我们需要确保所有协程都执行完成后,主程序才能结束。
主协程死亡后行为的随机性
当主协程提前结束时,其他协程的行为会有一定的 随机性,主要体现在以下几个方面:
- 主协程提前退出:如果主协程(通常是
main函数)没有等到其他协程完成就退出,Go 程序就会结束,所有其他正在运行的协程会被强制终止。 - 依赖调度器的随机性:Go 调度器(runtime scheduler)会调度并运行所有的协程,但调度顺序和执行时长并不能完全预知。主协程结束后,由于协程调度是并发进行的,可能会在不同的时间点终止,造成某些协程可能还没有完全执行完就被中止。
- 内存和资源的随机释放:如果主协程过早结束,它也可能会导致程序退出时对资源的释放不完全。具体来说,由于内存和资源的管理可能依赖于某些延迟释放机制(如垃圾回收、文件关闭等),所以如果主程序过早退出,有些资源可能还没有被正确释放,造成某些协程的执行受到影响。
示例:主协程死亡导致子协程随机执行
考虑以下代码:
package main
import (
"fmt"
"time"
)
func task(id int) {
fmt.Printf("Task %d started\n", id)
time.Sleep(time.Second * 2)
fmt.Printf("Task %d finished\n", id)
}
func main() {
// 启动多个 Goroutine
go task(1)
go task(2)
go task(3)
// 主协程提前退出
fmt.Println("Main goroutine is exiting!")
}运行结果(每次可能不同):
Main goroutine is exiting!
Task 1 started
Task 3 started
Task 2 started或者:
Main goroutine is exiting!
Task 1 started
Task 2 started
Task 3 started
Task 1 finished
Task 2 finished
Task 3 finished分析:
- 随机性:在这个例子中,由于
main协程提前退出,其他协程的执行顺序是不确定的,而且有时它们甚至没有完全执行完。Go 语言的调度器会在main协程退出后随机地调度正在执行的其他协程,这种行为可能在每次运行时表现出不同的结果。 - 提前退出:在没有同步机制(如
WaitGroup或Channel)的情况下,主协程提前退出时,其他协程仍然可以继续执行,但如果主协程在其还没有完成之前就退出了,程序会强制退出,导致所有正在运行的协程被中止。
解决方案:使用 sync.WaitGroup
为了确保主协程在所有子协程执行完毕后退出,我们可以使用 sync.WaitGroup 来同步所有 Goroutine 的执行。
package main
import (
"fmt"
"sync"
"time"
)
func task(id int, wg *sync.WaitGroup) {
defer wg.Done() // 在函数退出时调用 Done 来通知 WaitGroup
fmt.Printf("Task %d started\n", id)
time.Sleep(time.Second * 2)
fmt.Printf("Task %d finished\n", id)
}
func main() {
var wg sync.WaitGroup
// 启动多个 Goroutine
for i := 1; i <= 3; i++ {
wg.Add(1) // 每启动一个协程,增加一个计数
go task(i, &wg)
}
// 等待所有 Goroutine 完成
wg.Wait()
fmt.Println("Main goroutine is exiting!")
}运行结果(确保子协程都执行完毕):
Task 1 started
Task 2 started
Task 3 started
Task 1 finished
Task 2 finished
Task 3 finished
Main goroutine is exiting!分析:
- 在这个改进后的版本中,
sync.WaitGroup保证了main函数在所有协程执行完毕之前不会退出。main函数会在调用wg.Wait()时阻塞,直到所有 Goroutine 调用wg.Done()表示它们已经执行完成。 - 这样,无论主协程何时退出,所有的子协程都能正确地执行完毕。
1.5、锁
1.5.1、互斥锁
互斥锁是最常见的锁机制之一,通常用来保护共享资源,防止多个协程同时访问这些资源,造成数据竞争或不一致的问题。
sync.Mutex
- 作用:
sync.Mutex提供了一个基本的锁机制,通过Lock和Unlock方法来显式地加锁和解锁。 - 应用场景:当多个协程需要对共享资源进行修改时,可以使用互斥锁来保证每次只有一个协程能够访问该资源。
示例:使用 sync.Mutex 保护共享资源
package main
import (
"fmt"
"sync"
)
var (
counter int
mutex sync.Mutex
)
func increment(wg *sync.WaitGroup) {
defer wg.Done()
mutex.Lock() // 获取锁
counter++
mutex.Unlock() // 释放锁
}
func main() {
var wg sync.WaitGroup
for i := 0; i < 1000; i++ {
wg.Add(1)
go increment(&wg)
}
wg.Wait() // 等待所有 goroutine 执行完
fmt.Println("Final counter:", counter)
}解释:
mutex.Lock()和mutex.Unlock()用来保护对counter变量的访问,确保每次只有一个协程能修改它。sync.WaitGroup用来等待所有的协程完成。
1.5.2、读写锁
读写锁允许多个读操作并发执行,但写操作必须是独占的。这意味着,如果有一个写操作,所有的读操作都会被阻塞,直到写操作完成。
sync.RWMutex
- 作用:
sync.RWMutex提供了RLock(读锁)和Lock(写锁)两种锁操作。RLock:允许多个协程同时读取共享资源,但写操作时会阻塞。Lock:写锁会阻塞所有读操作和其他写操作。
- 应用场景:适用于读多写少的场景,例如缓存、配置管理等。
示例:使用 sync.RWMutex
package main
import (
"fmt"
"sync"
)
var (
data int
rwMutex sync.RWMutex
)
func read(wg *sync.WaitGroup) {
defer wg.Done()
rwMutex.RLock() // 获取读锁
fmt.Println("Reading:", data)
rwMutex.RUnlock() // 释放读锁
}
func write(wg *sync.WaitGroup, value int) {
defer wg.Done()
rwMutex.Lock() // 获取写锁
data = value
fmt.Println("Writing:", value)
rwMutex.Unlock() // 释放写锁
}
func main() {
var wg sync.WaitGroup
for i := 0; i < 5; i++ {
wg.Add(1)
go read(&wg)
}
for i := 0; i < 5; i++ {
wg.Add(1)
go write(&wg, i)
}
wg.Wait()
}解释:
rwMutex.RLock()用于获取读锁,允许多个协程同时读取数据。rwMutex.Lock()用于获取写锁,写操作会阻塞所有读操作和其他写操作。sync.WaitGroup用来等待所有的协程完成。
1.5.3、原子操作(Atomic Operations)
Go 提供了 sync/atomic 包,允许我们对整数值等进行原子操作,而无需使用互斥锁。这对于性能要求较高的场景非常有用,因为原子操作是无锁的。
示例:使用原子操作增加计数
package main
import (
"fmt"
"sync"
"sync/atomic"
)
var counter int64
func increment(wg *sync.WaitGroup) {
defer wg.Done()
atomic.AddInt64(&counter, 1) // 原子增加
}
func main() {
var wg sync.WaitGroup
for i := 0; i < 1000; i++ {
wg.Add(1)
go increment(&wg)
}
wg.Wait()
fmt.Println("Final counter:", counter)
}解释:
atomic.AddInt64是一个原子操作,它直接对counter进行加1操作,避免了使用锁的开销。sync.WaitGroup用来等待所有的协程完成。
1.5.3、 锁的竞争与死锁
在 Go 中,锁竞争(Lock Contention)和死锁(Deadlock)是并发编程中常见的问题。
- 锁竞争:多个协程争夺同一个锁时,可能会影响程序性能。如果多个协程频繁获取和释放锁,可能会导致上下文切换和 CPU 资源浪费。
- 死锁:多个协程在等待彼此释放锁时,程序就会进入死锁状态,导致所有协程无法继续执行。死锁通常是由于锁获取顺序不一致或者锁的使用不当引起的。
示例:死锁示例
package main
import "sync"
var mutex1, mutex2 sync.Mutex
func deadlock() {
mutex1.Lock()
mutex2.Lock() // 死锁:两个锁获取顺序相反
}
func main() {
go deadlock()
mutex1.Lock()
mutex2.Lock()
}解决死锁的方法:
- 避免嵌套锁:尽量减少多个锁的嵌套,可以考虑使用
sync.RWMutex或其他更高效的同步机制。 - 锁的获取顺序:确保所有协程获取锁的顺序一致,避免循环等待。
2、管道
2.1、基础
Go 语言提供了 Channel 来在 Goroutine 之间传递数据,从而实现 Goroutine 之间的通信和同步。Channel 是 Go 并发编程的核心工具。
创建 Channel
Channel 是通过 make 函数创建的,指定 Channel 的类型和容量:
ch := make(chan int) // 创建一个 int 类型的 Channel向 Channel 发送数据
使用 <- 运算符向 Channel 发送数据:
ch <- 42 // 向 Channel 发送数据 42从 Channel 接收数据
使用 <- 运算符从 Channel 接收数据:
value := <-ch // 从 Channel 接收数据使用 Channel 进行同步
Channel 不仅用于传输数据,也可以用于 Goroutine 之间的同步。例如,我们可以使用 Channel 等待多个 Goroutine 执行完毕:
package main
import "fmt"
func task(ch chan bool, id int) {
fmt.Printf("Task %d is running\n", id)
ch <- true // 向 Channel 发送信号,表示任务完成
}
func main() {
ch := make(chan bool)
// 启动多个 Goroutine
for i := 1; i <= 3; i++ {
go task(ch, i)
}
// 等待所有 Goroutine 完成
for i := 1; i <= 3; i++ {
<-ch // 从 Channel 接收信号,等待任务完成
}
fmt.Println("All tasks completed")
}输出:
Task 1 is running
Task 2 is running
Task 3 is running
All tasks completedChannel 的关闭
当我们不再需要向 Channel 发送数据时,可以使用 close() 来关闭 Channel:
close(ch)关闭的 Channel 不能再发送数据,但可以继续从中接收数据。通常,关闭 Channel 用于告知接收方所有数据已经发送完毕。
2.2、遍历
Go 语言中的 range 可以用来遍历一个管道中的数据,直到管道关闭。遍历操作会从管道中依次接收元素,直到管道被关闭并且没有更多的数据可以接收为止。
示例:遍历管道中的数据
package main
import "fmt"
func sendData(ch chan int) {
for i := 0; i < 5; i++ {
ch <- i // 将数据发送到管道
}
close(ch) // 发送数据完成后关闭管道
}
func main() {
ch := make(chan int)
// 启动协程发送数据
go sendData(ch)
// 遍历管道中的数据
for data := range ch {
fmt.Println(data) // 打印从管道接收到的数据
}
}解释:
range ch用来遍历ch管道中的数据。- 当管道关闭且所有数据被接收时,
range循环结束。 - 如果没有关闭管道且管道没有数据可接收,
range会阻塞。
2.3、无缓冲的channel
无缓冲的管道(unbuffered channel)是最基本的管道类型,它的特性是:
- 发送操作 会被阻塞,直到有一个接收操作准备好接收数据。
- 接收操作 会被阻塞,直到有一个发送操作向管道中发送数据。
这种类型的管道用于需要严格同步的场景,因为发送者和接收者必须保持同步。
示例:无缓冲管道
package main
import "fmt"
func sendData(ch chan int) {
fmt.Println("Sending data...")
ch <- 1 // 发送数据到管道,阻塞直到有接收者接收
fmt.Println("Data sent.")
}
func main() {
ch := make(chan int) // 无缓冲管道
go sendData(ch)
fmt.Println("Receiving data...")
data := <-ch // 从管道接收数据,阻塞直到数据到达
fmt.Println("Data received:", data)
}解释:
- 当
sendData协程中的ch <- 1执行时,程序会在ch <- 1处阻塞,直到主协程执行data := <-ch接收数据。
2.4、有缓冲的channel
有缓冲的管道(buffered channel)允许在管道中存储一定数量的数据,发送操作不会立即阻塞,直到缓冲区满时才会阻塞。接收操作也类似,只有在管道为空时才会阻塞。
有缓冲的管道非常适合于解耦生产者和消费者的速度差异。
示例:有缓冲管道
package main
import "fmt"
func sendData(ch chan int) {
for i := 0; i < 5; i++ {
ch <- i // 向管道发送数据,只有缓冲区未满时才会发送
fmt.Println("Sent:", i)
}
close(ch) // 完成数据发送后关闭管道
}
func main() {
ch := make(chan int, 3) // 有缓冲的管道,缓冲区大小为3
go sendData(ch)
// 接收数据
for data := range ch {
fmt.Println("Received:", data)
}
}解释:
- 管道
ch := make(chan int, 3)创建了一个缓冲区大小为 3 的管道。 - 在
sendData中,数据发送到管道,直到缓冲区满为止才会阻塞。 - 缓冲区满时,
sendData会阻塞,直到有接收者从管道中取出数据。
2.5、channel与range
在 Go 中,我们可以使用 range 来遍历管道中的数据。与普通的 range 遍历数组、切片或映射不同,遍历管道会持续接收数据,直到管道关闭。
示例:使用 range 遍历管道
package main
import "fmt"
func generateNumbers(ch chan int) {
for i := 0; i < 5; i++ {
ch <- i
}
close(ch) // 完成数据发送,关闭管道
}
func main() {
ch := make(chan int)
go generateNumbers(ch)
// 使用 range 遍历管道
for num := range ch {
fmt.Println(num) // 接收并打印数据
}
}解释:
range ch会不断从管道中接收数据,直到管道关闭并且所有数据都被接收。close(ch)用来关闭管道,表示没有更多数据会被发送。
2.6、channel与select
select 语句用于在多个管道操作之间进行选择,它类似于 switch 语句,能够同时监听多个管道的发送和接收操作,并且能够处理多个通道的阻塞。
示例:使用 select 同时处理多个管道
package main
import "fmt"
func sendData(ch chan int) {
for i := 0; i < 5; i++ {
ch <- i
}
close(ch)
}
func main() {
ch1 := make(chan int)
ch2 := make(chan int)
go sendData(ch1)
go sendData(ch2)
// 使用 select 监听多个管道
for i := 0; i < 10; i++ {
select {
case data := <-ch1:
fmt.Println("Received from ch1:", data)
case data := <-ch2:
fmt.Println("Received from ch2:", data)
}
}
}解释:
select语句会阻塞,直到它能够从其中一个管道接收到数据。- 在
select中,如果多个管道都可以操作,Go 会随机选择一个管道进行处理。 - 如果所有管道都没有准备好操作,
select会一直阻塞直到至少一个管道准备就绪。
示例:超时机制与 select
package main
import (
"fmt"
"time"
)
func sendData(ch chan int) {
time.Sleep(2 * time.Second)
ch <- 1
}
func main() {
ch := make(chan int)
go sendData(ch)
// 使用 select 设置超时
select {
case data := <-ch:
fmt.Println("Received data:", data)
case <-time.After(1 * time.Second): // 设置 1 秒超时
fmt.Println("Timeout!")
}
}解释:
time.After(1 * time.Second)会返回一个在 1 秒后发送数据的管道,因此如果管道在 1 秒内没有数据可接收,select会执行超时分支。
