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

Golang动态CSV转Parquet报错:空指针引用及文件过小问题求助

动态CSV转Parquet的空指针错误问题

我用Golang实现任意CSV文件转Parquet格式的功能时,遇到了空指针引用错误:Error finalizing Parquet file: runtime error: invalid memory address or nil pointer dereference,生成的Parquet文件仅4字节,小于Parquet要求的最小8字节footer。之前只能通过绑定特定结构体实现转换,现在需要适配任意CSV,尝试多种方法仍无法解决。

原代码

package converted

import (
    "bufio"
    "encoding/csv"
    "fmt"
    "io"
    "os"
    "path/filepath"
    "reflect"
    "strconv"
    "strings"

    "github.com/xitongsys/parquet-go-source/local"
    "github.com/xitongsys/parquet-go/parquet"
    "github.com/xitongsys/parquet-go/writer"
    "golang.org/x/text/encoding/unicode"
    "golang.org/x/text/transform"
)

// detectColumnType attempts to infer the type of a column value.
func detectColumnType(value string) reflect.Kind {
    if _, err := strconv.ParseBool(value); err == nil {
        return reflect.Bool
    }
    if _, err := strconv.ParseFloat(value, 64); err == nil {
        return reflect.Float64
    }
    if _, err := strconv.Atoi(value); err == nil {
        return reflect.Int
    }
    return reflect.String
}

// convertToUTF8 converts a string to UTF-8 encoding.
func convertToUTF8(input string) string {
    utf8Decoder := unicode.UTF8.NewDecoder()
    reader := transform.NewReader(io.NopCloser(strings.NewReader(input)), utf8Decoder)
    result, _ := io.ReadAll(reader)
    return string(result)
}

// CsvToParquet converts a CSV file to a Parquet file.
func CsvToParquet(csvFilePath, parquetFileName string) (string, error) {
    csvDir := filepath.Dir(csvFilePath)
    parquetFilePath := filepath.Join(csvDir, parquetFileName)
    fmt.Printf("Input CSV File Path: %s\n", csvFilePath)
    fmt.Printf("Output Parquet File Path: %s\n", parquetFilePath)

    // Open the CSV file for reading.
    csvFile, err := os.Open(csvFilePath)
    if err != nil {
        return "", fmt.Errorf("failed to open CSV file: %w", err)
    }
    defer csvFile.Close()

    // Create a CSV reader.
    reader := csv.NewReader(bufio.NewReader(csvFile))

    // Read the header row.
    header, err := reader.Read()
    if err != nil {
        return "", fmt.Errorf("failed to read CSV header: %w", err)
    }

    // Detect column types using the first data row.
    schema := make(map[string]reflect.Kind)
    firstDataRow, err := reader.Read()
    if err != nil && err != io.EOF {
        return "", fmt.Errorf("failed to read first data row: %w", err)
    }
    if len(firstDataRow) == 0 && len(header) > 0 {
        for _, col := range header {
            schema[col] = reflect.String
        }
    } else if len(firstDataRow) == 0 && len(header) == 0 {
        return "", fmt.Errorf("CSV file is empty")
    } else {
        for i, col := range header {
            if i < len(firstDataRow) {
                schema[col] = detectColumnType(firstDataRow[i])
            } else {
                schema[col] = reflect.String
            }
        }
    }

    // Create the Parquet file.
    file, err := os.Create(parquetFilePath)
    if err != nil {
        return "", fmt.Errorf("failed to create Parquet file: %w", err)
    }
    defer file.Close()

    // Initialize the local file writer for Parquet.
    source, err := local.NewLocalFileWriter(parquetFilePath)
    if err != nil {
        return "", fmt.Errorf("failed to create local file writer: %w", err)
    }

    // Build the Parquet schema.
    schemaElements := []*parquet.SchemaElement{}
    for col, kind := range schema {
        var typeValue parquet.Type
        switch kind {
        case reflect.Int:
            typeValue = parquet.Type_INT64
        case reflect.Float64:
            typeValue = parquet.Type_DOUBLE
        case reflect.Bool:
            typeValue = parquet.Type_BOOLEAN
        default:
            typeValue = parquet.Type_BYTE_ARRAY
        }
        repetitionType := parquet.FieldRepetitionType_REQUIRED
        schemaElements = append(schemaElements, &parquet.SchemaElement{
            Name:           col,
            Type:           ptrType(typeValue),
            NumChildren:    ptrInt32(0),
            RepetitionType: ptrFieldRepetitionType(repetitionType),
        })
    }

    // Log the generated schema for debugging purposes.
    for _, se := range schemaElements {
        fmt.Printf("Schema Element: Name=%s, Type=%v, RepetitionType=%v\n", se.Name, *se.Type, *se.RepetitionType)
    }

    // Create the Parquet file metadata.
    fileMetaData := &parquet.FileMetaData{
        Version: 1,
        Schema:  schemaElements,
    }

    // Initialize the Parquet writer.
    pw, err := writer.NewParquetWriter(source, fileMetaData, 8)
    if err != nil {
        return "", fmt.Errorf("failed to initialize Parquet writer: %w", err)
    }
    if pw == nil {
        return "", fmt.Errorf("parquet writer is nil")
    }
    defer func() {
        if pw != nil {
            if closeErr := pw.WriteStop(); closeErr != nil {
                fmt.Printf("Error finalizing Parquet file: %v\n", closeErr)
            }
        }
    }()

    // Write rows to the Parquet file.
    rowCount := 0
    for {
        row, err := reader.Read()
        if err == io.EOF {
            break
        }
        if err != nil {
            return "", fmt.Errorf("failed to read CSV row: %w", err)
        }

        // Convert the row into a map based on the detected schema.
        data := make(map[string]interface{})
        for i, col := range header {
            if i >= len(row) {
                continue
            }
            switch schema[col] {
            case reflect.Int:
                val, _ := strconv.Atoi(row[i])
                data[col] = val
            case reflect.Float64:
                val, _ := strconv.ParseFloat(row[i], 64)
                data[col] = val
            case reflect.Bool:
                val, _ := strconv.ParseBool(row[i])
                data[col] = val
            default:
                data[col] = convertToUTF8(row[i])
            }
        }

        // Write the row to the Parquet file.
        if err := pw.Write(data); err != nil {
            return "", fmt.Errorf("failed to write to Parquet: %w", err)
        }
        rowCount++
        fmt.Printf("Wrote row: %+v\n", data)
    }

    fmt.Printf("Wrote %d rows to the Parquet file.\n", rowCount)

    return parquetFilePath, nil
}

