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 }
问题分析与修复方案
核心问题
- 文件句柄冲突:先调用
os.Create创建Parquet文件,又用local.NewLocalFileWriter打开同一个文件,导致两个句柄操作同一文件,写入逻辑混乱。 - 第一行数据丢失:读取
firstDataRow用于类型检测后,未将其写入Parquet,且如果CSV只有一行数据,循环中不会写入任何内容。 - Parquet Schema结构错误:Parquet要求Schema必须有一个根节点(结构体类型),直接将列作为顶层元素会导致Writer初始化异常,最终
WriteStop时触发空指针。 - 忽略类型转换错误:类型转换时直接忽略错误,可能导致非法数据写入,引发后续异常。
修复后的代码
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 }
关键修改说明
- 移除重复文件创建:直接使用
local.NewLocalFileWriter创建并管理Parquet文件,避免句柄冲突。 - 读取所有行数据:一次性读取CSV所有行,既用于类型检测,也确保第一行数据被写入。
- 修正Schema结构:添加名为
Root的根节点,所有列作为根节点的子元素,符合Parquet的Schema规范。 - 处理缺失列数据:当行数据长度小于表头数量时
相关产品推荐
相关产品推荐

