如何使用Akka Streams或Alpakka读取S3中的Parquet文件并解决s3a报错
错误根因
java.io.IOException: No FileSystem for scheme: s3a 报错本质是Hadoop核心包默认未内置S3协议的文件系统实现,无法识别s3a前缀的路径,需要补充对应依赖并在Hadoop配置中显式注册实现类。
修复步骤
1. 补充依赖
首先需要引入和你环境Hadoop版本完全匹配的hadoop-aws依赖,以sbt配置为例:
libraryDependencies ++= Seq( // 请替换为和你项目中Hadoop版本完全一致的版本号 "org.apache.hadoop" % "hadoop-common" % "3.3.4", "org.apache.hadoop" % "hadoop-aws" % "3.3.4" )
如果使用Maven,对应调整pom.xml的依赖声明即可。
2. 修正Hadoop配置
在原有配置基础上新增s3a文件系统实现和凭证链配置:
val conf: Configuration = new Configuration() // 注册s3a协议对应的文件系统实现类 conf.set("fs.s3a.impl", "org.apache.hadoop.fs.s3a.S3AFileSystem") // 配置AWS凭证读取链,支持从环境变量、本地配置文件、EC2实例角色等位置自动读取凭证 conf.set("fs.s3a.aws.credentials.provider", "com.amazonaws.auth.DefaultAWSCredentialsProviderChain") // 保留原有Avro兼容配置 conf.setBoolean(AvroReadSupport.AVRO_COMPATIBILITY, true)
如果使用兼容S3的私有存储(如MinIO),可额外添加以下配置:
conf.set("fs.s3a.endpoint", "http://你的私有存储地址:端口") conf.set("fs.s3a.path.style.access", "true")
3. 对接Akka Streams完整示例
可以将ParquetReader封装为Akka Stream的Source,实现流式读取:
import akka.NotUsed import akka.stream.scaladsl.Source import org.apache.parquet.hadoop.ParquetReader import org.apache.avro.generic.GenericRecord import org.apache.parquet.avro.AvroParquetReader import org.apache.parquet.hadoop.util.HadoopInputFile import org.apache.hadoop.fs.Path import org.apache.hadoop.conf.Configuration // 将ParquetReader封装为Akka Source,自动关闭资源 def parquetSource(reader: ParquetReader[GenericRecord]): Source[GenericRecord, NotUsed] = { Source.unfold(reader) { reader => Option(reader.read()) match { case Some(record) => Some((reader, record)) case None => reader.close() None } } } // 调用逻辑 val path = "s3a://bucketName/path/to/foo/part-00000-656418ee-7cc0-42ee-93e-aaa69ee6f916.c000.snappy.parquet" val conf: Configuration = new Configuration() conf.set("fs.s3a.impl", "org.apache.hadoop.fs.s3a.S3AFileSystem") conf.set("fs.s3a.aws.credentials.provider", "com.amazonaws.auth.DefaultAWSCredentialsProviderChain") conf.setBoolean(AvroReadSupport.AVRO_COMPATIBILITY, true) val inputFile = HadoopInputFile.fromPath(new Path(path), conf) val reader: ParquetReader[GenericRecord] = AvroParquetReader.builder[GenericRecord](inputFile).withConf(conf).build() // 运行流处理数据 parquetSource(reader) .runForeach(record => println(record))
注意事项
- 必须保证
hadoop-common和hadoop-aws的版本完全一致,否则会出现类不兼容等异常 - 生产环境建议给流添加异常捕获逻辑,确保出现读取错误时也能正确释放ParquetReader和底层IO资源
内容的提问来源于stack exchange,提问作者igx
相关产品推荐
相关产品推荐

