Spark作业切换至AWS S3存储时出现IllegalArgumentException错误的解决咨询
这个错误的核心原因很明确:你的代码中部分FileSystem调用是直接获取了集群默认的HDFS文件系统实例,但当处理S3路径时,需要根据路径本身的scheme(也就是s3://)动态获取对应的S3文件系统客户端,而不是一直用HDFS的。
我们逐个看你列出的几个调用:
1. 自定义FileInputFormat的getSplits方法
FileSystem fs = path.getFileSystem(job.getConfiguration());
这个调用是正确的,不需要修改。因为path.getFileSystem()会根据传入的path的scheme自动匹配对应的文件系统(不管是HDFS还是S3),所以这个地方没问题。
2. 自定义RecordReader中的调用
final FileSystem fs = file.getFileSystem(job);
同样,这个也是正确的逻辑,file是一个Path对象,调用它的getFileSystem会自动适配对应的文件系统,不用改动。
3. 自行编写的RecordReader的nextKeyValue()方法
FileSystem fs = FileSystem.get(jc);
这里就是问题所在了!FileSystem.get(jc)会直接返回配置中默认的文件系统(你的集群默认是HDFS),当你在这个方法里处理S3路径时,就会出现"Wrong FS"的错误。
修改方案:你需要根据当前操作的具体Path来获取文件系统。比如如果你是在处理某个文件路径targetPath,就改成:
FileSystem fs = targetPath.getFileSystem(jc);
这样就会根据targetPath的s3://前缀自动获取S3的文件系统客户端,而不是默认的HDFS。
4. 检测文件夹文件数量的Scala代码
val fs = FileSystem.get(sc.hadoopConfiguration); val status = fs.listStatus(new Path(path))
这里和上面的问题一样,FileSystem.get(sc.hadoopConfiguration)拿到的是默认HDFS实例,当path是S3路径时就会报错。
修改成根据路径动态获取:
val targetPath = new Path(path) val fs = targetPath.getFileSystem(sc.hadoopConfiguration) val status = fs.listStatus(targetPath)
这样fs就会是对应S3路径的文件系统客户端,能正确处理s3://开头的路径。
额外注意事项
在EMR集群上,确保你的Hadoop配置已经正确配置了S3相关的实现类(EMR通常会默认配置好,比如fs.s3a.impl指向org.apache.hadoop.fs.s3a.S3AFileSystem),不过只要你的代码里是根据路径动态获取FS,一般不需要额外手动配置这个,EMR的环境会自动处理。
内容的提问来源于stack exchange,提问作者osk

