【发布时间】: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