// Helper functions to create pointers.
func ptrInt32(val int32) *int32 {
    v := val
    return &v
}

func ptrType(val parquet.Type) *parquet.Type {
    v := val
    return &v
}

func ptrFieldRepetitionType(val parquet.FieldRepetitionType) *parquet.FieldRepetitionType {
    v := val
    return &v
}

问题分析与修复方案

核心问题

  1. 文件句柄冲突:先调用os.Create创建Parquet文件,又用local.NewLocalFileWriter打开同一个文件,导致两个句柄操作同一文件,写入逻辑混乱。
  2. 第一行数据丢失:读取firstDataRow用于类型检测后,未将其写入Parquet,且如果CSV只有一行数据,循环中不会写入任何内容。
  3. Parquet Schema结构错误:Parquet要求Schema必须有一个根节点(结构体类型),直接将列作为顶层元素会导致Writer初始化异常,最终WriteStop时触发空指针。
  4. 忽略类型转换错误:类型转换时直接忽略错误,可能导致非法数据写入,引发后续异常。

修复后的代码

package converted

import (
    "bufio"
    "encoding/csv"
    "fmt"
    "io"
    "os"
    "path/filepath"
    "reflect"
    "strconv"
    "strings"

    "github.com/xitongsys/parquet-go-source/local"
    "github.com/xitongsys/parquet-go/parquet"
    "github.com/xitongsys/parquet-go/writer"
    "golang.org/x/text/encoding/unicode"
    "golang.org/x/text/transform"
)

// detectColumnType 推断列值类型
func detectColumnType(value string) reflect.Kind {
    if value == "" {
        return reflect.String
    }
    if _, err := strconv.ParseBool(value); err == nil {
        return reflect.Bool
    }
    if _, err := strconv.ParseFloat(value, 64); err == nil {
        return reflect.Float64
    }
    if _, err := strconv.Atoi(value); err == nil {
        return reflect.Int
    }
    return reflect.String
}

// convertToUTF8 转换字符串为UTF-8编码
func convertToUTF8(input string) string {
    utf8Decoder := unicode.UTF8.NewDecoder()
    reader := transform.NewReader(strings.NewReader(input), utf8Decoder)
    result, err := io.ReadAll(reader)
    if err != nil {
        return input
    }
    return string(result)
}

