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

Java结合Hadoop用ParquetReader读S3 Parquet文件报FileNotFoundException

Java读取S3存储Parquet文件排错与实现方案

问题复现

初始Hadoop实现代码(抛出FileNotFoundException)

public static void main(String[] args) throws IOException, URISyntaxException {
        Path path = new Path("s3", "batch-dev", "/aman/part-e52b.c000.snappy.parquet");
        Configuration conf = new Configuration();
        conf.set("fs.s3.awsAccessKeyId", "xxx");
        conf.set("fs.s3.awsSecretAccessKey", "xxxxx");
        InputFile file = HadoopInputFile.fromPath(path, conf);
        ParquetFileReader reader2 = ParquetFileReader.open(conf, path);
        
        //MessageType schema = reader2.getFooter().getFileMetaData().getSchema();
        //System.out.println(schema);
}

异常信息

Exception in thread "main" java.io.FileNotFoundException: s3://batch-dev/aman/part-e52b.c000.snappy.parquet: No such file or directory.
    at org.apache.hadoop.fs.s3.S3FileSystem.getFileStatus(S3FileSystem.java:334)
    at org.apache.parquet.hadoop.util.HadoopInputFile.fromPath(HadoopInputFile.java:39)
    at com.bidgely.cloud.core.cass.gb.S3GBRawDataHandler.main(S3GBRawDataHandler.java:505)

注意:当前使用s3协议而非s3a协议,不确定Hadoop是否支持s3协议。

S3客户端验证代码(可正常获取对象)

public static void main(String args[]) {
        AWSCredentials credentials = new BasicAWSCredentials("XXXXX", "XXXXX");
        AmazonS3 s3Client = AmazonS3ClientBuilder.standard().withRegion("us-west-2").withCredentials(new AWSStaticCredentialsProvider(credentials)).build();
        S3Object object = s3Client.getObject(new GetObjectRequest("batch-dev", "/aman/part-e52b.c000.snappy.parquet"));
        System.out.println(object.getObjectContent());
}

该代码可正常获取S3对象流,但返回的输入流无法直接解析Parquet数据。

问题根因

  • Hadoop旧版s3://协议(对应org.apache.hadoop.fs.s3.S3FileSystem)已被官方废弃,路径解析逻辑与S3原生API不兼容:代码中传入的对象键以/开头,S3实际存储的对象键为aman/part-e52b.c000.snappy.parquet(无前导斜杠),旧实现不会自动裁剪多余斜杠,直接查询会返回对象不存在。
  • Parquet为列式存储格式,读取时需要支持随机定位(seek)的输入流,S3 SDK直接返回的S3ObjectInputStream仅支持顺序读取,无法满足Parquet解析要求。

前置Maven依赖

<!-- Parquet 核心依赖 -->
<dependency>
    <groupId>org.apache.parquet</groupId>
    <artifactId>parquet-hadoop</artifactId>
    <version>1.12.3</version>
</dependency>
<!-- AWS S3 SDK 依赖 -->
<dependency>
    <groupId>com.amazonaws</groupId>
    <artifactId>aws-java-sdk-s3</artifactId>
    <version>1.12.400</version>
</dependency>
<!-- 仅Hadoop s3a方案需要引入,纯SDK方案不需要 -->
<!--
<dependency>
    <groupId>org.apache.hadoop</groupId>
    <artifactId>hadoop-aws</artifactId>
    <version>3.3.4</version>
</dependency>
-->

纯SDK方案仅需引入parquet-hadoop和aws-java-sdk-s3两个依赖,不需要引入整套Hadoop依赖,可避免大部分依赖版本冲突问题。

方案1:纯S3 SDK实现(无Hadoop依赖,推荐)

核心是实现Parquet的InputFile和SeekableInputStream接口,通过S3 Range GET请求实现随机读取能力,完全规避Hadoop文件系统的协议兼容问题。

import com.amazonaws.auth.AWSStaticCredentialsProvider;
import com.amazonaws.auth.BasicAWSCredentials;
import com.amazonaws.services.s3.AmazonS3;
import com.amazonaws.services.s3.AmazonS3ClientBuilder;
import com.amazonaws.services.s3.model.GetObjectRequest;
import com.amazonaws.services.s3.model.S3Object;
import org.apache.parquet.example.data.Group;
import org.apache.parquet.hadoop.ParquetReader;
import org.apache.parquet.hadoop.example.GroupReadSupport;
import org.apache.parquet.io.InputFile;
import org.apache.parquet.io.SeekableInputStream;
import java.io.IOException;
import java.io.InputStream;

public class S3ParquetReader {

