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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 18:06:02