Kafka Streams中ValueJoiner为何触发String类型转换异常?
Kafka Streams ClassCastException 异常求助
我基于之前的问题扩展了示例代码如下:
package content.streams import content.models.Models import org.apache.kafka.common.serialization.Serdes import org.apache.kafka.streams.{KafkaStreams, KeyValue, StreamsBuilder, StreamsConfig} import org.apache.kafka.streams.kstream.{JoinWindows, KStream, KeyValueMapper, Printed, ValueJoiner, ValueJoinerWithKey} import java.time.Duration import java.util.Properties object KafkaStreams { private val config: Properties = { val p = new Properties() p.put(StreamsConfig.APPLICATION_ID_CONFIG, "kafka-streams-test") p.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9094") p.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass) p.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass) p } private object Topics { val streamInputA = "streamA" val streamInputB = "streamB" } private val builder = new StreamsBuilder() // can't chain, gives cryptic error .. private val streamR1: KStream[String, String] = builder.stream[String, String](Topics.streamInputA) private val streamR2: KStream[String, Models.Event] = streamR1.mapValues(str => Models.loads(str)) private val streamR3: KStream[String, Models.Event] = streamR2.selectKey( (_, value) => value.sessionId ) private val streamC1: KStream[String, String] = builder.stream[String, String](Topics.streamInputB) private val streamC2: KStream[String, Models.Event] = streamC1.mapValues(str => Models.loads(str)) private val streamC3: KStream[String, Models.Event] = streamC2.selectKey((_, value) => value.sessionId) // works streamR3.peek((key, value) => println("key:" + key + " value:" + value)) streamC3.peek((key, value) => println("key:" + key + " value:" + value)) private val joinWindow = JoinWindows.ofTimeDifferenceWithNoGrace(Duration.ofSeconds(10)) private val valueJoinerWithKey = new ValueJoinerWithKey[String, Models.Event, Models.Event, Models.Event] { override def apply(key: String, evt1: Models.Event, evt2: Models.Event): Models.Event = { // doesn't get to here // println("valueJoinerWithKey " + evt1 + " " + evt2) evt1 } } streamC3.join(streamR3, valueJoinerWithKey, joinWindow) private val streams: KafkaStreams = new KafkaStreams(builder.build(), config) def main(args: Array[String]) = { streams.start() } }
Models 文件内容如下:
... object Models { implicit val formats: Formats = DefaultFormats case class Event( sessionId: String, eventType: String, eventTime: String, ) ...
程序在ValueJoinerWithKey执行过程中抛出如下异常:
Caused by: java.lang.ClassCastException: class content.models.Models$Event cannot be cast to class java.lang.String (content.models.Models$Event is in unnamed module of loader 'app'; java.lang.String is in module java.base of loader 'bootstrap')
我无法理解该异常产生的原因,特此求助。
内容的提问来源于stack exchange,提问作者kev
相关产品推荐
相关产品推荐

