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

如何使用C#读取Parquet文件并将数据存入数据库表?

.NET实现Parquet文件数据插入数据库表

1. 安装依赖包

需要两个核心类库:

  • Parquet.Net:用于解析读取Parquet文件
  • 根据数据库类型选择数据访问类库,例如:
    • SQL Server:Microsoft.EntityFrameworkCore.SqlServer 或 Dapper
    • MySQL:MySql.EntityFrameworkCore 或 Dapper

通过NuGet CLI安装示例:

# 安装Parquet解析库
Install-Package Parquet.Net

# SQL Server + EF Core 示例
Install-Package Microsoft.EntityFrameworkCore.SqlServer

2. 定义实体类

创建与Parquet列、数据库表结构匹配的实体类,保证属性名/类型与目标结构兼容:

public class Employee
{
    public int Id { get; set; }
    public string Name { get; set; }
    public int Age { get; set; }
    public string Department { get; set; }
    public decimal Salary { get; set; }
}

如果Parquet列名与实体属性名不一致,后续读取时需手动映射对应关系

3. 读取Parquet文件数据

Parquet.Net提供两种读取方式,按需选择:

方式A:自动序列化(推荐,列名与实体属性匹配时使用)

using Parquet;

public static async Task<List<Employee>> ReadParquetFileAsync(string filePath)
{
    using var stream = File.OpenRead(filePath);
    // 自动将Parquet数据序列化为实体列表
    return await ParquetConvert.DeserializeAsync<Employee>(stream);
}

方式B:手动列映射(列名与实体属性不匹配时使用)

using Parquet;
using Parquet.Data;

public static List<Employee> ReadParquetFile(string filePath)
{
    var employees = new List<Employee>();

    using (var stream = File.OpenRead(filePath))
    {
        using (var reader = new ParquetReader(stream))
        {
            using (var rowGroupReader = reader.OpenRowGroupReader(0))
            {
                // 按Parquet文件的实际列名读取数据
                var idColumn = rowGroupReader.ReadColumn((DataField)reader.Schema.GetDataField("emp_id"));
                var nameColumn = rowGroupReader.ReadColumn((DataField)reader.Schema.GetDataField("emp_name"));
                var ageColumn = rowGroupReader.ReadColumn((DataField)reader.Schema.GetDataField("emp_age"));
                var departmentColumn = rowGroupReader.ReadColumn((DataField)reader.Schema.GetDataField("emp_dept"));
                var salaryColumn = rowGroupReader.ReadColumn((DataField)reader.Schema.GetDataField("emp_salary"));

                // 遍历行组装实体
                for (int i = 0; i < rowGroupReader.RowCount; i++)
                {
                    employees.Add(new Employee
                    {
                        Id = idColumn.Data.GetValue<int>(i),
                        Name = nameColumn.Data.GetValue<string>(i),
                        Age = ageColumn.Data.GetValue<int>(i),
                        Department = departmentColumn.Data.GetValue<string>(i),
                        Salary = salaryColumn.Data.GetValue<decimal>(i)
                    });
                }
            }
        }
    }

    return employees;
}

4. 批量插入到数据库

提供两种常用实现方案,适配不同数据量场景:

方案一:EF Core(中小数据量)

先定义DbContext:

public class AppDbContext : DbContext
{
    public DbSet<Employee> Employees { get; set; }

    protected override void OnConfiguring(DbContextOptionsBuilder optionsBuilder)
    {
        optionsBuilder.UseSqlServer("你的数据库连接字符串");
    }
}

执行批量插入:

public static async Task InsertEmployeesToDb(List<Employee> employees)
{
    using (var context = new AppDbContext())
    {
        context.Employees.AddRange(employees);
        await context.SaveChangesAsync();
    }
}

数据量超10万条时,可搭配EFCore.BulkExtensions提升插入性能

方案二:Dapper + SqlBulkCopy(大数据量)

适合超大规模数据的高效批量插入:

using Dapper;
using System.Data.SqlClient;
using System.Data;

public static async Task BulkInsertEmployees(string connectionString, List<Employee> employees)
{
    // 将实体列表转换为DataTable
    var dataTable = new DataTable();
    dataTable.Columns.Add("Id", typeof(int));
    dataTable.Columns.Add("Name", typeof(string));
    dataTable.Columns.Add("Age", typeof(int));
    dataTable.Columns.Add("Department", typeof(string));
    dataTable.Columns.Add("Salary", typeof(decimal));

    foreach (var emp in employees)
    {
        dataTable.Rows.Add(emp.Id, emp.Name, emp.Age, emp.Department, emp.Salary);
    }

    // 使用SqlBulkCopy批量写入
    using (var connection = new SqlConnection(connectionString))
    {
        await connection.OpenAsync();
        using (var bulkCopy = new SqlBulkCopy(connection))
        {
            bulkCopy.DestinationTableName = "Employees"; // 目标数据库表名
            // 列映射(DataTable列名与表字段名一致时可省略)
            bulkCopy.ColumnMappings.Add("Id", "Id");
            bulkCopy.ColumnMappings.Add("Name", "Name");
            bulkCopy.ColumnMappings.Add("Age", "Age");
            bulkCopy.ColumnMappings.Add("Department", "Department");
            bulkCopy.ColumnMappings.Add("Salary", "Salary");

            await bulkCopy.WriteToServerAsync(dataTable);
        }
    }
}

注意事项

  • 确保Parquet列的数据类型与数据库表字段类型兼容,避免类型转换错误
  • 处理超大文件时,建议分批次读取插入,防止内存溢出
  • 确保应用程序拥有目标数据库的写入权限

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 05:15:24