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

Flink 1.13.5自托管集群读取S3 ORC文件报错求助

环境与代码

使用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)

当前配置与依赖

  1. 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版本依赖
  2. 集群插件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
  3. 集群状态:未安装完整Hadoop包(无/etc/hadoop/conf目录),未设置HADOOP_HOME环境变量。

疑问

读取S3上的ORC文件是否必须安装Hadoop?或是遗漏了其他配置?


解决方案分析

是否必须安装完整Hadoop包?

不需要安装完整的Hadoop发行版,但Flink的ORC读取和S3文件系统依赖Hadoop的核心API,因此需要确保这些API的依赖能被Flink正确加载。

问题原因分析

  1. 类加载冲突:Fat Jar中打包的Hadoop依赖(如hadoop-common、hadoop-client)与Flink插件目录中的s3-fs-hadoop插件依赖的Hadoop类可能存在版本冲突,或者类加载器隔离导致org.apache.hadoop.fs.FileSystem无法被正确找到。Flink的插件类加载器与应用类加载器是隔离的,重复打包依赖容易引发类加载问题。
  2. Hadoop配置缺失:虽然未安装Hadoop,但需要显式配置S3文件系统的实现类,否则Flink无法正确初始化S3文件系统的Hadoop API。

具体解决步骤

  1. 调整依赖打包策略
    将Hadoop相关依赖(hadoop-client、hadoop-common、hadoop-aws、hadoop-hdfs-client)设置为provided scope(Maven/Gradle),避免打包进Fat Jar。因为Flink的s3-fs-hadoop插件已经依赖这些类,重复打包会导致类加载冲突。

  2. 配置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
    
  3. 调整类加载顺序
    如果仍存在类加载问题,在flink-conf.yaml中设置类加载优先级为父类加载器优先:

    classloader.resolve-order: parent-first
    

    确保Hadoop的核心类由Flink的系统类加载器加载,而非应用类加载器。

  4. 修正代码中的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")
    
  5. 清理冗余插件
    如果不使用Presto的S3文件系统实现,删除集群节点上的/opt/flink/plugins/s3-fs-presto/目录下的Jar包,避免不必要的类加载冲突。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 19:40:55