惯性聚合 高效追踪和阅读你感兴趣的博客、新闻、科技资讯
阅读原文 在惯性聚合中打开

推荐订阅源

钛媒体:引领未来商业与生活新知
钛媒体:引领未来商业与生活新知
U
Unit 42
GbyAI
GbyAI
M
MIT News - Artificial intelligence
美团技术团队
罗磊的独立博客
雷峰网
雷峰网
量子位
博客园 - 【当耐特】
Last Week in AI
Last Week in AI
D
Docker
小众软件
小众软件
S
SegmentFault 最新的问题
Blog — PlanetScale
Blog — PlanetScale
阮一峰的网络日志
阮一峰的网络日志
宝玉的分享
宝玉的分享
T
Tailwind CSS Blog
WordPress大学
WordPress大学
V
V2EX
博客园_首页
腾讯CDC
The Cloudflare Blog
A
About on SuperTechFans
Cyber Security Advisories - MS-ISAC
Cyber Security Advisories - MS-ISAC

OhYee 博客

小鹏辅助驾驶测评|OhYee 博客 小鹏非支持手机开启自动解锁|OhYee 博客 使用函数计算实现 301 重定向|OhYee 博客 针对 HTML 内容使用 Ant Design 图片弹框|OhYee 博客 博客进程泄露及僵尸进程解决|OhYee 博客 蓝易云服务器体验|OhYee 博客 SSH 调起本地 VSCode|OhYee 博客 【2022 秋招内推】阿里云后端研发工程师|OhYee 博客 使用函数计算获取 IP 地址信息|OhYee 博客 正确获取客户端 IP/HTTP Header 也可能重复|OhYee 博客 评测 Oculus Quest2 及 BigScreen|OhYee 博客 NextJS 热重载保留状态|OhYee 博客 如何优雅地贴 gist 代码|OhYee 博客 Linux 精细化文件权限|OhYee 博客 VSCode 容器开发环境|OhYee 博客 Clash 的不兼容更新排查|OhYee 博客 Zeek 导出 PCAP|OhYee 博客 记一次 ssh 配置问题|OhYee 博客 Git Commit 规范化工具|OhYee 博客 谈谈《星之卡比-探索发现》|OhYee 博客 VSCode 快捷键绑定 Shell 命令|OhYee 博客 ASN.1 语法及 X.509 证书格式解析解析|OhYee 博客 腾讯企业邮箱忽略 MX 记录发信|OhYee 博客 Chrome/Edge 标签组插件|OhYee 博客 【应届内推】阿里云后端研发工程师|OhYee 博客 损坏的 Typecho 备份处理为 JSON|OhYee 博客 VS Code VIM 插件高效使用|OhYee 博客 SSH 正反向代理|OhYee 博客 Let's Encrypt 根证书过期引发的问题|OhYee 博客 OpenWRT 忽略内核依赖|OhYee 博客
Go io 流管道连接|OhYee 博客
2021-10-10 · via OhYee 博客

这是一篇最后编辑于 5 年前 的文章,其内容可能与目前实际情况差异较大,请注意甄别

Go io 流管道连接

Go 一系列的 io 操作的接口非常棒,可以将一切东西归为 Read()Write(),以及其他少数几个接口。完美诠释了 Linux 万物皆文件的思想
(对比其他语言,虽然也是同样的思路,但是还暴露的更多的接口,反而增大了使用复杂度)

在一些场景下,需要实现将两个流拼接起来,比如提供一个服务器上的 gRPC 服务,可以通过 K8S 或是 Docker 的接口,连接到容器内执行命令(exec)。Docker 本身提供的是一个 Reader & Writer 的接口,CRI 接口则要求的是 Reader、Writer 流。

而 gRPC 流接口则给出的是被阻塞的 Message。
因此需要将 gRPC 的数据转换为流,目前内置的各种相关功能都无法满足需要:

  • bytes.Buffer: 读完缓冲区就会直接结束(这时可能 gRPC 还没开始写入)
  • io.Copy: 复制到 EOF 就会结束
  • io.Pipe: 返回的是新的 Reader 和 Writer,无法直接连接 Reader 和 Writer

因此,应该仿照 io.Pipe 实现一个可以将 Write()Read() 函数连接起来,同时可以通过 ReadFrom()WriteTo() 将其连接到其他 Reader 和 Writer 中

使用样例

两个流直接连接

将 stdin 连接到 stdout

pipe.RWPipe(stdin, stdout)

将输入流同步到其他流中

stdin := pipe.New()
go stdin.WriteTo(os.Stdout)

将非流数据转换为流

in := pipe.New()
go in.WriteTo(os.Stdout)
for {
    msg := getMessage()
    in.Write(msg)
}

将流数据转换为非流

out := pipe.New()
go out.ReadFrom(os.Stdin)
reader := bufio.NewReader(out)
for {
    b, err := reader.ReadBytes('\n')
    send(b)
}

代码

代码中 WriteTo()ReadFrom() 需要在单独的 goroutine 中运行

package pipe

import (
	"io"
	"sync"
)

func RWPipe(r io.Reader, w io.Writer) (n int64, err error) {
	buf := make([]byte, 512)
	var n1, n2 int
	for {
		n1, err = r.Read(buf)
		if err != nil {
			break
		}
		n2, err = w.Write(buf[:n1])
		if err != nil {
			break
		}
		n += int64(n2)
	}
	return
}

type Pipe struct {
	wrMu sync.Mutex
	wrCh chan []byte
	rdCh chan int

	once sync.Once
	done chan struct{}
}

func New() *Pipe {
	return &Pipe{
		wrMu: sync.Mutex{},
		wrCh: make(chan []byte),
		rdCh: make(chan int),
		once: sync.Once{},
		done: make(chan struct{}),
	}
}

func (p *Pipe) Read(b []byte) (int, error) {
	select {
	case bw := <-p.wrCh:
		nr := copy(b, bw)
		p.rdCh <- nr
		return nr, nil
	case <-p.done:
		return 0, io.EOF
	}
}

func (p *Pipe) Write(b []byte) (int, error) {
	select {
	case <-p.done:
		return 0, io.EOF
	default:
		p.wrMu.Lock()
		defer p.wrMu.Unlock()
	}

	n := 0
	for once := true; once || len(b) > 0; once = false {
		select {
		case p.wrCh <- b:
			nw := <-p.rdCh
			b = b[nw:]
			n += nw
		case <-p.done:
			return n, io.EOF
		}
	}
	return n, nil
}

func (p *Pipe) ReadFrom(r io.Reader) (n int64, err error) {
	defer p.Close()
	return RWPipe(r, p)
}

func (p *Pipe) WriteTo(w io.Writer) (n int64, err error) {
	defer p.Close()
	return RWPipe(p, w)
}

func (p *Pipe) Close() { p.once.Do(func() { close(p.done) }) }