【发布时间】:2018-04-18 23:24:00
【问题描述】:
我无法使用文档 (https://docs.spring.io/spring-cloud-stream/docs/Elmhurst.RELEASE/reference/htmlsingle/#_configuration_options_3) 中指定的语法更改通道(或绑定)的 Serde。
假设我的频道是pcin,我知道我应该使用以下属性spring.cloud.stream.kafka.streams.bindings.pcin.producer.valueSerde 和spring.cloud.stream.kafka.streams.bindings.pcin.producer.keySerde 指定valueSerde 和keySerde。
但是,我收到一个异常:
Caused by: org.apache.kafka.streams.errors.StreamsException: A serializer (key: org.apache.kafka.common.serialization.StringSerializer / value: org.apache.kafka.common.serialization.StringSerializer) is not compatible to the actual key or value type (key type: java.lang.String / value type: java.lang.Long). Change the default Serdes in StreamConfig or provide correct Serdes via method parameters.
我正在尝试改编 Josh Long 的 Spring Tips 中的示例:https://github.com/spring-tips/spring-cloud-stream-kafka-streams
我刚刚将PageViewEventProcessor类改成如下:
@Component
public static class PageViewEventProcessor {
@StreamListener
@SendTo(AnalyticsBinding.PAGE_COUNT_OUT)
public KStream<String, Long> process(@Input(AnalyticsBinding.PAGE_VIEWS_IN) KStream<String, PageViewEvent> events) {
return events
.filter((key, value) -> value.getDuration() > 10)
.map((key, value) -> new KeyValue<>(value.getPage(), value.getDuration()))
.groupByKey()
.aggregate(()-> 0L,
(cle, val, valAgregee) -> valAgregee + val,
Materialized.as(AnalyticsBinding.PAGE_COUNT_MV))
.toStream();
}
}
我不计算事件(页面访问)的数量,而是计算每次访问的持续时间总和。
这是 application.properties 的摘录(来自 Spring 提示示例):
# page counts out
spring.cloud.stream.bindings.pcout.destination=pcs
spring.cloud.stream.bindings.pcout.producer.use-native-encoding=true
spring.cloud.stream.kafka.streams.bindings.pcout.producer.key-serde=org.apache.kafka.common.serialization.Serdes$StringSerde
spring.cloud.stream.kafka.streams.bindings.pcout.producer.value-serde=org.apache.kafka.common.serialization.Serdes$LongSerde
#
# page counts in
spring.cloud.stream.bindings.pcin.destination=pcs
spring.cloud.stream.bindings.pcin.consumer.use-native-decoding=true
spring.cloud.stream.bindings.pcin.group=pcs
spring.cloud.stream.bindings.pcin.content-type=application/json
spring.cloud.stream.bindings.pcin.consumer.header-mode=raw
spring.cloud.stream.kafka.streams.bindings.pcin.consumer.key-serde=org.apache.kafka.common.serialization.Serdes$StringSerde
spring.cloud.stream.kafka.streams.bindings.pcin.consumer.value-serde=org.apache.kafka.common.serialization.Serdes$LongSerde
还有其他必要的更改吗?
【问题讨论】: