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
相关产品推荐
相关产品推荐

