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

Spring Cloud Stream与Kafka Streams Join时ClassCastException问题排查

问题

本人是Spring Cloud Stream和Kafka Streams新手,正尝试基于输入流的映射键构建新KStream。在对Join后的流进行反序列化时出现ClassCastException错误(错误日志如下)。我推测是由于Spring Cloud Stream绑定配置中设置的GenericAvroSerde被应用到了Join后的KStream,若该推测正确,请问如何为Join的KStream指定不同的Serde?若推测错误,该问题的根本原因是什么?

错误日志

org.apache.kafka.streams.errors.StreamsException: ClassCastException invoking processor: second-first-join-other-join. Do the Processor's input types match the deserialized types? Check the Serde setup and change the default Serdes in StreamConfig or provide correct Serdes via method parameters. Make sure the Processor can accept the deserialized input of type key: java.lang.String, and value: my.package.MyFirstTopicValue.
Note that although incorrect Serdes are a common cause of error, the cast exception might have another cause (in user code, for example). For example, if a processor wires in a store, but casts the generics incorrectly, a class cast exception could be raised during processing, but the cause would not be wrong Serdes.
    at org.apache.kafka.streams.processor.internals.ProcessorNode.process(ProcessorNode.java:150) ~[kafka-streams-3.1.2.jar:na]
    at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forwardInternal(ProcessorContextImpl.java:253) ~[kafka-streams-3.1.2.jar:na]
    at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forward(ProcessorContextImpl.java:232) ~[kafka-streams-3.1.2.jar:na]
    at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forward(ProcessorContextImpl.java:191) ~[kafka-streams-3.1.2.jar:na]
    at org.apache.kafka.streams.kstream.internals.KStreamJoinWindow$KStreamJoinWindowProcessor.process(KStreamJoinWindow.java:55) ~[kafka-streams-3.1.2.jar:na]
    at org.apache.kafka.streams.processor.internals.ProcessorNode.process(ProcessorNode.java:146) ~[kafka-streams-3.1.2.jar:na]
    at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forwardInternal(ProcessorContextImpl.java:253) ~[kafka-streams-3.1.2.jar:na]
    at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forward(ProcessorContextImpl.java:232) ~[kafka-streams-3.1.2.jar:na]
    at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forward(ProcessorContextImpl.java:191) ~[kafka-streams-3.1.2.jar:na]
    at org.apache.kafka.streams.processor.internals.SourceNode.process(SourceNode.java:84) ~[kafka-streams-3.1.2.jar:na]
    at org.apache.kafka.streams.processor.internals.StreamTask.lambda$process$1(StreamTask.java:731) ~[kafka-streams-3.1.2.jar:na]
    at org.apache.kafka.streams.processor.internals.metrics.StreamsMetricsImpl.maybeMeasureLatency(StreamsMetricsImpl.java:809) ~[kafka-streams-3.1.2.jar:na]
    at org.apache.kafka.streams.processor.internals.StreamTask.process(StreamTask.java:731) ~[kafka-streams-3.1.2.jar:na]
    at org.apache.kafka.streams.processor.internals.TaskManager.process(TaskManager.java:1296) ~[kafka-streams-3.1.2.jar:na]
    at org.apache.kafka.streams.processor.internals.StreamThread.runOnce(StreamThread.java:784) ~[kafka-streams-3.1.2.jar:na]
    at org.apache.kafka.streams.processor.internals.StreamThread.runLoop(StreamThread.java:604) ~[kafka-streams-3.1.2.jar:na]
    at org.apache.kafka.streams.processor.internals.StreamThread.run(StreamThread.java:576) ~[kafka-streams-3.1.2.jar:na]
Caused by: java.lang.ClassCastException: class java.util.LinkedHashMap cannot be cast to class my.package.MySecondTopicValue (java.util.LinkedHashMap is in module java.base of loader 'bootstrap'; my.package.MySecondTopicValue is in unnamed module of loader 'app')
    at org.apache.kafka.streams.kstream.internals.AbstractStream.lambda$reverseJoinerWithKey$1(AbstractStream.java:106) ~[kafka-streams-3.1.2.jar:na]
    at org.apache.kafka.streams.kstream.internals.KStreamKStreamJoin$KStreamKStreamJoinProcessor.process(KStreamKStreamJoin.java:177) ~[kafka-streams-3.1.2.jar:na]
    at org.apache.kafka.streams.processor.internals.ProcessorNode.process(ProcessorNode.java:146) ~[kafka-streams-3.1.2.jar:na]
    ... 16 common frames omitted

