Flink 1.13.5自托管集群读取S3 ORC文件报错求助
NoClassDefFoundError 环境与代码
使用Flink 1.13.5自托管集群,尝试读取AWS S3上的ORC文件,代码如下:
import org.apache.flink.streaming.api.scala._ import org.apache.flink.streaming.api.functions.source.FileProcessingMode import org.apache.hadoop.conf.Configuration import java.util.concurrent.TimeUnit import java.time.LocalDateTime import org.apache.flink.core.fs.Path import org.apache.flink.orc.OrcRowInputFormat val en = StreamExecutionEnvironment.getExecutionEnvironment val config = new Configuration config.setInt("fs.s3.connection.maximum", 15000) config.setString("fs.s3.readahead.range", "1M") val Schema = "struct<empId:int,name:string,salary:double>" val orcRowInput = new OrcRowInputFormat("s3a://NotApplicable", Schema, config) val fileRead = en.readFile(orcRowInput, "s3a://bucket/folder/", FileProcessingMode.PROCESS_CONTINUOUSLY, TimeUnit.HOURS.toMillis(24)).name("s3a://bucket/folder/").uid("s3a://bucket/folder/") env.execute("Flink APP")
运行异常
应用以Jar包形式运行时,出现如下错误:
RUNNING to FAILED with failure cause: java.lang.NoClassDefFoundError: org/apache/hadoop/fs/FileSystem
at java.lang.ClassLoader.defineClass(Native Method)
当前配置与依赖
Fat Jar包含依赖:
- flink-orc_2.11 1.13.5
- hadoop-hdfs-client 3.3.6
- hadoop-client 3.3.6
- hadoop-common 3.3.6
- hadoop-aws 3.3.6
- flink-s3-fs-hadoop 1.13.5
- 其他Flink 1.13.5版本依赖
集群插件Jar:
/opt/flink/plugins/s3-fs-hadoop/flink-s3-fs-hadoop-1.13.5.jar/opt/flink/plugins/s3-fs-presto/flink-s3-fs-presto-1.13.5.jar
集群状态:未安装完整Hadoop包(无
/etc/hadoop/conf目录),未设置HADOOP_HOME环境变量。
疑问
读取S3上的ORC文件是否必须安装Hadoop?或是遗漏了其他配置?
是否必须安装完整Hadoop包?
不需要安装完整的Hadoop发行版,但Flink的ORC读取和S3文件系统依赖Hadoop的核心API,因此需要确保这些API的依赖能被Flink正确加载。
问题原因分析
- 类加载冲突:Fat Jar中打包的Hadoop依赖(如hadoop-common、hadoop-client)与Flink插件目录中的
s3-fs-hadoop插件依赖的Hadoop类可能存在版本冲突,或者类加载器隔离导致org.apache.hadoop.fs.FileSystem无法被正确找到。Flink的插件类加载器与应用类加载器是隔离的,重复打包依赖容易引发类加载问题。 - Hadoop配置缺失:虽然未安装Hadoop,但需要显式配置S3文件系统的实现类,否则Flink无法正确初始化S3文件系统的Hadoop API。
具体解决步骤
调整依赖打包策略
将Hadoop相关依赖(hadoop-client、hadoop-common、hadoop-aws、hadoop-hdfs-client)设置为providedscope(Maven/Gradle),避免打包进Fat Jar。因为Flink的s3-fs-hadoop插件已经依赖这些类,重复打包会导致类加载冲突。配置Flink的S3文件系统
在集群所有节点的flink-conf.yaml中添加以下配置:# 指定S3文件系统实现 fs.s3a.impl: org.apache.flink.fs.s3hadoop.S3HadoopFileSystem # AWS认证信息(根据实际情况填写) fs.s3a.access.key: YOUR_AWS_ACCESS_KEY fs.s3a.secret.key: YOUR_AWS_SECRET_KEY # S3端点(根据你的AWS区域调整,如us-east-1为s3.amazonaws.com) fs.s3a.endpoint: s3.amazonaws.com调整类加载顺序
如果仍存在类加载问题,在flink-conf.yaml中设置类加载优先级为父类加载器优先:classloader.resolve-order: parent-first确保Hadoop的核心类由Flink的系统类加载器加载,而非应用类加载器。
修正代码中的Configuration初始化
显式在代码的Configuration对象中指定S3文件系统实现,避免依赖默认配置:val config = new Configuration() // 添加S3文件系统实现配置 config.set("fs.s3a.impl", "org.apache.flink.fs.s3hadoop.S3HadoopFileSystem") config.setInt("fs.s3.connection.maximum", 15000) config.setString("fs.s3.readahead.range", "1M")清理冗余插件
如果不使用Presto的S3文件系统实现,删除集群节点上的/opt/flink/plugins/s3-fs-presto/目录下的Jar包,避免不必要的类加载冲突。
内容的提问来源于stack exchange,提问作者nirmal