// CsvToParquet 将CSV文件转换为Parquet文件
func CsvToParquet(csvFilePath, parquetFileName string) (string, error) {
    csvDir := filepath.Dir(csvFilePath)
    parquetFilePath := filepath.Join(csvDir, parquetFileName)
    fmt.Printf("Input CSV File Path: %s\n", csvFilePath)
    fmt.Printf("Output Parquet File Path: %s\n", parquetFilePath)

    // 打开CSV文件
    csvFile, err := os.Open(csvFilePath)
    if err != nil {
        return "", fmt.Errorf("打开CSV文件失败: %w", err)
    }
    defer csvFile.Close()

    reader := csv.NewReader(bufio.NewReader(csvFile))
    // 读取表头
    header, err := reader.Read()
    if err != nil {
        return "", fmt.Errorf("读取CSV表头失败: %w", err)
    }

    // 读取所有行数据,包括第一行
    allRows := make([][]string, 0)
    for {
        row, err := reader.Read()
        if err == io.EOF {
            break
        }
        if err != nil {
            return "", fmt.Errorf("读取CSV行失败: %w", err)
        }
        allRows = append(allRows, row)
    }

    // 检测列类型
    schema := make(map[string]reflect.Kind)
    if len(allRows) == 0 {
        // 只有表头,所有列设为字符串
        for _, col := range header {
            schema[col] = reflect.String
        }
    } else {
        firstRow := allRows[0]
        for i, col := range header {
            if i < len(firstRow) {
                schema[col] = detectColumnType(firstRow[i])
            } else {
                schema[col] = reflect.String
            }
        }
    }

    // 创建Parquet文件写入源,移除多余的os.Create
    source, err := local.NewLocalFileWriter(parquetFilePath)
    if err != nil {
        return "", fmt.Errorf("创建Parquet文件写入源失败: %w", err)
    }
    defer source.Close()

    // 构建Parquet Schema:添加根节点
    var schemaElements []*parquet.SchemaElement
    // 根节点(结构体类型)
    rootElement := &parquet.SchemaElement{
        Name:           "Root",
        NumChildren:    ptrInt32(int32(len(header))),
        RepetitionType: ptrFieldRepetitionType(parquet.FieldRepetitionType_REQUIRED),
    }
    schemaElements = append(schemaElements, rootElement)

    // 添加列节点
    for col, kind := range schema {
        var typeValue parquet.Type
        switch kind {
        case reflect.Int:
            typeValue = parquet.Type_INT64
        case reflect.Float64:
            typeValue = parquet.Type_DOUBLE
        case reflect.Bool:
            typeValue = parquet.Type_BOOLEAN
        default:
            typeValue = parquet.Type_BYTE_ARRAY
        }
        schemaElements = append(schemaElements, &parquet.SchemaElement{
            Name:           col,
            Type:           ptrType(typeValue),
            NumChildren:    ptrInt32(0),
            RepetitionType: ptrFieldRepetitionType(parquet.FieldRepetitionType_REQUIRED),
        })
    }

    // 打印Schema用于调试
    for _, se := range schemaElements {
        fmt.Printf("Schema Element: Name=%s, Type=%v, RepetitionType=%v\n", se.Name, se.Type, se.RepetitionType)
    }

    // 创建文件元数据
    fileMetaData := &parquet.FileMetaData{
        Version: 1,
        Schema:  schemaElements,
    }

    // 初始化Parquet Writer
    pw, err := writer.NewParquetWriter(source, fileMetaData, 8)
    if err != nil {
        return "", fmt.Errorf("初始化Parquet Writer失败: %w", err)
    }
    if pw == nil {
        return "", fmt.Errorf("Parquet Writer为空")
    }
    defer func() {
        if err := pw.WriteStop(); err != nil {
            fmt.Printf("Finalizing Parquet文件失败: %v\n", err)
        }
    }()

    // 写入所有行数据
    rowCount := 0
    for _, row := range allRows {
        data := make(map[string]interface{})
        for i, col := range header {
            if i >= len(row) {
                data[col] = getDefaultValue(schema[col])
                continue
            }
            valStr := row[i]
            switch schema[col] {
            case reflect.Int:
                val, err := strconv.Atoi(valStr)
                if err != nil {
                    return "", fmt.Errorf("列%s转换为整数失败: %w", col, err)
                }
                data[col] = val
            case reflect.Float64:
                val, err := strconv.ParseFloat(valStr, 64)
                if err != nil {
                    return "", fmt.Errorf("列%s转换为浮点数失败: %w", col, err)
                }
                data[col] = val
            case reflect.Bool:
                val, err := strconv.ParseBool(valStr)
                if err != nil {
                    return "", fmt.Errorf("列%s转换为布尔值失败: %w", col, err)
                }
                data[col] = val
            default:
                data[col] = convertToUTF8(valStr)
            }
        }

        if err := pw.Write(data); err != nil {
            return "", fmt.Errorf("写入Parquet行失败: %w", err)
        }
        rowCount++
        fmt.Printf("已写入行: %+v\n", data)
    }

    fmt.Printf("共写入%d行到Parquet文件\n", rowCount)
    return parquetFilePath, nil
}

// getDefaultValue 获取对应类型的默认值
func getDefaultValue(kind reflect.Kind) interface{} {
    switch kind {
    case reflect.Int:
        return 0
    case reflect.Float64:
        return 0.0
    case reflect.Bool:
        return false
    default:
        return ""
    }
}

// ptrInt32 创建int32指针
func ptrInt32(val int32) *int32 {
    v := val
    return &v
}

// ptrType 创建parquet.Type指针
func ptrType(val parquet.Type) *parquet.Type {
    v := val
    return &v
}

// ptrFieldRepetitionType 创建FieldRepetitionType指针
func ptrFieldRepetitionType(val parquet.FieldRepetitionType) *parquet.FieldRepetitionType {
    v := val
    return &v
}

关键修改说明

  1. 移除重复文件创建:直接使用local.NewLocalFileWriter创建并管理Parquet文件,避免句柄冲突。
  2. 读取所有行数据:一次性读取CSV所有行,既用于类型检测,也确保第一行数据被写入。
  3. 修正Schema结构:添加名为Root的根节点,所有列作为根节点的子元素,符合Parquet的Schema规范。
  4. 处理缺失列数据:当行数据长度小于表头数量时
相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 04:21:30