Go处理大CSV文件:分块读取+多线程处理+有序结果聚合
超大CSV分块并行处理方案
针对你的需求,完全可以通过goroutines实现分块并行处理,同时保证结果按原始顺序写入。下面是具体的设计思路和代码实现:
核心改造思路
原来的全量操作接口是针对整个内存中的表,分块场景下需要将操作改为针对单个数据块,同时给每个分块标记索引以保证顺序。我们需要:
- 重新定义操作接口,支持分块处理、初始化和结果聚合
- 用通道(channel)实现操作间的流水线数据传递
- 通过索引暂存+顺序写入保证输出顺序
1. 接口与结构体设计
分块操作接口
import ( "encoding/csv" "fmt" "io" "os" "sort" "strconv" ) // ChunkOperation 定义分块操作的核心接口 type ChunkOperation interface { // Apply 处理单个数据块,chunkIndex为分块的全局序号(从0开始) Apply(chunk [][]string, chunkIndex int) ([][]string, error) // Init 操作初始化(可选,比如校验参数、初始化全局状态) Init() error // Aggregate 聚合所有分块的结果(可选,比如全局求和、统计) Aggregate([]interface{}) (interface{}, error) }
CSV处理器结构体(链式调用核心)
type CSVProcessor struct { inputPath string outputPath string chunkSize int // 每个分块的行数,默认1000 ops []ChunkOperation err error // 链式调用中提前捕获的错误 } // Read 创建CSV处理器实例,指定输入文件路径 func Read(path string) *CSVProcessor { return &CSVProcessor{ inputPath: path, chunkSize: 1000, } } // SetChunkSize 自定义分块行数 func (p *CSVProcessor) SetChunkSize(size int) *CSVProcessor { if size <= 0 { p.err = fmt.Errorf("chunk size must be positive, got %d", size) return p } p.chunkSize = size return p } // With 添加分块操作到处理器 func (p *CSVProcessor) With(op ChunkOperation) *CSVProcessor { if p.err != nil { return p } if err := op.Init(); err != nil { p.err = err return p } p.ops = append(p.ops, op) return p }
2. 核心流程实现
分块读取CSV
// readChunks 分块读取CSV文件,返回带索引的分块通道和错误通道 func (p *CSVProcessor) readChunks() (<-chan struct { index int chunk [][]string }, <-chan error) { chunkChan := make(chan struct { index int chunk [][]string }, 5) // 缓冲通道,平衡读取和处理速度 errChan := make(chan error, 1) go func() { defer close(chunkChan) file, err := os.Open(p.inputPath) if err != nil { errChan <- err return } defer file.Close() reader := csv.NewReader(file) chunk := make([][]string, 0, p.chunkSize) index := 0 for { row, err := reader.Read() if err == io.EOF { break } if err != nil { errChan <- err return } chunk = append(chunk, row) if len(chunk) == p.chunkSize { chunkChan <- struct { index int chunk [][]string }{index, chunk} chunk = make([][]string, 0, p.chunkSize) index++ } } // 处理最后一个不足分块大小的剩余数据 if len(chunk) > 0 { chunkChan <- struct { index int chunk [][]string }{index, chunk} } errChan <- nil }() return chunkChan, errChan }
流水线并行处理
每个操作阶段启动独立goroutine,从输入通道读取分块,处理后发送到下一个通道:
// processPipeline 构建操作流水线,并行处理每个分块 func (p *CSVProcessor) processPipeline(inputChan <-chan struct { index int chunk [][]string }) <-chan struct { index int chunk [][]string } { currentChan := inputChan for _, op := range p.ops { op := op // 闭包捕获变量,避免循环变量复用问题 nextChan := make(chan struct { index int chunk [][]string }, 5) go func() { defer close(nextChan) for item := range currentChan { processedChunk, err := op.Apply(item.chunk, item.index) if err != nil { panic(fmt.Sprintf("process chunk %d failed: %v", item.index, err)) } nextChan <- struct { index int chunk [][]string }{item.index, processedChunk} } }() currentChan = nextChan } return currentChan }
按原始顺序写入结果
用map暂存处理后的分块,按索引顺序连续写入:
// Write 执行所有操作并将结果按原始顺序写入输出文件 func (p *CSVProcessor) Write(path string) error { if p.err != nil { return p.err } p.outputPath = path // 1. 分块读取 chunkChan, errChan := p.readChunks() if err := <-errChan; err != nil { return err } // 2. 流水线并行处理 processedChan := p.processPipeline(chunkChan) // 3. 按顺序写入 file, err := os.Create(p.outputPath) if err != nil { return err } defer file.Close() writer := csv.NewWriter(file) defer writer.Flush() resultMap := make(map[int][][]string) nextIndex := 0 for item := range processedChan { resultMap[item.index] = item.chunk // 写入连续的分块 for { chunk, exists := resultMap[nextIndex] if !exists { break } for _, row := range chunk { if err := writer.Write(row); err != nil { return fmt.Errorf("write row failed: %v", err) } } delete(resultMap, nextIndex) nextIndex++ } } // 检查是否有遗漏的分块 if len(resultMap) > 0 { missing := make([]int, 0, len(resultMap)) for k := range resultMap { missing = append(missing, k) } sort.Ints(missing) return fmt.Errorf("missing chunks: %v", missing) } return nil }
3. 具体操作实现示例
GetColumns操作
type GetColumnsOperation struct { start, end int // 列索引(从0开始,左闭右开) } func GetColumns(start, end int) ChunkOperation { return &GetColumnsOperation{start: start, end: end} } func (op *GetColumnsOperation) Init() error { if op.start < 0 || op.end <= op.start { return fmt.Errorf("invalid column range: start=%d, end=%d", op.start, op.end) } return nil } func (op *GetColumnsOperation) Apply(chunk [][]string, _ int) ([][]string, error) { processed := make([][]string, 0, len(chunk)) for _, row := range chunk { if op.end > len(row) { return nil, fmt.Errorf("row has only %d columns, requested end=%d", len(row), op.end) } processed = append(processed, row[op.start:op.end]) } return processed, nil } func (op *GetColumnsOperation) Aggregate(_ []interface{}) (interface{}, error) { return nil, nil // 无需聚合 }
SumRow操作(每行指定列求和,添加到行末尾)
type SumRowOperation struct { sumStart, sumEnd int // 要求和的列范围 } func SumRow(sumStart, sumEnd int) ChunkOperation { return &SumRowOperation{sumStart: sumStart, sumEnd: sumEnd} } func (op *SumRowOperation) Init() error { if op.sumStart < 0 || op.sumEnd <= op.sumStart { return fmt.Errorf("invalid sum column range: start=%d, end=%d", op.sumStart, op.sumEnd) } return nil } func (op *SumRowOperation) Apply(chunk [][]string, _ int) ([][]string, error) { processed := make([][]string, 0, len(chunk)) for _, row := range chunk { sum := 0.0 end := op.sumEnd if end > len(row) { end = len(row) } for _, col := range row[op.sumStart:end] { num, err := strconv.ParseFloat(col, 64) if err != nil { return nil, fmt.Errorf("invalid number in row: %s", col) } sum += num } newRow := append(row[:], strconv.FormatFloat(sum, 'f', 2, 64)) processed = append(processed, newRow) } return processed, nil } func (op *SumRowOperation) Aggregate(_ []interface{}) (interface{}, error) { return nil, nil // 无需聚合 }
4. 使用示例
func main() { err := Read("input.csv"). SetChunkSize(2000). With(GetColumns(3, 5)). With(SumRow(0, 2)). Write("output.csv") if err != nil { panic(err) } }
关键问题解答
是否可以用goroutines并行处理?
完全可以。每个操作阶段启动独立goroutine,处理输入通道中的分块,实现流水线式并行。注意:如果操作需要全局状态(比如全局求和),可通过Aggregate方法汇总所有分块的中间结果。如何在各操作间传递数据?
通过带缓冲的通道传递带索引的分块。每个分块携带原始的全局序号,既保证了操作间的数据流转,又为后续顺序写入提供依据。如何按原始顺序写入文件?
用map暂存处理后的分块,维护一个nextIndex变量记录下一个要写入的分块序号。每当收到一个分块,就检查是否可以连续写入当前及后续的分块,直到没有连续的分块为止。这种方式即使goroutine处理速度不同,也能保证输出顺序和原始文件一致。
内容的提问来源于stack exchange,提问作者Adam G
相关产品推荐
相关产品推荐

