Kubernetes中Flink s3-fs-hadoop插件无法使用的问题排查
问题背景
使用Flink 1.18.1读取S3文件,已在Dockerfile中将s3-fs-hadoop插件jar包移动到/opt/flink/plugins/s3-fs-hadoop/目录,路径确认正确,但运行Pod时抛出类找不到异常:
Caused by: java.lang.ClassNotFoundException: Class org.apache.hadoop.fs.s3a.S3AFileSystem not found
同时日志显示插件已成功加载:
[] - Plugin loader with ID found, reusing it: s3-fs-hadoop [] - Delegation token receiver s3-hadoop loaded and initialized
添加hadoop-aws依赖后问题解决,但认为这属于冗余,插件应无需额外依赖即可工作,疑问:这是否是类加载问题?该如何解决?
使用环境
- Flink 1.18.1
- Scala 2.12.6
相关配置
YAML配置(FlinkDeployment及关联资源)
apiVersion: flink.apache.org/v1beta1 kind: FlinkDeployment metadata: name: flink-test namespace: mynamespace finalizers: - flinkdeployments.flink.apache.org/finalizer spec: image: .../flink/test:v1 flinkVersion: v1_18 flinkConfiguration: taskmanager.numberOfTaskSlots: "4" taskmanager.memory.managed.fraction: "0.1" classloader.resolve-order: parent-first metrics.reporters: prom metrics.reporter.prom.factory.class: org.apache.flink.metrics.prometheus.PrometheusReporterFactory metrics.reporter.prom.port: "9250" serviceAccount: flink jobManager: resource: memory: 2048m cpu: 1 taskManager: replicas: 1 resource: memory: 4096m cpu: 4 job: jarURI: local:///tmp/myjar-0.0.1.jar parallelism: 2 upgradeMode: stateless # stateless or savepoint or last-state entryClass: com.org.MyMainClass args: [...] podTemplate: apiVersion: v1 kind: Pod metadata: name: flink-test spec: containers: - name: flink-main-container securityContext: runAsUser: 9999 # UID of a non-root user runAsNonRoot: true ports: - name: metrics containerPort: 9250 protocol: TCP volumeMounts: - mountPath: /etc/keystore name: user-cert readOnly: true - mountPath: /etc/truststore name: ca-cert readOnly: true volumes: - name: user-cert secret: secretName: the-user - name: ca-cert secret: secretName: cluster-cert --- apiVersion: monitoring.coreos.com/v1 kind: PodMonitor metadata: name: flink-test labels: release: prometheus spec: selector: matchLabels: app: test podMetricsEndpoints: - port: metrics --- apiVersion: networking.k8s.io/v1 kind: Ingress metadata: name: flink-test-ingress annotations: nginx.ingress.kubernetes.io/rewrite-target: /$2 spec: ingressClassName: ingress-nginx-iec rules: - host: "hostname" http: paths: - pathType: Prefix path: "/flink/flink-test(/|$)(.*)" backend: service: name: flink-test-rest port: number: 8081 # Replace with your service port tls: - hosts: - host secretName: verysecretname
Dockerfile
FROM flink:1.18.1-scala_2.12 USER root RUN mkdir -p /opt/flink/log/ /opt/flink/conf/ /opt/flink/plugins/s3-fs-hadoop && \ cp /opt/flink/opt/flink-s3-fs-hadoop-1.18.1.jar /opt/flink/plugins/s3-fs-hadoop/ && \ chown -R flink:flink /opt/flink/ &&\ chmod -R 755 /opt/flink/ # You can choose whatever directory you want WORKDIR /opt/flink/lib # Copy your JAR files COPY ../target/scala-2.12/test-0.0.1.jar /tmp/test-0.0.1.jar
Assembly配置(Scala)
assemblyMergeStrategy := { case PathList("META-INF", xs@_*) => MergeStrategy.discard case PathList("META-INF", "services", xs@_*) => MergeStrategy.concat case PathList("reference.conf") => MergeStrategy.concat case _ => MergeStrategy.first }
问题分析与解决方案
1. 本质原因:插件的依赖传递性
flink-s3-fs-hadoop插件本身不包含hadoop-aws的依赖,它只是Flink与Hadoop S3A文件系统的适配层,实际的S3A实现类org.apache.hadoop.fs.s3a.S3AFileSystem属于hadoop-aws包。官方镜像中的/opt/flink/opt/目录下只有插件本身,没有附带其依赖的Hadoop AWS组件,所以当插件加载后,找不到底层的实现类。
2. 是否属于类加载问题?
是,但不是插件加载失败,而是插件依赖的底层类无法被类加载器找到。日志显示插件已成功加载,但当实际执行S3文件操作时,需要加载S3AFileSystem,而这个类不在插件的类路径中,也不在Flink的默认类路径里。
3. 解决方案(无需在业务Jar中添加冗余依赖)
方案一:在Docker镜像中添加hadoop-aws及相关依赖
修改Dockerfile,下载与Flink兼容版本的hadoop-aws和hadoop-common(Flink 1.18.1默认适配Hadoop 3.3.4),放入插件目录:
FROM flink:1.18.1-scala_2.12 USER root # 创建插件目录并复制Flink S3插件 RUN mkdir -p /opt/flink/plugins/s3-fs-hadoop && \ cp /opt/flink/opt/flink-s3-fs-hadoop-1.18.1.jar /opt/flink/plugins/s3-fs-hadoop/ && \ # 下载hadoop-aws和依赖的hadoop-common curl -L https://repo1.maven.org/maven2/org/apache/hadoop/hadoop-aws/3.3.4/hadoop-aws-3.3.4.jar -o /opt/flink/plugins/s3-fs-hadoop/hadoop-aws-3.3.4.jar && \ curl -L https://repo1.maven.org/maven2/org/apache/hadoop/hadoop-common/3.3.4/hadoop-common-3.3.4.jar -o /opt/flink/plugins/s3-fs-hadoop/hadoop-common-3.3.4.jar && \ chown -R flink:flink /opt/flink/ &&\ chmod -R 755 /opt/flink/ # 复制业务Jar COPY ../target/scala-2.12/test-0.0.1.jar /tmp/test-0.0.1.jar
将依赖放入插件目录的好处是,这些类会和插件一起被插件类加载器加载,避免类路径冲突。
方案二:使用Flink提供的完整Hadoop镜像
直接使用flink:1.18.1-scala_2.12-hadoop3镜像,该镜像已包含完整的Hadoop依赖,包括hadoop-aws,无需额外添加。
方案三:调整类加载顺序(不推荐)
当前配置中classloader.resolve-order: parent-first,如果改为child-first,业务Jar中的hadoop-aws依赖会被优先加载,但这可能引发其他类冲突问题,不建议在生产环境使用。
4. 为什么官方文档没有明确说明?
Flink的插件机制是解耦的,flink-s3-fs-hadoop作为插件,依赖于Hadoop的生态组件,官方默认假设用户会自行提供Hadoop环境或依赖,尤其是在容器化场景下,需要用户根据需求补充依赖。
内容的提问来源于stack exchange,提问作者Niko

