Join KStream with KTable got No headers available and no default type provided

Viewed 17

I'm trying to join KStream with KTable but still got

java.lang.IllegalStateException: No headers available and no default type provided 

Actually I can't realize what I'm doing bad, can you help me please..? Have same number of partition and replicas on topics, same message key.

Using JsonSerdes on producer side.

Code which doing this, (if I use one of commented returns it works fine):

@StreamListener
@SendTo("simulations-result-output-channel")
public KStream<UUID, CandleWithParameters> process(
        @Input("simulations-parameters-input-channel") KTable<UUID, SimulationDto> simulationDtoKTable,
        @Input("simulations-data-input-channel") KStream<UUID, CandleDto> candleDtoKStream
) {

  //  return candleDtoKStream
  //            .mapValues((c) -> new CandleWithParameters(c, null));

  //  return simulationDtoKTable
  //         .mapValues((p) -> new CandleWithParameters(null, p.getParameters()))
  //         .toStream();


    return candleDtoKStream
            .peek((k, v) -> log.info("Key: {}, Candle {}", k, v))
            .join(
                    simulationDtoKTable,
                    (candle, simulation) -> new CandleWithParameters(candle, simulation.getParameters()),
                    Joined.with(
                            new Serdes.UUIDSerde(),
                            new JsonSerde<>(CandleDto.class),
                            new JsonSerde<>(SimulationDto.class)
                    )
            )
            .peek((k, v) -> log.info("Key: {}, CandleWithParams {}", k, v));

Config:

spring:
    cloud:
        stream:
            bindings:
                simulations-parameters-input-channel:
                    destination: simulation.parameters
                simulations-data-input-channel:
                    destination: simulation.data
                simulations-result-output-channel:
                    destination: simulation.results
            kafka:
                streams:
                    binder:
                        brokers: broker-1:29092
                        configuration:
                            commit.interval.ms: 10
                            security.protocol: SSL
                            state.dir: my-store
                            default:
                                key.serde: org.apache.kafka.common.serialization.Serdes$UUIDSerde
                                value.serde: org.springframework.kafka.support.serializer.JsonSerde
                            spring.json.trusted.packages: '*'
                            ssl:
                                truststore:
                                    location: classpath:security/client.ts.p12
                                    password: xxx
                                    type: PKCS12
                                keystore:
                                    location: classpath:security/robot.ks.p12
                                    password: xxx
                                    type: PKCS12

Edit, if I update to functional style error is the same:

Bean:

@Configuration
@Slf4j
public class StreamListenerConfig {

    @Bean
    public BiFunction<KStream<UUID, CandleDto>, KTable<UUID, SimulationDto>, KStream<UUID, CandleWithParameters>> simulationsProcess() {
        return (candleDtoKStream, simulationDtoKTable) ->
                candleDtoKStream
                    .peek((k, v) -> log.info("Key: {}, Candle {}", k, v))
                    .join(
                            simulationDtoKTable,
                            (candle, simulation) -> new CandleWithParameters(candle, simulation.getParameters()),
                            Joined.with(
                                    new Serdes.UUIDSerde(),
                                    new JsonSerde<>(CandleDto.class),
                                    new JsonSerde<>(SimulationDto.class)
                            )
                    )
                    .peek((k, v) -> log.info("Key: {}, CandleWithParams {}", k, v));
        }
}

Config update:

spring:
    cloud:
    stream:
        bindings:
            simulationsProcess-in-0:
                destination: simulation.data
            simulationsProcess-in-1:
                destination: simulation.parameters
            simulationsProcess-out-0:
                destination: simulation.results
        kafka:
            streams:
                binder:
                    ... Same config as above...
                    functions:
                        simulationsProcess.applicationId: simulation-robot

Trace:

java.lang.IllegalStateException: No headers available and no default type provided
2022-09-24T11:11:49.141715700Z  at org.springframework.util.Assert.state(Assert.java:76) ~[spring-core-5.3.22.jar!/:5.3.22]
2022-09-24T11:11:49.141721200Z  at org.springframework.kafka.support.serializer.JsonDeserializer.deserialize(JsonDeserializer.java:605) ~[spring-kafka-2.8.8.jar!/:2.8.8]
2022-09-24T11:11:49.141724700Z  at org.apache.kafka.streams.state.internals.ValueAndTimestampDeserializer.deserialize(ValueAndTimestampDeserializer.java:58) ~[kafka-streams-3.1.1.jar!/:na]
2022-09-24T11:11:49.141728000Z  at org.apache.kafka.streams.state.internals.ValueAndTimestampDeserializer.deserialize(ValueAndTimestampDeserializer.java:31) ~[kafka-streams-3.1.1.jar!/:na]
2022-09-24T11:11:49.141732000Z  at org.apache.kafka.streams.state.StateSerdes.valueFrom(StateSerdes.java:163) ~[kafka-streams-3.1.1.jar!/:na]
2022-09-24T11:11:49.141736000Z  at org.apache.kafka.streams.state.internals.MeteredKeyValueStore$1.apply(MeteredKeyValueStore.java:183) ~[kafka-streams-3.1.1.jar!/:na]
2022-09-24T11:11:49.141739900Z  at org.apache.kafka.streams.state.internals.MeteredKeyValueStore$1.apply(MeteredKeyValueStore.java:178) ~[kafka-streams-3.1.1.jar!/:na]
2022-09-24T11:11:49.141743600Z  at org.apache.kafka.streams.state.internals.CachingKeyValueStore.putAndMaybeForward(CachingKeyValueStore.java:107) ~[kafka-streams-3.1.1.jar!/:na]
2022-09-24T11:11:49.141747000Z  at org.apache.kafka.streams.state.internals.CachingKeyValueStore.lambda$initInternal$0(CachingKeyValueStore.java:87) ~[kafka-streams-3.1.1.jar!/:na]
2022-09-24T11:11:49.141750800Z  at org.apache.kafka.streams.state.internals.NamedCache.flush(NamedCache.java:151) ~[kafka-streams-3.1.1.jar!/:na]
2022-09-24T11:11:49.141754400Z  at org.apache.kafka.streams.state.internals.NamedCache.flush(NamedCache.java:109) ~[kafka-streams-3.1.1.jar!/:na]
2022-09-24T11:11:49.141758200Z  at org.apache.kafka.streams.state.internals.ThreadCache.flush(ThreadCache.java:136) ~[kafka-streams-3.1.1.jar!/:na]
2022-09-24T11:11:49.141766400Z  at org.apache.kafka.streams.state.internals.CachingKeyValueStore.flushCache(CachingKeyValueStore.java:345) ~[kafka-streams-3.1.1.jar!/:na]
2022-09-24T11:11:49.141770900Z  at org.apache.kafka.streams.state.internals.WrappedStateStore.flushCache(WrappedStateStore.java:71) ~[kafka-streams-3.1.1.jar!/:na]
2022-09-24T11:11:49.141775000Z  at org.apache.kafka.streams.processor.internals.ProcessorStateManager.flushCache(ProcessorStateManager.java:491) ~[kafka-streams-3.1.1.jar!/:na]
2022-09-24T11:11:49.141778800Z  at org.apache.kafka.streams.processor.internals.StreamTask.prepareCommit(StreamTask.java:402) ~[kafka-streams-3.1.1.jar!/:na]
2022-09-24T11:11:49.141782800Z  at org.apache.kafka.streams.processor.internals.TaskManager.commitTasksAndMaybeUpdateCommittableOffsets(TaskManager.java:1112) ~[kafka-streams-3.1.1.jar!/:na]
2022-09-24T11:11:49.141788800Z  at org.apache.kafka.streams.processor.internals.TaskManager.commit(TaskManager.java:1084) ~[kafka-streams-3.1.1.jar!/:na]
2022-09-24T11:11:49.141802700Z  at org.apache.kafka.streams.processor.internals.StreamThread.maybeCommit(StreamThread.java:1071) ~[kafka-streams-3.1.1.jar!/:na]
2022-09-24T11:11:49.141807200Z  at org.apache.kafka.streams.processor.internals.StreamThread.runOnce(StreamThread.java:817) ~[kafka-streams-3.1.1.jar!/:na]
2022-09-24T11:11:49.141811000Z  at org.apache.kafka.streams.processor.internals.StreamThread.runLoop(StreamThread.java:604) ~[kafka-streams-3.1.1.jar!/:na]
2022-09-24T11:11:49.141814100Z  at org.apache.kafka.streams.processor.internals.StreamThread.run(StreamThread.java:576) ~[kafka-streams-3.1.1.jar!/:na]
0 Answers
Related