Spring Kafka Stream - No type information in headers and no default type provided when using BiFunction

Viewed 2662

I am trying to join 2 topics and produce output in to a 3rd topic using BiFunction. I am facing issue with resolving the type for incoming message. My left side message is getting deserialized successfully, but right side it throws "No type information in headers and no default type provided".

When I step through the code I could see it fails in the line org.springframework.kafka.support.serializer.JsonDeserializer

Assert.state(localReader != null, "No headers available and no default type provided");

The Messages are produced by spring boot with Kafka binder. And it has below properties.

spring.kafka.producer.key-serializer=org.apache.kafka.common.serialization.StringSerializer
spring.kafka.producer.value-serializer=org.springframework.kafka.support.serializer.JsonSerializer
##
spring.kafka.producer.properties.spring.json.type.mapping=type1:com.demo.domain.type2,type1:com.demo.domain.type2
spring.kafka.producer.properties.spring.json.trusted.packages=com.demo.domain
spring.kafka.producer.properties.spring.json.add.type.headers=true

And on the Kafka Stream binder consumer side

# kafka stream setting
spring.cloud.stream.bindings.joinProcess-in-0.destination=local-stream-process-type1
spring.cloud.stream.bindings.joinProcess-in-1.destination=local-stream-process-type2
spring.cloud.stream.bindings.joinProcess-out-0.destination=local-stream-process-type3

spring.cloud.stream.kafka.streams.binder.functions.joinProcess.applicationId=local-stream-process

spring.cloud.stream.kafka.streams.binder.configuration.default.key.serde=org.apache.kafka.common.serialization.Serdes$StringSerde
spring.cloud.stream.kafka.streams.binder.configuration.default.value.serde=org.springframework.kafka.support.serializer.JsonSerde

spring.kafka.streams.properties.spring.json.trusted.packages=*
spring.kafka.properties.spring.json.type.mapping=type1:com.demo.domain.type2,type1:com.demo.domain.type2
spring.kafka.streams.properties.spring.json.use.type.headers=true

And My Bifunction looks like


@Configuration
public class StreamsConfig {

    @Bean
    public RecordMessageConverter converter() {
        return new StringJsonMessageConverter();
    }


    @Bean
    public BiFunction<KStream<String, type1>, KStream<String, type2>, KStream<String, type3>> joinProcess() {
        return (type1, type2) ->
                type1.join(type2, joiner(),
                        JoinWindows.of(Duration.ofDays(1)));
    }

    private ValueJoiner<type1, type2, type3> joiner() {
         return (type1, type2) -> { new type3("test"); 
     };    
    }
}

I have pretty much went through all the previous questions and none of them were Bifunction. The one thing i havent tried is set VALUE_TYPE_METHOD.

###Update###

I resolved my issue with explicitly providing the serdes and disabling auto type conversion.

    @Bean
    public BiFunction<KStream<String, Type1>, KStream<String, Type2>, KStream<String, Type3>> joinStream() {
        return (type1, type2) ->
                type1.join(type2, myValueJoiner(),
                        JoinWindows.of(Duration.ofMinutes(1)), StreamJoined.with(Serdes.String(), new Type1Serde(), new Type2Serde()));
    }

And I disabled the automatic Deserialization like below

spring.cloud.stream.bindings.joinStream-in-0.consumer.use-native-decoding=false
spring.cloud.stream.bindings.joinStream-in-1.consumer.use-native-decoding=false

0 Answers
Related