如何使用C#读取Parquet文件并将数据存入数据库表?
.NET实现Parquet文件数据插入数据库表
1. 安装依赖包
需要两个核心类库:
Parquet.Net:用于解析读取Parquet文件- 根据数据库类型选择数据访问类库,例如:
- SQL Server:
Microsoft.EntityFrameworkCore.SqlServer或Dapper - MySQL:
MySql.EntityFrameworkCore或Dapper
- SQL Server:
通过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
相关产品推荐
相关产品推荐

