You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.28 07:35:18