2022-10-31 15:55:03.913 ERROR 25964 --- [-StreamThread-1] org.apache.kafka.streams.KafkaStreams    : stream-client [...-processors-group-161db237-fe34-4ed8-a6b5-b57ed2f09a73] Encountered the following exception during processing and Kafka Streams opted to SHUTDOWN_CLIENT. The streams client is going to shut down now. 

org.apache.kafka.streams.errors.StreamsException: ClassCastException invoking processor: second-first-join-other-join. Do the Processor's input types match the deserialized types? Check the Serde setup and change the default Serdes in StreamConfig or provide correct Serdes via method parameters. Make sure the Processor can accept the deserialized input of type key: java.lang.String, and value: my.package.MyFirstTopicValue.
Note that although incorrect Serdes are a common cause of error, the cast exception might have another cause (in user code, for example). For example, if a processor wires in a store, but casts the generics incorrectly, a class cast exception could be raised during processing, but the cause would not be wrong Serdes.
Caused by: java.lang.ClassCastException: class java.util.LinkedHashMap cannot be cast to class my.package.MySecondTopicValue (java.util.LinkedHashMap is in module java.base of loader 'bootstrap'; my.package.MySecondTopicValue is in unnamed module of loader 'app')

application.properties配置片段

...

spring.cloud.stream.kafka.streams.bindings.myEventProcessor-in-0.consumer.keySerde=io.confluent.kafka.streams.serdes.avro.GenericAvroSerde
spring.cloud.stream.kafka.streams.bindings.myEventProcessor-in-0.consumer.valueSerde=io.confluent.kafka.streams.serdes.avro.GenericAvroSerde
spring.cloud.stream.kafka.streams.bindings.myEventProcessor-in-1.consumer.keySerde=io.confluent.kafka.streams.serdes.avro.GenericAvroSerde
spring.cloud.stream.kafka.streams.bindings.myEventProcessor-in-1.consumer.valueSerde=io.confluent.kafka.streams.serdes.avro.GenericAvroSerde

...

相关Java代码

@Bean
public BiConsumer<KStream<GenericData.Record, GenericData.Record>, KStream<GenericData.Record, GenericData.Record>> myEventProcessor() {
    return (mySecondGenericKStream, myFirstGenericKStream) -> {

        // Building my first typed KStream from AvroGenericRecord
        KStream<String, MyFirstTopicValue> myFirstWithStringKeyKStream = myFirstGenericKStream
                .map((key, value) -> {
                    CdcRecordAdapter adapter = new CdcRecordAdapter(key, value);
                    MyFirstTopicValue myFirstTopicValue = adapter.getValueAs(MyFirstTopicValue.class);
                    String stringKey = myFirstTopicValue.getTransactionId();
                    return KeyValue.pair(stringKey, myFirstTopicValue);
                })
                .repartition(Repartitioned.with(new Serdes.StringSerde(), newFirstJsonSerde()).withName("myFirstWithStringKeyKStream"))
                .peek((key, value) -> printOnConsole("myFirstWithStringKeyKStream", key, value.toString()));

        // Building my second typed KStream from AvroGenericRecord
        KStream<String, MySecondTopicValue> mySecondWithStringKeyKStream = mySecondGenericKStream
                .map((key, value) -> {
                    CdcRecordAdapter adapter = new CdcRecordAdapter(key, value);
                    MySecondTopicKey mySecondTopicKey = adapter.getKeyAs(MySecondTopicKey.class);
                    MySecondTopicValue mySecondTopicValue = adapter.getValueAs(MySecondTopicValue.class);
                    String stringKey = mySecondTopicKey.getTransactionId();
                    return KeyValue.pair(stringKey, mySecondTopicValue);
                })
                .repartition(Repartitioned.with(new Serdes.StringSerde(), newSecondJsonSerde()).withName("mySecondWithStringKeyKStream"))
                .peek((key, value) -> printOnConsole("mySecondWithStringKeyKStream", key, value.toString()));

        // Join two mapped KStreams with new keys
        mySecondWithStringKeyKStream
                .filter((key, value) -> value.getCrudOperation() == CrudOperation.INSERT)
                .leftJoin(
                        myFirstWithStringKeyKStream,
                        (readOnlyKey, mySecondValue, myFirstValue) -> new MySecondFirstBuilder(readOnlyKey, mySecondValue, myFirstValue).build(),
                        JoinWindows.ofTimeDifferenceAndGrace(Duration.ofSeconds(15L), Duration.ofSeconds(5)),
                        StreamJoined.with(Serdes.String(), newSecondJsonSerde(), newFirstJsonSerde())
                                .withName("second-first-topics-joined")
                                .withStoreName("second-firs-topics-joined"))
                .peek((key, value) -> printOnConsole("second-first-topics-joined", key, value.toString()));
        
        //.process(() -> new MyProcessor(myEventPublisher))
    };
}
问题分析与解决

