项目场景:
最近在学习flink相关内容,测试flink写出到kafka的时候,抛出如下异常:
Exception in thread "main" java.lang.NoClassDefFoundError: org/apache/kafka/common/errors/InvalidTxnStateException
at com.blue.flink.sink.KafkaSinkTest$.main(KafkaSinkTest.scala:34)
at com.blue.flink.sink.KafkaSinkTest.main(KafkaSinkTest.scala)
Caused by: java.lang.ClassNotFoundException: org.apache.kafka.common.errors.InvalidTxnStateException
at java.net.URLClassLoader.findClass(URLClassLoader.java:381)
at java.lang.ClassLoader.loadClass(ClassLoader.java:424)
at sun.misc.Launcher$AppClassLoader.loadClass(Launcher.java:338)
at java.lang.ClassLoader.loadClass(ClassLoader.java:357)
... 2 more
**经查,是因为引入了
spark-streaming-kafka-0-10_2.11
和
flink-connector-kafka_2.11
两个pom依赖中的
kafka-client
版本冲突导致。**
有时候这些小问题还是会引起一些困惑,在此记录一下,也算是我自己排查jar冲突的一种简单手段
问题描述
相关代码:
val ds: DataStream[String]= env.readTextFile("./data/wcdata.txt")
ds.addSink(new FlinkKafkaProducer[String]("topic_bw",new KafkaSerializationSchema[String]{overridedef serialize(element:String, timestamp: lang.Long): ProducerRecord[Array[Byte], Array[Byte]]={new ProducerRecord[Array[Byte], Array[Byte]]("topic_bw",element.getBytes())}},
properties,
Semantic.EXACTLY_ONCE //恰好一次))
原因分析:
从报错可以很明显看出是缺少jar包,排查引入的pom文件和依赖项
从项目中查看依赖项,发现是kafka-client有冲突
pom分析插件也印证了这一点,而且下载的源码在idea中居然还报红 ,明显是依赖对不上了
解决方案:
去除低版本的依赖项
版权归原作者 蓝木 所有, 如有侵权,请联系我们删除。