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

Spring Batch是否支持Parquet文件读写?Java读取Parquet遇阻

Spring Batch 对 Parquet 文件读写的支持方案

Spring Batch本身并没有内置专门用于Parquet文件读写的组件,但它提供了灵活的扩展机制,允许你基于Apache Parquet库自定义ItemReader和ItemWriter来实现Parquet文件的读写操作——这刚好适配你已经尝试过Apache Parquet库的情况。

具体实现思路

1. 自定义Parquet ItemReader

基于Apache Parquet的ParquetReader封装Spring Batch的ItemReader,负责读取Parquet文件并将记录映射为业务POJO:

import org.apache.parquet.hadoop.ParquetReader;
import org.apache.parquet.hadoop.api.ReadSupport;
import org.springframework.batch.item.ItemReader;
import org.springframework.batch.item.NonTransientResourceException;
import org.springframework.batch.item.ParseException;
import org.springframework.batch.item.UnexpectedInputException;

import java.io.IOException;

public class ParquetItemReader<T> implements ItemReader<T> {

    private final ParquetReader<T> parquetReader;

    public ParquetItemReader(String filePath, ReadSupport<T> readSupport) throws IOException {
        this.parquetReader = ParquetReader.builder(readSupport, new org.apache.hadoop.fs.Path(filePath)).build();
    }

    @Override
    public T read() throws Exception, UnexpectedInputException, ParseException, NonTransientResourceException {
        return parquetReader.read();
    }

    // 确保资源关闭
    public void close() throws IOException {
        if (parquetReader != null) {
            parquetReader.close();
        }
    }
}

如果需要将Parquet记录映射为POJO,可以使用AvroReadSupport或者自定义ReadSupport实现Schema到POJO的映射。

2. 自定义Parquet ItemWriter

同样基于Apache Parquet的ParquetWriter封装ItemWriter,负责将业务POJO写入Parquet文件:

import org.apache.parquet.hadoop.ParquetWriter;
import org.apache.parquet.hadoop.api.WriteSupport;
import org.springframework.batch.item.ItemWriter;

import java.io.IOException;
import java.util.List;

public class ParquetItemWriter<T> implements ItemWriter<T> {

    private final ParquetWriter<T> parquetWriter;

    public ParquetItemWriter(String filePath, WriteSupport<T> writeSupport) throws IOException {
        this.parquetWriter = ParquetWriter.builder(writeSupport, new org.apache.hadoop.fs.Path(filePath)).build();
    }

    @Override
    public void write(List<? extends T> items) throws Exception {
        for (T item : items) {
            parquetWriter.write(item);
        }
    }

    // 确保资源关闭并完成写入
    public void close() throws IOException {
        if (parquetWriter != null) {
            parquetWriter.close();
        }
    }
}

3. 依赖配置

需要在项目中引入Apache Parquet相关依赖(以Maven为例):

<dependency>
    <groupId>org.apache.parquet</groupId>
    <artifactId>parquet-hadoop</artifactId>
    <version>1.14.2</version>
</dependency>
<!-- 如果使用Avro映射POJO,需引入 -->
<dependency>
    <groupId>org.apache.parquet</groupId>
    <artifactId>parquet-avro</artifactId>
    <version>1.14.2</version>
</dependency>

关键注意事项

  • Schema定义:Parquet文件依赖Schema,你需要提前定义好对应业务数据的Parquet Schema,或通过Avro等工具从POJO生成Schema。
  • 资源管理:在Spring Batch的Step结束后,务必调用自定义Reader/Writer的close()方法,确保Parquet流正确关闭,避免资源泄漏。
  • 数据类型映射:注意Java类型与Parquet数据类型的对应关系,比如Integer对应Parquet的INT32,String对应UTF8等。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 23:53:09