    public static void main(String[] args) throws IOException {
        // 初始化S3客户端,替换为实际的AK、SK、区域
        BasicAWSCredentials credentials = new BasicAWSCredentials("你的AK", "你的SK");
        AmazonS3 s3Client = AmazonS3ClientBuilder.standard()
                .withRegion("us-west-2")
                .withCredentials(new AWSStaticCredentialsProvider(credentials))
                .build();

        String bucket = "batch-dev";
        // 注意:对象键不要加前导斜杠
        String key = "aman/part-e52b.c000.snappy.parquet";

        // 封装S3对象为Parquet支持的InputFile
        InputFile s3InputFile = new S3InputFile(s3Client, bucket, key);

        // 初始化ParquetReader读取数据
        try (ParquetReader<Group> reader = ParquetReader.builder(new GroupReadSupport(), s3InputFile).build()) {
            Group record;
            while ((record = reader.read()) != null) {
                // 替换为实际的业务处理逻辑
                System.out.println(record.toString());
            }
        }
    }

    // 自定义S3 InputFile实现
    static class S3InputFile implements InputFile {
        private final AmazonS3 s3Client;
        private final String bucket;
        private final String key;
        private final long contentLength;

        public S3InputFile(AmazonS3 s3Client, String bucket, String key) {
            this.s3Client = s3Client;
            this.bucket = bucket;
            this.key = key;
            this.contentLength = s3Client.getObjectMetadata(bucket, key).getContentLength();
        }

        @Override
        public long getLength() {
            return contentLength;
        }

        @Override
        public SeekableInputStream newStream() {
            return new S3SeekableInputStream(s3Client, bucket, key, contentLength);
        }
    }

    // 自定义支持seek的输入流,通过S3 range请求实现随机读取
    static class S3SeekableInputStream extends SeekableInputStream {
        private final AmazonS3 s3Client;
        private final String bucket;
        private final String key;
        private final long contentLength;
        private long currentPos = 0;
        private InputStream currentStream;

        public S3SeekableInputStream(AmazonS3 s3Client, String bucket, String key, long contentLength) {
            this.s3Client = s3Client;
            this.bucket = bucket;
            this.key = key;
            this.contentLength = contentLength;
            this.currentStream = openStream(0, contentLength - 1);
        }

        private InputStream openStream(long start, long end) {
            GetObjectRequest req = new GetObjectRequest(bucket, key)
                    .withRange(start, end);
            S3Object object = s3Client.getObject(req);
            return object.getObjectContent();
        }

        @Override
        public void seek(long newPos) throws IOException {
            if (newPos < 0 || newPos >= contentLength) {
                throw new IOException("无效的跳转位置: " + newPos);
            }
            if (currentStream != null) currentStream.close();
            currentPos = newPos;
            currentStream = openStream(newPos, contentLength - 1);
        }

        @Override
        public long getPos() {
            return currentPos;
        }

        @Override
        public void readFully(byte[] bytes, int off, int len) throws IOException {
            int read = 0;
            while (read < len) {
                int r = currentStream.read(bytes, off + read, len - read);
                if (r == -1) throw new IOException("流意外结束");
                read += r;
                currentPos += r;
            }
        }

        @Override
        public int read() throws IOException {
            int b = currentStream.read();
            if (b != -1) currentPos++;
            return b;
        }

        @Override
        public int read(byte[] b, int off, int len) throws IOException {
            int read = currentStream.read(b, off, len);
            if (read > 0) currentPos += read;
            return read;
        }

        @Override
        public void close() throws IOException {
            if (currentStream != null) currentStream.close();
        }
    }
}

方案2:修正Hadoop s3a协议实现(兼容现有Hadoop生态)

如果必须使用Hadoop文件系统API,需要将协议替换为官方维护的s3a://,修正配置项和路径格式。

import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.Path;
import org.apache.parquet.hadoop.ParquetFileReader;
import org.apache.parquet.hadoop.util.HadoopInputFile;
import org.apache.parquet.io.InputFile;
import java.io.IOException;

public class HadoopS3ParquetReader {
    public static void main(String[] args) throws IOException {
        // 替换为s3a协议,路径不要加前导斜杠
        Path path = new Path("s3a://batch-dev/aman/part-e52b.c000.snappy.parquet");
        Configuration conf = new Configuration();
        // s3a配置前缀为fs.s3a,不是旧版的fs.s3
        conf.set("fs.s3a.access.key", "你的AK");
        conf.set("fs.s3a.secret.key", "你的SK");
        conf.set("fs.s3a.endpoint", "s3.us-west-2.amazonaws.com");
        // 自建S3服务需要开启路径式访问,AWS公有云可省略
        // conf.set("fs.s3a.path.style.access", "true");
        
        InputFile file = HadoopInputFile.fromPath(path, conf);
        try (ParquetFileReader reader = ParquetFileReader.open(file)) {
            // 替换为实际的schema读取、数据解析逻辑
            System.out.println(reader.getFooter().getFileMetaData().getSchema());
        }
    }
}

关键注意事项

  • S3对象键不要加前导斜杠,所有S3 API访问都需要使用无前导斜杠的键名。
  • 旧版s3://协议不再维护,存在大量兼容、性能问题,禁止生产环境使用。
  • 读取Parquet必须使用支持seek的输入流,不可直接使用普通顺序输入流解析。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 15:15:33