ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Java读写Parquet完全指南:从列式存储原理到代码实践

Java读写Parquet完全指南:从列式存储原理到代码实践 1. 内容整体设计与思路拆解1.1 为什么要单独写一篇“Java读写Parquet”的笔记先说个现象。如果你现在去网上搜“Java读写Parquet”搜出来的结果五花八门有教SPARK的有贴Python/Pandas的还有一堆没头没尾的代码碎片甚至混进来FPGA的DDR读写、EEPROM的Verilog程序。可能你只是想搞清楚一个问题在纯Java项目里到底怎么把一个DataFrame或者记录列表落成Parquet文件又怎么读回来。我刚接触Parquet的时候也有同样的困惑。当时接手一个离线数仓模块上游数据源是HDFS上的Parquet文件下游业务方要求我们输出Parquet格式而不是以前用的CSV或JSON。CSV简单但遇到嵌套对象、日期时间类型就非常痛苦JSON灵活但没有强schema约束离线任务跑着跑着就字段漂移了而Parquet天然是“列式存储强类型schema”的格式压缩率高、扫描性能好和Hive/Spark等框架对接也顺畅。这篇笔记不是教你写一个完整的数仓而是聚焦一件事在Java开发环境下把Parquet的读写链路彻底跑通。我踩过的坑、验证过的配置、最后能稳定运行的代码都会完整写出来。无论你是做后端服务写数据接入程序还是维护离线ETL这篇文章都能直接当手册用。1.2 为什么要用Parquet列式存储的价值不止是“省空间”理解Parquet先理解“列式存储”到底解决了什么问题。传统行式存储比如CSV在读取时哪怕你只需要某一张表的“age”字段也要把整行读出来、解析、再过滤。数据量小的时候无所谓但一旦表一天产出几个GB查询侧就会明显感觉到IO浪费。Parquet按列组织数据块每一列的数据在物理上连续存放读取时可以做“列剪枝”只加载需要的列配合每列的统计信息最小值、最大值等还可以做“谓词下推”跳过完全不符合过滤条件的行组。说白了它就是用“物理布局的合理设计”来换取“查询效率和存储成本的优化”。再叠加压缩算法Snappy、Gzip、Zstd等同样一份数据Parquet的体积通常只有CSV的30%~50%。所以我现在看到有人还在用CSV存数仓中间表都会建议他看一眼Parquet。这个选择不是“赶时髦”而是实打实能让下游跑得更快、磁盘用得更省。1.3 适合谁来读这篇文章我也把这条写在前面好让读者判断要不要继续往下看。你已经有Java基础会用Maven/Gradle管理依赖但没碰过Parquet你在纯Java服务端动手写文件读写程序不想为了一个数据格式就引入全套Spark你在Hive/Spark周边做数据管道虽然平时用SQL居多但有需要手写Java代码操纵Parquet文件的时刻你面试时被问到“Parquet为什么快”希望从底层原理到代码实现都能说清。如果你属于以上任何一种这篇文章的内容会非常合适。我尽量把原理讲得简单把代码写得可以直接跑。2. 核心细节解析与实操要点2.1 Parquet的核心概念RowGroup、ColumnChunk、Page刚上手Parquet建议先记住三个概念RowGroup、ColumnChunk、Page。它们从大到小划分了Parquet的物理结构。RowGroup逻辑上是一批行的集合Parquet写入时按行分批落盘。每个RowGroup包含这一批每一列的ColumnChunk。RowGroup的大小直接影响读取时的并行度——太小则元数据膨胀太大则内存压力大。ColumnChunk一个RowGroup里某一列的所有数据合在一起就是一个ColumnChunk。它独立压缩、独立存储统计信息。PageColumnChunk内部继续切成Page这是Parquet编解码的最小单元。读取时会按Page加载并基于Page内部的encoding做解码。这个结构对写出高效代码的影响是写入时RowGroup大小要合适太小会让元数据比例变高文件膨胀太大则单个写任务占用的内存高。我常用的经验值是不低于64MB但也要配合实际数据量和底层存储块大小来调整。读取时则尽量利用“按需加载列”避免整文件扫描。2.2 Schema与Java类型的映射关系Parquet文件里的Schema用Thrift定义支持嵌套结构。Java代码里最常用的映射表大概是这样的Parquet类型Java类型Avro说明BOOLEANboolean/Boolean布尔INT32int/Integer32位整数INT64long/Long64位整数FLOATfloat/Float单精度DOUBLEdouble/Double双精度BYTE_ARRAYbyte[] / String二进制或UTF-8字符串INT96不推荐使用旧版Timestamp只有Hive历史数据里才见到TIMESTAMPINT64逻辑类型java.time.Instant/LocalDateTime新版本推荐很多新人把“INT96”当成Timestamp的默认写法其实这是老Hive序列化的历史遗留问题Spark/Hive早期版本写出的Timestamp经常是INT96。新项目里如果自己控制Schema尽量避免INT96推荐用INT64配合逻辑类型TIMESTAMP_MILLIS/TIMESTAMP_MICROS来存时间。否则跨框架读数据时时间字段经常出现几小时或日期错位的奇怪现象。2.3 时间字段隐藏的时区陷阱Parquet的标准里特意提到一个标志位isAdjustedToUTC。我读过的很多中文资料都没说明白这件事但它在跨集群读写时特别关键。简单说如果写入方比如Spark把时间先转成UTC再存到Parquet的INT64中isAdjustedToUTCtrue读取方读到这个值时应该再按本地时区转换回显示时间。如果写入方存的是本地时间、标志位又是false那读取方就不能再做一次时区切换否则就会出现“时间前后相差8小时”的经典问题。我在生产系统里遇到过一次——上游用Spark的to_utc_timestamp处理时间下游Java程序直接把这个INT64当成“自纪元以来的毫秒数”塞进new Date(value)结果展示出来的时间比真实值早了8小时。排查了大半天最后看Parquet元数据里的isAdjustedToUTC标志才发现问题。所以你在写Java读取逻辑时要明确约定要么全链路都统一存UTC时间戳读取时一次性转为本地时区要么全链路都存已转换好的字符串时间或本地时间戳千万不要混用。这是一个兼容性约定不是某一个库能独自解决的。2.4 工具选型原生Parquet库还是Spark这是“Java里读写Parquet”最常见的一个分叉口。如果你只是在Java进程里读写本地/HDFS上的Parquet文件优先选择Parquet原生Java库parquet-hadoopparquet-avro不需要Spark那套庞大的依赖。它足够底层可控性强适合写独立的ETL程序、数据接入服务。如果你本来就在Spark作业里运行那没必要绕过Spark自己操作Parquet直接df.write().parquet()最方便。Spark内部也是调Parquet的原生库但它替你处理了时间类型、压缩配置、动态分区的很多细节。第三种情况需要提一下如果你已经有Avro/Protobuf对象想让Parquet文件带上完整的schema信息那么用parquet-avro库天然契合因为Parquet对Avro的支持最成熟。如果只有普通Java Bean就得手工构造MessageType或者让Avro负责转换。我的建议是场景推荐方案Java服务端/离线脚本读写本地文件parquet-hadoopparquet-avro用Avro schema定义结构数据量大、字段多、需要分区Spark SQL/DataFrame API 写入已有现成的Avro POJO不想重复造schemaparquet-avro的AvroParquetWriter只需要快速读列、不想引入POJOparquet-hadoop的ParquetReaderGroup另外如果你的项目已经用了Arrow生态也可以考虑arrow-vector读取Parquet但在Java场景下它比Avro路径更复杂除非你有向量化内存计算的需求否则不必一开始就用。3. 实操过程与核心环节实现3.1 环境准备与依赖引入我先用Maven项目做示例。JDK用11或17都行我这边实测JDK 8也能跑但建议至少8以上因为部分较新的Parquet版本开始要求JDK 8且时间API用java.time会更舒服。核心依赖如下这是我实际验证过的版本组合properties parquet.version1.14.1/parquet.version hadoop.version3.3.6/hadoop.version /properties dependencies dependency groupIdorg.apache.parquet/groupId artifactIdparquet-hadoop/artifactId version${parquet.version}/version /dependency dependency groupIdorg.apache.parquet/groupId artifactIdparquet-avro/artifactId version${parquet.version}/version /dependency dependency groupIdorg.apache.hadoop/groupId artifactIdhadoop-common/artifactId version${hadoop.version}/version /dependency dependency groupIdorg.apache.hadoop/groupId artifactIdhadoop-mapreduce-client-core/artifactId version${hadoop.version}/version /dependency /dependencies为什么需要hadoop-*依赖因为Parquet底层把文件系统抽象成了Hadoop的FileSystemAPI即使你写的是本地文件它也会通过Hadoop的LocalFileSystem来访问。没有这个依赖类加载阶段就会因为找不到org.apache.hadoop.fs.Path而直接报错。这个坑几乎是所有新手都会踩的。如果你只处理本地文件不需要hadoop-mapreduce-client-core但hadoop-common和hadoop-hdfs如果走HDFS是少不了的。3.2 用Avro Parquet完成第一次写入先定义一个Avro schema或者直接用Java类。我习惯先写一个Avro schema的JSON文件放在src/main/resources/avro/UserRecord.avsc{ type: record, name: UserRecord, namespace: com.example.parquet, fields: [ {name: id, type: long}, {name: name, type: string}, {name: age, type: int}, {name: email, type: [null, string], default: null}, {name: createdAt, type: {type: long, logicalType: timestamp-millis}} ] }这里有个细节createdAt用long类型 logicalType: timestamp-millis这样Parquet内部会用INT64存储并记录它是“毫秒精度的时间戳”。比用INT96干净得多。然后写工具类分三步完成写入加载Schema、构造数据、调Writer。import org.apache.avro.Schema; import org.apache.avro.generic.GenericData; import org.apache.avro.generic.GenericRecord; import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.Path; import org.apache.parquet.avro.AvroParquetWriter; import org.apache.parquet.hadoop.ParquetWriter; import org.apache.parquet.hadoop.metadata.CompressionCodecName; import org.apache.parquet.hadoop.util.HadoopOutputFile; import java.io.InputStream; import java.nio.file.Files; import java.nio.file.Paths; import java.util.ArrayList; import java.util.List; public class ParquetWriterDemo { public static void main(String[] args) throws Exception { // 1. 读取并解析 Avro Schema Schema schema; try (InputStream in Files.newInputStream(Paths.get(src/main/resources/avro/UserRecord.avsc))) { schema new Schema.Parser().parse(in); } // 2. 生成示例数据 ListGenericRecord records new ArrayList(); for (long i 1; i 1000; i) { GenericRecord record new GenericData.Record(schema); record.put(id, i); record.put(name, user_ i); record.put(age, (int) (20 (i % 40))); record.put(email, user i example.com); record.put(createdAt, System.currentTimeMillis()); records.add(record); } // 3. 写入 Parquet 文件 Path outputPath new Path(./output/users.parquet); try (ParquetWriterGenericRecord writer AvroParquetWriter .GenericRecordbuilder(HadoopOutputFile.fromPath(outputPath, new Configuration())) .withSchema(schema) .withCompressionCodec(CompressionCodecName.SNAPPY) .withRowGroupSize(64 * 1024 * 1024L) // 64MB .withPageSize(1024 * 1024) // 1MB .build()) { for (GenericRecord record : records) { writer.write(record); } } System.out.println(写入完成文件大小: Files.size(Paths.get(./output/users.parquet))); } }这里解释几个参数的选择CompressionCodecName.SNAPPYSnappy压缩速度快、CPU开销小适合大多数离线场景。想要更高压缩比可以用Gzip或Zstd但写入耗时会上涨。注意如果下游是Hive/SparkSnappy兼容性最省心Gzip次之Zstd需要版本匹配。withRowGroupSize(64MB)RowGroup太大写入时内存压力大太小则元数据膨胀。我一般从64MB起步数据量大的场景再调大到128MB或256MB。withPageSize(1MB)PageSize影响读取时的随机IO粒度。默认值通常够用不必频繁调整。3.3 用ParquetReader读回数据读完再读用ParquetReader配合GenericRecord读取。因为我们的schema是Avro的所以读取侧也用Avro的解析方式这样能自动完成类型转换。import org.apache.avro.generic.GenericRecord; import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.Path; import org.apache.parquet.avro.AvroParquetReader; import org.apache.parquet.hadoop.ParquetReader; public class ParquetReaderDemo { public static void main(String[] args) throws Exception { Path path new Path(./output/users.parquet); Configuration conf new Configuration(); try (ParquetReaderGenericRecord reader AvroParquetReader .GenericRecordbuilder(path) .withConf(conf) .build()) { GenericRecord record; int count 0; while ((record reader.read()) ! null) { if (count 3) { System.out.println(record); } count; } System.out.println(读取记录数: count); } } }读取端的关键是尽量别在Java里手动逐字段转换让Avro和Parquet配合自动生成类型映射。如果你读出的field是Utf8类型而不是String也不用慌调用record.get(name).toString()即可取得字符串。这是Avro内部对字符串的表示方式。3.4 用原生MessageType读任意Parquet文件有时候你读的Parquet文件不是Avro生成也没有对应的Avro Schema只有Parquet自带的schema。这时就不能用AvroParquetReader可以用ParquetReaderGroup配合SimpleGroup读取。import org.apache.parquet.example.data.Group; import org.apache.parquet.example.data.simple.SimpleGroupFactory; import org.apache.parquet.hadoop.ParquetReader; import org.apache.parquet.schema.MessageType; import org.apache.parquet.schema.MessageTypeParser; import org.apache.hadoop.fs.Path; public class RawParquetReaderDemo { public static void main(String[] args) throws Exception { Path path new Path(./output/users.parquet); MessageType schema MessageTypeParser.parseMessageType( message UserRecord { required int64 id; required binary name (UTF8); required int32 age; optional binary email (UTF8); required int64 createdAt (TIMESTAMP_MILLIS); } ); try (ParquetReaderGroup reader ParquetReader .Groupbuilder(new HadoopInputFile(path, new Configuration())) .build()) { Group group; while ((group reader.read()) ! null) { long id group.getLong(0, 0); String name group.getBinary(1, 0).toStringUsingUTF8(); int age group.getInteger(2, 0); String email group.getBinary(3, 0).toStringUsingUTF8(); long createdAt group.getLong(4, 0); System.out.println(id , name , age , email , createdAt); } } } }这里有一个高频坑binary (UTF8)在Parquet内部的物理类型是BYTE_ARRAY但读取Java API的时候不能直接当String用必须先拿Binary实例再调用toStringUsingUTF8()转换。之前看到有人直接用group.getString(...)编译不报错但运行时报ClassCastException。另外如果字段是optional可空取值前最好用group.getFieldRepetitionCount(index)判断一下是否有值避免越界。3.5 直接对接Spark DataFrame一种“偷懒但好用”的方案如果手头有Spark环境并且数据量已经大到单机内存扛不住那就没有必要写底层API。直接用Spark读写Parquet代码量最少而且自动处理schema演化。// 写 dataFrame .write .mode(overwrite) .option(compression, snappy) .partitionBy(dt) .parquet(/path/to/output) // 读 val df spark.read.parquet(/path/to/output) df.filter($age 18).select(id, name).show()这段代码背后的原理还是我们前文提到的那套Spark把DataFrame的schema映射成Parquet schema用ParquetWriter落盘。但它帮你把复杂的RowGroup大小、PageSize、时间戳时区等参数都做了合理默认。所以如果你不需要在业务代码里精细控制底层行为能走Spark就尽量走Spark别自己重复造轮子。3.6 写入分区数据从单文件走向多文件目录实际生产里很少只有一个Parquet文件。更多时候是按日期、按业务维度做分区目录结构类似/path/to/table/dt2025-01-01/part-00001.parquet /path/to/table/dt2025-01-02/part-00002.parquet用Java原生API在代码里实现分区也很直接先根据记录的dt字段把数据分组每个分组写到对应的目录路径下目录名用dt2025-01-01这种Hive风格。需要注意两点目录名必须遵循“列名值”的格式下游Hive/Spark才能自动识别分区列。不要用字符拼接日期路径无法避免目录层次错误先构造Path再拼接这才是安全且可维护的写法。MapString, ListGenericRecord groups records.stream() .collect(Collectors.groupingBy(r - r.get(dt).toString())); for (Map.EntryString, ListGenericRecord entry : groups.entrySet()) { Path partitionPath new Path(./output/table/dt entry.getKey()); // 继续用 AvroParquetWriter将 Path 改为 partitionPath 下的文件名 }这块不算复杂但面试时经常被问到“Parquet分区与Hive分区有什么关系”。简单回答是Hive分区就是在表的目录下建“列名值”的子目录Parquet本身不强制要求分区分区是文件系统层面的组织方式两者协同实现数据裁剪。4. 常见问题与排查技巧实录4.1 依赖冲突NoClassDefFoundError、jackson版本冲突Java生态里只要粘上Hadoop相关依赖依赖冲突就是一个绕不开的话题。Parquet引用了jackson-mapper-asl、snappy-java、commons-compress等一堆库如果你的项目里其他组件比如Spring Boot自带Jackson 2.x也用到这些库就很容易出现不同版本同时存在的“俄罗斯套娃”问题。排查思路很简单用mvn dependency:tree看看实际的依赖树把冲突的版本理出来在pom.xml里用exclusion排除不需要的旧版本或者用dependencyManagement统一版本如果项目里已经有Spring Boot建议用它的spring-boot-starter-parent作为parent版本仲裁能省不少事。我踩过最典型的坑是引入parquet-hadoop后原来正常的ObjectMapper序列化突然报错最后发现是Hadoop的一个传递依赖把jackson-databind版本拽低了。遇到这种情况别慌直接在POM里显式声明一个更高版本的jackson-databind即可。4.2 写入本地文件失败AccessDeniedException / C0602有人写Parquet到本地路径时遇到类似这样的报错java.io.IOException: Cannot create file ... (AccessDeniedException)多数情况不是权限问题而是Hadoop的LocalFileSystem对路径的处理与Java原生的Files API不一致。它会把./output/users.parquet解析成file:/当前工作目录/output/users.parquet如果当前目录没有写权限或者路径解析出现问题就会报错。解决方案有两种先调用Files.createDirectories(Paths.get(output))确保目录存在使用Path outputPath new Path(new java.io.File(output/users.parquet).getAbsolutePath());拼成绝对路径避免相对路径歧义。另外如果用Windows开发要特别注意Hadoop对Windows本地文件系统支持不完善最常见的是缺少winutils.exe。解决方法是下载对应Hadoop版本的winutils.exe放到某个路径然后设置环境变量HADOOP_HOMED:\hadoop PATH%HADOOP_HOME%\bin;%PATH%或者在代码里早期就把hadoop.home.dir属性设好System.setProperty(hadoop.home.dir, D:\\hadoop);4.3 类型转换异常expected INT32 but got LONG这种事非常典型。上游Spark/Hive表里的字段在底层可能是INT64Long但下游Java代码里按int去取就会得到org.apache.parquet.io.ParquetDecodingException: Can not read value at 0 in field age排查思路先用ParquetFileReader.readFooter()读一下文件的schema确认物理类型到底是什么再决定Java侧用getInteger还是getLong。一般情况下跨框架读取时强类型转换一定要以文件schema为准不要以业务约定为准。4.4 压缩算法不兼容Can not read value at ... due to corrupt stream如果Parquet文件用了LZO压缩而下游环境没有引入对应的hadoop-lzo库读取时会报类似java.io.IOException: invalid stream header或Corrupt错误。排查技巧检查写入时用的是哪种压缩parquet-tools命令行或ParquetFileReader可以读元数据确认读端的类库里包含了对应的解压器。Snappy兼容性最好Zstd其次。如果团队跨多个版本框架协作我建议统一用Snappy不要为了压缩率冒险换小众压缩算法。压缩率虽然重要但兼容性和稳定性永远排在第一位。4.5 时间字段比实际少8小时这个前文已经点过名。再给一个具体的排查路径用parquet-tools meta查看Parquet文件内时间字段的logicalType和isAdjustedToUTC对比写入端的时区设置比如Spark的spark.sql.session.timeZone统一约定所有时间字段要么存毫秒时间戳读写双方都做时区转换要么存字符串无时区歧义。我最终在生产规范里做了两条硬性要求第一Parquet中时间全部用TIMESTAMP_MILLIS存储UTC毫秒第二读取展示时统一转回北京时间。从此再没因为这个吵过架。4.6 常见问题速查表现象核心原因解法找不到org.apache.hadoop.fs.Path缺Hadoop Common引入hadoop-common依赖AccessDeniedException写本地文件相对路径/HDFS权限先创建目录或改用绝对路径NoClassDefFoundError: org/xerial/snappy/SnappyInputStreamSnappy版本冲突检查依赖树统一snappy-java版本读取时类型不匹配schema与实际Parquet类型不同用ParquetFileReader查看footer真实schema时间差8小时isAdjustedToUTC约定不一致全链路统一存UTC毫秒无法读取LZO压缩文件缺hadoop-lzo库换Snappy或引入对应压缩库Windows下本地读写失败缺winutils.exe配置HADOOP_HOME4.7 书写性能调优经验如果写数据量很大比如一次写几百万条可以做三件事使用多线程并行写多个文件。Parquet单文件写入不是CPU密集型的瓶颈通常在压缩和IO并行度建议控制在核心数的1~2倍。每个线程写独立文件最后再合并目录。合理设置withRowGroupSize和withPageSize。RowGroup Size如果太小比如默认1MB会生成海量的小ColumnChunk元数据开销大、读性能差。如果你只做“一次性归档”可以用更小的RowGroup来降低写内存如果主要用于查询扫描尽量用大RowGroup。注意GC压力。ParquetWriter内部会缓冲一个RowGroup的数据如果你的RowGroup设置太大比如1GB而堆内存又不够就会频繁Full GC。先把RowGroupSize和堆内存匹配好再看实际写入速度。读性能优化的核心则在于“只读需要的列”尤其在你用ParquetReader遍历数据的时候别每个字段都get一遍用不到的列就别取。Parquet的列剪枝在底层是自动跳过的但你代码里如果不注意把不需要的字段都取出来依然会造成额外开销。5. 面试场景下的Parquet知识点清单因为关于“java读写parquet”的热搜词里混着“java面试题”之类的内容我也顺手整理一组高频面试问答。这部分不是凑字数而是帮你把学习到的实践往“能讲清原理”上拔高。5.1 为什么Parquet比CSV/JSON更适合数据分析从存储结构上回答分三点。第一列式存储天然压缩率高相同数据量下体积更小第二查询时可以列剪枝只读取目标列减少IO第三每列自带统计信息配合谓词下推可以跳过大量无关数据块。反观CSV整行读取IO占用大JSON又缺少强schema解析开销更高。5.2 Parquet的RowGroup和Page分别是什么RowGroup是若干行的集合是Parquet读写的基本批次单位。Page是ColumnChunk内部更细粒度的编码单元。RowGroup越大单批次读出的有效数据越多越适合顺序扫描Page大小影响随机读取效率。了解这两个层级就能解释为什么Parquet在大数据场景下优于普通文件。5.3 写Parquet时如何选择压缩算法具体场景具体分析。Snappy速度最快CPU开销最低适合写入频繁的实时/准实时链路Gzip压缩率更高适合长期归档Zstd是近几年的新选择压缩率高且速度不差但要注意框架版本支持。兼容性优先级Snappy Gzip Zstd LZO。5.4 Spark写Parquet和原生Java写Parquet有什么不同本质是一样的底层都是Parquet格式的Writer。区别是Spark屏蔽了很多细节并且帮你处理了分区、动态schema、时间时区、并行任务调度等场景。原生Java API更底层、更灵活可以嵌入到服务端程序里但不具备分布式task调度能力。面试时说清楚“两者底层格式相同只是使用场景和抽象层级不同”基本就能过关。5.5 如果读一个Parquet文件发现schema和预期不一样怎么办这个问题很实际。先承认“Parquet是有schema的文件格式读之前应该先确认schema”然后说操作用ParquetFileReader.readFooter读取底层schema对比预期如果只是字段顺序不同重新建一个映射转换如果是物理类型不同做显式Cast如果是字段缺失检查写入时的schema演化策略。这个回答展示的是“有具体排查经验”而不是“背书”。6. 最后再分享几个我自己一直在用的小经验第一点如果你只是偶尔写个脚本处理Parquet直接用SPARK SQL最省心不要为了“纯粹”而硬上原生API。原生API适合嵌入业务系统、适合跟别人对接依赖、适合你在面试时展示底层理解但日常数据分析Spark已经替你处理了九成麻烦。第二点Parquet文件的“血缘”信息是隐含在schema里的如果你们团队没有统一的schema管理方案建议在写Parquet之前就把Avro Schema文件纳入版本库。这样以后不管谁接手都能从Schema文件看出这个Parquet文件应该长什么样。否则文件越堆越多最后变成“谁都不敢动”的黑洞。第三点存储和压缩参数不要照搬网上的默认值。我见过有人从网上复制withRowGroupSize(128L * 1024 * 1024)在本地小文件上跑结果写几万条数据就OOM。参数永远要配合数据和内存来调整。先测一条数据的大小再估算RowGroup占用的内存比凭感觉调参靠谱得多。第四点时间类型的坑比想象中多。我强烈建议你在团队的接口文档里明确写上“时间字段在Parquet文件里如何存储”并且加一个单元测试专门验证“写入-读出”后时间字段和原值严格一致。这个测试成本极低但能省掉的排查时间远超你的预期。我个人在实际操作中的体会是Parquet本身不难难的是它夹在Hadoop、Avro、Spark这些生态之间任何一个“自以为是的默认值”都可能让你踩坑。这篇笔记不是唯一答案但我把那些坑都标了出来。如果你照着代码跑通了恭喜你Java读写Parquet这个能力你已经真正拿到手了。后面再遇到Hive、Spark或数据湖相关的问题你会发现很多概念在这里已经打通了。
返回列表