You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

关键问题解答

  1. 是否可以用goroutines并行处理?
    完全可以。每个操作阶段启动独立goroutine,处理输入通道中的分块,实现流水线式并行。注意:如果操作需要全局状态(比如全局求和),可通过Aggregate方法汇总所有分块的中间结果。

  2. 如何在各操作间传递数据?
    通过带缓冲的通道传递带索引的分块。每个分块携带原始的全局序号,既保证了操作间的数据流转,又为后续顺序写入提供依据。

  3. 如何按原始顺序写入文件?
    用map暂存处理后的分块,维护一个nextIndex变量记录下一个要写入的分块序号。每当收到一个分块,就检查是否可以连续写入当前及后续的分块,直到没有连续的分块为止。这种方式即使goroutine处理速度不同,也能保证输出顺序和原始文件一致。

内容的提问来源于stack exchange,提问作者Adam G

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.24 05:50:55