【问题标题】:Flink fails to load ProducerRecord class with LinkageError at runtimeFlink 在运行时无法加载带有 LinkageError 的 ProducerRecord 类
【发布时间】:2020-08-24 19:04:28
【问题描述】:

使用 Scala 2.12 运行 Flink 1.9.0 并尝试使用 flink-connector-kafka 将数据发布到 Kafka,在本地调试时一切正常。将作业提交到集群后,我会在运行时收到以下 java.lang.LinkageError,但无法运行作业:

java.lang.LinkageError: loader constraint violation: loader (instance of org/apache/flink/util/ChildFirstClassLoader) previously initiated loading for a different type with name "org/apache/kafka/clients/producer/ProducerRecord"
    at java.lang.ClassLoader.defineClass1(Native Method)
    at java.lang.ClassLoader.defineClass(ClassLoader.java:763)
    at java.security.SecureClassLoader.defineClass(SecureClassLoader.java:142)
    at java.net.URLClassLoader.defineClass(URLClassLoader.java:468)
    at java.net.URLClassLoader.access$100(URLClassLoader.java:74)
    at java.net.URLClassLoader$1.run(URLClassLoader.java:369)
    at java.net.URLClassLoader$1.run(URLClassLoader.java:363)
    at java.security.AccessController.doPrivileged(Native Method)
    at java.net.URLClassLoader.findClass(URLClassLoader.java:362)
    at org.apache.flink.util.ChildFirstClassLoader.loadClass(ChildFirstClassLoader.java:66)
    at java.lang.ClassLoader.loadClass(ClassLoader.java:357)
    at java.lang.Class.getDeclaredMethods0(Native Method)
    at java.lang.Class.privateGetDeclaredMethods(Class.java:2701)
    at java.lang.Class.getDeclaredMethod(Class.java:2128)
    at java.io.ObjectStreamClass.getPrivateMethod(ObjectStreamClass.java:1629)
    at java.io.ObjectStreamClass.access$1700(ObjectStreamClass.java:79)
    at java.io.ObjectStreamClass$3.run(ObjectStreamClass.java:520)
    at java.io.ObjectStreamClass$3.run(ObjectStreamClass.java:494)
    at java.security.AccessController.doPrivileged(Native Method)
    at java.io.ObjectStreamClass.<init>(ObjectStreamClass.java:494)
    at java.io.ObjectStreamClass.lookup(ObjectStreamClass.java:391)
    at java.io.ObjectStreamClass.initNonProxy(ObjectStreamClass.java:681)
    at java.io.ObjectInputStream.readNonProxyDesc(ObjectInputStream.java:1885)
    at java.io.ObjectInputStream.readClassDesc(ObjectInputStream.java:1751)
    at java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2042)
    at java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1573)
    at java.io.ObjectInputStream.defaultReadFields(ObjectInputStream.java:2287)
    at java.io.ObjectInputStream.defaultReadObject(ObjectInputStream.java:561)
    at org.apache.flink.streaming.connectors.kafka.FlinkKafkaProducer.readObject(FlinkKafkaProducer.java:1202)
    at sun.reflect.GeneratedMethodAccessor358.invoke(Unknown Source)
    at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
    at java.lang.reflect.Method.invoke(Method.java:498)
    at java.io.ObjectStreamClass.invokeReadObject(ObjectStreamClass.java:1170)
    at java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:2178)
    at java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2069)
    at java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1573)
    at java.io.ObjectInputStream.defaultReadFields(ObjectInputStream.java:2287)
    at java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:2211)
    at java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2069)
    at java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1573)
    at java.io.ObjectInputStream.defaultReadFields(ObjectInputStream.java:2287)
    at java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:2211)
    at java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2069)
    at java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1573)
    at java.io.ObjectInputStream.readObject(ObjectInputStream.java:431)
    at org.apache.flink.util.InstantiationUtil.deserializeObject(InstantiationUtil.java:576)
    at org.apache.flink.util.InstantiationUtil.deserializeObject(InstantiationUtil.java:562)
    at org.apache.flink.util.InstantiationUtil.deserializeObject(InstantiationUtil.java:550)
    at org.apache.flink.util.InstantiationUtil.readObjectFromConfig(InstantiationUtil.java:511)
    at org.apache.flink.streaming.api.graph.StreamConfig.getStreamOperatorFactory(StreamConfig.java:235)
    at org.apache.flink.streaming.runtime.tasks.OperatorChain.createChainedOperator(OperatorChain.java:427)
    at org.apache.flink.streaming.runtime.tasks.OperatorChain.createOutputCollector(OperatorChain.java:354)
    at org.apache.flink.streaming.runtime.tasks.OperatorChain.createChainedOperator(OperatorChain.java:418)
    at org.apache.flink.streaming.runtime.tasks.OperatorChain.createOutputCollector(OperatorChain.java:354)
    at org.apache.flink.streaming.runtime.tasks.OperatorChain.createChainedOperator(OperatorChain.java:418)
    at org.apache.flink.streaming.runtime.tasks.OperatorChain.createOutputCollector(OperatorChain.java:354)
    at org.apache.flink.streaming.runtime.tasks.OperatorChain.createChainedOperator(OperatorChain.java:418)
    at org.apache.flink.streaming.runtime.tasks.OperatorChain.createOutputCollector(OperatorChain.java:354)
    at org.apache.flink.streaming.runtime.tasks.OperatorChain.createChainedOperator(OperatorChain.java:418)
    at org.apache.flink.streaming.runtime.tasks.OperatorChain.createOutputCollector(OperatorChain.java:354)
    at org.apache.flink.streaming.runtime.tasks.OperatorChain.createChainedOperator(OperatorChain.java:418)
    at org.apache.flink.streaming.runtime.tasks.OperatorChain.createOutputCollector(OperatorChain.java:354)
    at org.apache.flink.streaming.runtime.tasks.OperatorChain.createChainedOperator(OperatorChain.java:418)
    at org.apache.flink.streaming.runtime.tasks.OperatorChain.createOutputCollector(OperatorChain.java:354)
    at org.apache.flink.streaming.runtime.tasks.OperatorChain.createChainedOperator(OperatorChain.java:418)
    at org.apache.flink.streaming.runtime.tasks.OperatorChain.createOutputCollector(OperatorChain.java:354)
    at org.apache.flink.streaming.runtime.tasks.OperatorChain.<init>(OperatorChain.java:144)
    at org.apache.flink.streaming.runtime.tasks.StreamTask.invoke(StreamTask.java:370)
    at org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:705)
    at org.apache.flink.runtime.taskmanager.Task.run(Task.java:530)
    at java.lang.Thread.run(Thread.java:748)

