【发布时间】:2021-09-08 06:32:37
【问题描述】:
我正在处理一个涉及生产者-消费者模式的问题。我有一个生产任务的生产者和消费任务的“n”个消费者。消费者任务是从文件中读取一些数据,然后将该数据上传到 S3。一位消费者可以读取多达 xMB(8/16/32) 的数据,然后将其上传到 s3。将所有数据保存在内存中导致内存消耗超过程序的预期,所以我切换到从文件中读取数据,然后将其写入某个临时文件,然后将文件上传到 S3,尽管这在内存方面,但CPU受到了打击。我想知道是否有任何方法可以一次分配固定大小的内存,然后在不同的 goroutines 中使用它? 我想要的是,如果我有 4 个 goroutine,那么我可以分配 4 个不同的 xMB 数组,然后在每个 goroutine 调用中使用相同的数组,这样 goroutine 就不会每次都分配内存,也不依赖于GC 释放内存?
编辑:添加我的代码的症结所在。我的 go 消费者看起来像:
type struct Block {
offset int64
size int64
}
func consumer (blocks []Block) {
var dataArr []byte
for _, block := range blocks {
data := file.Read(block.offset, block.size)
dataArr = append(dataArr, data)
}
upload(dataArr)
}
我基于Blocks从文件中读取数据,这个块可以包含几个xMB限制的小块或一个xMB的大块。
Edit2:根据评论中的建议尝试了 sync.Pool。但我没有看到内存消耗有任何改善。我做错了吗?
var pool *sync.Pool
func main() {
pool = &sync.Pool{
New: func()interface{} {
return make([]byte, 16777216)
},
}
for i:=0; i < 4; i++ {
// blocks is 2-d array each index contains array of blocks.
go consumer(blocks[i])
}
}
go consumer(blocks []Blocks) {
var dataArr []byte
d := pool.(Get).([]byte)
for _, block := range blocks {
file.Read(block.offset,block.size,d[block.offset:block.size])
}
upload(data)
pool.put(data)
}
【问题讨论】:
-
您可以使用小缓冲区流式传输文件。
-
对于我们看不到的代码,我们无能为力,但几乎每个这样的 API 都提供了一种流式传输数据的方法,以防止您不得不缓冲整个文件。
-
您不需要读取整个文件然后再写入。应该有一个获取 io.Reader 的 s3 上传 API,以便您可以流式传输文件。您的网络带宽受到限制,因此如果将其分配给多个 goroutine,您可能不会获得太多的吞吐量提升。
-
使用 io.Pipe。分块读取数据,写入管道。使用管道的阅读器端上传。
-
每个人都提到流媒体绝对是正确的。只是为了回答“我想知道是否有任何方法可以分配一次固定大小的内存,然后在不同的 goroutines 中使用它?”,您正在寻找的可能是
sync.Pool。
标签: go memory-management heap-memory goroutine memory-pool