Skip to content

文件协程与管道

十一、文件操作

十二、协程和管道

基本概念

  • 并发(Concurrency):并发是指程序中多个操作可以同时进行。在Go中,goroutine使得并发执行变得简单。
  • 并行(Parallelism):并行是指程序中多个操作同时在多个CPU核心上执行。并发不一定意味着并行,但并行是并发的一种形式。
  • goroutine:goroutine是Go语言中的并发执行单元,它比传统的线程更轻量级,开销更小。

1、协程

Goroutine 是 Go 语言中的一种轻量级线程,由 Go 运行时管理。每个 Goroutine 都由 Go 运行时的调度器(scheduler)调度执行,而不像操作系统的线程一样由操作系统直接调度。Goroutine 允许我们在程序中并发地执行多个任务。

1.1、协程的创建

go
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 的特性

  1. 轻量性:相比于操作系统线程,Goroutine 更加轻量。每个 Goroutine 的栈空间非常小,通常只有 2KB 左右,而操作系统线程的栈通常需要几 MB 的内存。Goroutine 会在运行时动态扩展和收缩栈空间,因此可以更高效地使用内存。
  2. 并发性:Goroutine 通过 Go 运行时的调度器进行调度,多个 Goroutine 可能被映射到少数几个操作系统线程上执行。调度器负责将 Goroutine 分配到可用的 CPU 核心上执行,实现高效的并发。
  3. 非阻塞执行:当启动一个 Goroutine 后,主程序会继续执行,而不会等待该 Goroutine 执行完毕。换句话说,go 关键字会让指定的函数以并发的方式执行,不会阻塞主程序。
  4. 调度机制:Go 运行时的调度器使用的是 M:N 调度模型(即多个 Goroutine 运行在多个操作系统线程上)。Go 调度器通过抢占式调度和协作式调度来管理 Goroutine 的执行。这个调度模型是 Go 并发模型的核心。

示例:启动多个 Goroutine

通过 go 关键字,我们可以启动多个 Goroutine 来并发执行多个任务。以下是一个例子:

go
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 通常会与其他同步机制配合使用(例如,使用 ChannelWaitGroup)来等待所有任务完成。

1.3、 协程的调度机制

M:N调度模型

Go运行时采用M:N调度模型,其中:

  • M 表示操作系统线程。
  • N 表示协程(goroutine)。

这种模型允许多个协程在少数几个操作系统线程上多路复用,从而提高资源利用率和并发性能。当某个协程阻塞时,调度器会将其挂起,并将其他就绪的协程调度到当前线程上执行,从而避免线程阻塞。

调度器的工作原理

Go的调度器负责将协程分配到可用的操作系统线程上执行。调度器会根据协程的状态(如就绪、运行、阻塞等)进行动态调度。当某个协程阻塞时(例如等待I/O操作),调度器会将其挂起,并将其他就绪的协程调度到当前线程上执行。

1.4、主死从随

在 Go 中,协程(goroutine)是非常轻量的线程,主程序中的主协程并不会因为启动了其他协程而阻塞。主协程结束后,程序就会退出,无论其他协程是否执行完毕。因此,通常我们需要确保所有协程都执行完成后,主程序才能结束。

主协程死亡后行为的随机性

当主协程提前结束时,其他协程的行为会有一定的 随机性,主要体现在以下几个方面:

  1. 主协程提前退出:如果主协程(通常是 main 函数)没有等到其他协程完成就退出,Go 程序就会结束,所有其他正在运行的协程会被强制终止。
  2. 依赖调度器的随机性:Go 调度器(runtime scheduler)会调度并运行所有的协程,但调度顺序和执行时长并不能完全预知。主协程结束后,由于协程调度是并发进行的,可能会在不同的时间点终止,造成某些协程可能还没有完全执行完就被中止。
  3. 内存和资源的随机释放:如果主协程过早结束,它也可能会导致程序退出时对资源的释放不完全。具体来说,由于内存和资源的管理可能依赖于某些延迟释放机制(如垃圾回收、文件关闭等),所以如果主程序过早退出,有些资源可能还没有被正确释放,造成某些协程的执行受到影响。
示例:主协程死亡导致子协程随机执行

考虑以下代码:

go
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!")
}
运行结果(每次可能不同):
go
Main goroutine is exiting!
Task 1 started
Task 3 started
Task 2 started

或者:

go
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 协程退出后随机地调度正在执行的其他协程,这种行为可能在每次运行时表现出不同的结果。
  • 提前退出:在没有同步机制(如 WaitGroupChannel)的情况下,主协程提前退出时,其他协程仍然可以继续执行,但如果主协程在其还没有完成之前就退出了,程序会强制退出,导致所有正在运行的协程被中止。
解决方案:使用 sync.WaitGroup

为了确保主协程在所有子协程执行完毕后退出,我们可以使用 sync.WaitGroup 来同步所有 Goroutine 的执行。

go
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!")
}
运行结果(确保子协程都执行完毕):
go
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 提供了一个基本的锁机制,通过 LockUnlock 方法来显式地加锁和解锁。
  • 应用场景:当多个协程需要对共享资源进行修改时,可以使用互斥锁来保证每次只有一个协程能够访问该资源。
示例:使用 sync.Mutex 保护共享资源
go
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
go
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 包,允许我们对整数值等进行原子操作,而无需使用互斥锁。这对于性能要求较高的场景非常有用,因为原子操作是无锁的。

示例:使用原子操作增加计数
go
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 资源浪费。
  • 死锁:多个协程在等待彼此释放锁时,程序就会进入死锁状态,导致所有协程无法继续执行。死锁通常是由于锁获取顺序不一致或者锁的使用不当引起的。
示例:死锁示例
go
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 的类型和容量:

go
ch := make(chan int) // 创建一个 int 类型的 Channel
向 Channel 发送数据

使用 <- 运算符向 Channel 发送数据:

go
ch <- 42 // 向 Channel 发送数据 42
从 Channel 接收数据

使用 <- 运算符从 Channel 接收数据:

go
value := <-ch // 从 Channel 接收数据
使用 Channel 进行同步

Channel 不仅用于传输数据,也可以用于 Goroutine 之间的同步。例如,我们可以使用 Channel 等待多个 Goroutine 执行完毕:

go
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")
}

输出:

go
Task 1 is running
Task 2 is running
Task 3 is running
All tasks completed
Channel 的关闭

当我们不再需要向 Channel 发送数据时,可以使用 close() 来关闭 Channel:

go
close(ch)

关闭的 Channel 不能再发送数据,但可以继续从中接收数据。通常,关闭 Channel 用于告知接收方所有数据已经发送完毕。

2.2、遍历

Go 语言中的 range 可以用来遍历一个管道中的数据,直到管道关闭。遍历操作会从管道中依次接收元素,直到管道被关闭并且没有更多的数据可以接收为止。

示例:遍历管道中的数据
go
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)是最基本的管道类型,它的特性是:

  • 发送操作 会被阻塞,直到有一个接收操作准备好接收数据。
  • 接收操作 会被阻塞,直到有一个发送操作向管道中发送数据。

这种类型的管道用于需要严格同步的场景,因为发送者和接收者必须保持同步。

示例:无缓冲管道
go
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)允许在管道中存储一定数量的数据,发送操作不会立即阻塞,直到缓冲区满时才会阻塞。接收操作也类似,只有在管道为空时才会阻塞。

有缓冲的管道非常适合于解耦生产者和消费者的速度差异。

示例:有缓冲管道
go
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 遍历管道
go
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 同时处理多个管道
go
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
go
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 会执行超时分支。