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]