根本原因

你的推测不对,问题核心是Join操作的窗口存储没有使用你指定的JSON Serde,而是 fallback 到了全局默认Serde,导致存储中的值被反序列化为LinkedHashMap,无法转换成自定义的MySecondTopicValue类型。

从错误日志能看到,系统在Join处理时把LinkedHashMap强转成MySecondTopicValue失败,说明Join的窗口存储读取数据时,没应用你在StreamJoined里配置的Serde。你在repartition时指定了Serde,但Join的窗口存储是独立组件,不会自动继承上游流的Serde设置。

解决方法

1. 确保JSON Serde是带泛型的类型化实例

手动创建JSON Serde时必须指定泛型,否则反序列化只会返回LinkedHashMap:

// 正确创建对应类型的JSON Serde
Serde<MySecondTopicValue> newSecondJsonSerde = new JsonSerde<>(MySecondTopicValue.class);
Serde<MyFirstTopicValue> newFirstJsonSerde = new JsonSerde<>(MyFirstTopicValue.class);

如果用Spring自动配置的Serde,也要通过泛型或Serdes.serdeFrom()明确类型。

2. 强化Join操作的Serde配置

虽然你用了StreamJoined.with(),但要确保窗口存储完全应用指定的Serde:

  • 核对StreamJoined参数顺序:StreamJoined.with(keySerde, 左流valueSerde, 右流valueSerde),你的代码顺序是正确的(StringSerde对应key,newSecondJsonSerde对应左流mySecondWithStringKeyKStream的value,newFirstJsonSerde对应右流)。
  • 可设置全局默认Serde避免 fallback,在application.properties中添加:
spring.cloud.stream.kafka.streams.binder.configuration.default.key.serde=org.apache.kafka.common.serialization.Serdes$StringSerde
# 注意:全局用JsonSerde时,仍需确保每个流处理环节明确泛型,否则还是会出现LinkedHashMap问题
spring.cloud.stream.kafka.streams.binder.configuration.default.value.serde=org.springframework.kafka.support.serializer.JsonSerde

更稳妥的方式是每个操作都明确指定Serde,不要依赖全局默认。

3. 若Join后的流要输出到Topic,需配置输出绑定的Serde

如果Join后的结果要发送到某个Kafka Topic,必须为输出绑定配置对应Serde:

spring.cloud.stream.kafka.streams.bindings.myEventProcessor-out-0.producer.keySerde=org.apache.kafka.common.serialization.Serdes$StringSerde
# 替换成你Join后生成的结果类型对应的Serde
spring.cloud.stream.kafka.streams.bindings.myEventProcessor-out-0.producer.valueSerde=my.package.MyJoinedResultSerde

4. 验证Repartition环节的Serde是否生效

可以查看repartition对应的Kafka Topic(比如myFirstWithStringKeyKStream)的消息内容,确认序列化是否正确,避免在repartition环节就出现数据格式错误。


内容的提问来源于stack exchange,提问作者nanachimi

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 23:40:36