处理异步数据流的通用模式

处理异步数据流的通用模式

当使用通道来串行化任务时,异步流通常就诞生了。

异步编程很容易搞砸,尤其需求丰富,健壮度高时,比如拥有超时取消等。 一个经典的例子,是从文件逐行读取,交给另一个异步过程解析。通常I/O慢而计算快,为了最大化效率,一般会并行读取,放入缓存,再顺序处理,此时就需要管道来协调;如遇到读取错误,超时等,需终止整个流程。

我找到了一个通用模式,其内在逻辑很容易理解:

异步慢过程 后跟 同步快过程。

如果慢过程是 provider,则模式是:

1go provider()
2consumer()

一个完整的极简框架:

 1var (
 2    err              error
 3    provider         = make(chan int)
 4    consumer         = make(chan int)
 5    ctx, cancel      = context.WithTimeout(context.Background, time.Second)
 6)
 7defer cancel()
 8
 9// 生产者
10go func() {
11    defer close(consumer)
12
13    for {
14        select {
15        case <-ctx.Done():
16            return
17        case v, ok := <-provider:
18            if !ok {
19                return
20            }
21            consumer <- v
22        }
23    }
24}()
25
26consume:
27for {
28    select {
29    case <-ctx.Done():
30        return
31    case v, ok := <-consumer:
32        if !ok {
33            break consume
34        }
35
36        if err = process(v); err != nil {
37            break consume
38        }
39    }
40}
41
42// 后续逻辑
43if err != nil {
44    log.Error(err)
45}

可见,provider 部分是一个 goroutine,它从 provider 通道中读取数据并将其传递给 consumer 通道。consumer 部分则是一个阻塞的 for loop,保证程序按顺序完成。至于消息是如何进provider的,这完全由独立逻辑控制,比如一个额外的异步函数,逐行读取文件,直到结束。

如果是 provider 的过程快而 consumer 过程慢呢?比如先通过计算,快速生成消息,再缓慢的写入到文件,则需要颠倒启动的顺序:

1go consumer()
2provider()

先启动的 consumer 始终在监听 provider 的通道获取消息,provider 则是一个阻塞的同步过程,生成消息。

祝我们在异步编程中好运!