使用-verbose:class 查看加载的类时,我看到该类被加载了几次:

taskmanager [Loaded org.apache.kafka.clients.producer.ProducerRecord from file:/tmp/blobStore-8cf95113-e767-4073-9b1b-e579d46c0283/job_f0c3db8b84dd38e83f92ecf1bc61b698/blob_p-c327eb8f4333a638b2b7049049368f23254aeb9c-03045e6d6a9c8f3c7dacdded8cb97d6e]
taskmanager [Loaded org.apache.kafka.clients.producer.ProducerRecord from file:/tmp/blobStore-8cf95113-e767-4073-9b1b-e579d46c0283/job_f0c3db8b84dd38e83f92ecf1bc61b698/blob_p-c327eb8f4333a638b2b7049049368f23254aeb9c-03045e6d6a9c8f3c7dacdded8cb97d6e]

类是从我提交给 Flink 的同一个 Uber-JAR 加载的。此外,没有加载 ProducerRecord 的多个传递依赖项,我的 JAR 是该依赖项的唯一提供者。

build.sbt:

lazy val flinkVersion = "1.9.0"

libraryDependencies ++= Seq(
    "org.apache.flink"                 %% "flink-table-planner"              % flinkVersion,
    "org.apache.flink"                 %% "flink-table-api-scala-bridge"     % flinkVersion,
    "org.apache.flink"                 % "flink-s3-fs-hadoop"                % flinkVersion,
    "org.apache.flink"                 %% "flink-container"                  % flinkVersion,
    "org.apache.flink"                 %% "flink-connector-kafka"            % flinkVersion,
    "org.apache.flink"                 %% "flink-scala"                      % flinkVersion % "provided",
    "org.apache.flink"                 %% "flink-streaming-scala"            % flinkVersion % "provided",
    "org.apache.flink"                 % "flink-json"                        % flinkVersion % "provided",
    "org.apache.flink"                 % "flink-avro"                        % flinkVersion % "provided",
    "org.apache.flink"                 %% "flink-parquet"                    % flinkVersion % "provided",
    "org.apache.flink"                 %% "flink-runtime-web"                % flinkVersion % "provided",
    "org.apache.flink"                 %% "flink-cep"                        % flinkVersion
)

【问题讨论】:

    标签: scala apache-flink


    【解决方案1】:

    由于未知原因,将classloader.resolve-order 属性设置为parent-first(如Apache Flink mailing list 中所述)可以解决此问题。我仍然对 为什么 感到困惑,因为在加载此依赖项的不同版本的子类加载器和父类加载器之间不应该存在依赖项冲突(因为它没有与 @987654326 一起提供开箱即用@我正在使用)。

    "Debugging Classloading" 下的Flink 文档中,有一个section which talks about this parent-child relationship

    在涉及动态类加载的设置中(插件组件, 会话设置中的 Flink 作业),通常有两个层次结构 ClassLoaders:(1)Java的应用类加载器,它拥有所有 类路径中的类,以及 (2) 动态插件/用户代码 类加载器。用于从插件或用户代码加载类 罐子。动态类加载器将应用程序类加载器作为其 父母。

    默认情况下,Flink 会反转类加载顺序,这意味着它会查看 首先是动态类加载器,并且只查看父级 (应用程序类加载器)如果类不是动态的 加载代码。

    反向类加载的好处是插件和作业可以使用 与 Flink 的核心本身不同的库版本,这非常 当库的不同版本不可用时很有用 兼容的。该机制有助于避免常见的依赖 IllegalAccessError 或 NoSuchMethodError 等冲突错误。 代码的不同部分仅具有单独的类副本 (Flink 的核心或其依赖项之一可以使用不同的副本 用户代码或插件代码)。在大多数情况下,这工作得很好,没有 需要用户进行额外配置。

    我还没有理解为什么加载ProducerRecord 会多次发生,或者异常消息中的这个“不同类型”指的是什么(对-verbose:class 的结果进行greping 只为ProducerRecord 产生了一个路径) .

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2012-04-12
      • 2017-12-20
      • 2019-09-17
      • 2023-03-12
      • 1970-01-01
      • 2019-01-14
      相关资源
      最近更新 更多