【问题标题】:Flink application ClassCastExceptionFlink 应用程序 ClassCastException
【发布时间】:2022-09-28 14:56:45
【问题描述】:

我有一个 flink 应用程序,它从 kafka 读取并将其下沉到 kafka。

当我从 Intellij IDEA 运行应用程序时,应用程序运行没有问题,但是当我将 shadowJar 提交到 flink 集群时会给出 ClassCastException。我可以得到一些帮助来弄清楚我在这里做错了什么吗?

异常跟踪:

Caused by: java.lang.ClassCastException: cannot assign instance of org.apache.kafka.clients.consumer.OffsetResetStrategy to field org.apache.flink.connector.kafka.source.enumerator.initializer.ReaderHandledOffsetsInitializer.offsetResetStrategy of type org.apache.kafka.clients.consumer.OffsetResetStrategy in instance of org.apache.flink.connector.kafka.source.enumerator.initializer.ReaderHandledOffsetsInitializer
    at java.base/java.io.ObjectStreamClass$FieldReflector.setObjFieldValues(ObjectStreamClass.java:2205)
    at java.base/java.io.ObjectStreamClass$FieldReflector.checkObjectFieldValueTypes(ObjectStreamClass.java:2168)
    at java.base/java.io.ObjectStreamClass.checkObjFieldValueTypes(ObjectStreamClass.java:1422)
    at java.base/java.io.ObjectInputStream.defaultCheckFieldValues(ObjectInputStream.java:2517)
    at java.base/java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:2424)
    at java.base/java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2233)
    at java.base/java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1692)
    at java.base/java.io.ObjectInputStream.defaultReadFields(ObjectInputStream.java:2501)
    at java.base/java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:2395)
    at java.base/java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2233)

使用的代码:

KafkaSource<String> source = KafkaSource.<String>builder()
                    .setBootstrapServers(\"localhost:9092\")
                    .setTopics(\"topic\")
                    .setGroupId(\"grp\")
                    .setStartingOffsets(OffsetsInitializer.earliest())
                    .setValueOnlyDeserializer(new SimpleStringSchema())
                    .build();


            DataStream<String> eventStream = env.fromSource(source, WatermarkStrategy.noWatermarks(), \"Kafka Source\")
                    .name(\"event-stream\").sinkTo(\"kafka\");

构建文件:flinkVersion = 1.15.0

 //flinkShadowJar \"org.apache.flink:flink-connector-kafka:${flinkVersion}\"

    // https://mvnrepository.com/artifact/org.apache.flink/flink-streaming-java
    implementation group: \'org.apache.flink\', name: \'flink-streaming-java\', version: \"${flinkVersion}\"

    // https://mvnrepository.com/artifact/org.apache.flink/flink-java
    implementation group: \'org.apache.flink\', name: \'flink-java\', version: \"${flinkVersion}\"

    // https://mvnrepository.com/artifact/org.apache.flink/flink-core
    implementation group: \'org.apache.flink\', name: \'flink-core\', version: \"${flinkVersion}\"

    // https://mvnrepository.com/artifact/org.apache.flink/flink-clients
    implementation group: \'org.apache.flink\', name: \'flink-clients\', version: \"${flinkVersion}\"

    // https://mvnrepository.com/artifact/org.apache.kafka/kafka
   // flinkShadowJar group: \'org.apache.kafka\', name: \'kafka_2.12\', version: \"${kafkaVersion}\"

    flinkShadowJar \"org.apache.avro:avro:1.11.0\"
    flinkShadowJar group: \'org.apache.flink\', name: \'flink-avro-confluent-registry\', version: \"${flinkVersion}\"

    // https://mvnrepository.com/artifact/org.apache.flink/flink-connector-kafka
    flinkShadowJar group: \'org.apache.flink\', name: \'flink-connector-kafka\', version: \"${flinkVersion}\"
  //  flinkShadowJar group: \'org.apache.flink\', name: \'flink-connector-base\', version: \"${flinkVersion}\"

    // https://mvnrepository.com/artifact/org.apache.logging.log4j/log4j-core
    implementation group: \'org.apache.logging.log4j\', name: \'log4j-core\', version: \'2.17.2\'

    // https://mvnrepository.com/artifact/org.apache.logging.log4j/log4j-slf4j-impl
    implementation group: \'org.apache.logging.log4j\', name: \'log4j-slf4j-impl\', version: \'2.17.2\'

    // https://mvnrepository.com/artifact/org.apache.logging.log4j/log4j-api
    implementation group: \'org.apache.logging.log4j\', name: \'log4j-api\', version: \'2.17.2\'
  • 我认为您收到的错误表明已流式传输的 org.apache.kafka.clients.consumer.OffsetResetStrategy 对象是在接收上下文中加载的同一类的不同版本。我会查看您的开发环境中的依赖关系,并确保它们与您的集群兼容。
  • 我正在使用 flink 集群版本 1.15 并在我的代码中使用相同的版本。如果有帮助,我已经复制了我的构建脚本
  • 您的应用程序 jar 是否包含在 /opt/flink/usrlib/classpath 的 Flink 容器中,您是否在 /opt/flink/lib 中拥有 flink 提供的库和在 /opt/flink/plugins 中的 flink 插件?我遇到了同样的问题 - 在本地工作,在我的 k3s 集群上失败,出现同样的错误。我检查了我正在构建的容器,它似乎具有正确版本的所有内容(1.15.0),并且 intellij 使用的 jar 似乎与 Flink 容器的 /opt/flink/lib 文件夹中的 jar 相同.

标签: java apache-flink


【解决方案1】:

你解决了吗?我遇到了同样的问题。

【讨论】:

    猜你喜欢
    • 2015-03-06
    • 2021-09-30
    • 2016-10-31
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2022-01-27
    相关资源
    最近更新 更多