【发布时间】:2018-12-06 07:04:18
【问题描述】:
我的 Kafka Streams 应用程序正在使用使用以下键值布局的 kafka 主题:
String.class -> HistoryEvent.class
打印我当前的主题时,可以确认这一点:
bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic flow-event-stream-file-service-test-instance --property print.key=true --property key.separator=" -- " --from-beginning
flow1 -- SUCCESS #C:\Daten\file-service\in\crypto.p12
“flow1”是String键,--之后的部分是序列化值。
我的流程是这样设置的:
KStream<String, HistoryEvent> eventStream = builder.stream(applicationTopicName, Consumed.with(Serdes.String(),
historyEventSerde));
eventStream.selectKey((key, value) -> new HistoryEventKey(key, value.getIdentifier()))
.groupByKey()
.reduce((e1, e2) -> e2,
Materialized.<HistoryEventKey, HistoryEvent, KeyValueStore<Bytes, byte[]>>as(streamByKeyStoreName)
.withKeySerde(new HistoryEventKeySerde()));
据我所知,我告诉它使用 String 和 HistoryEvent serde 来使用主题,因为这是主题中的内容。然后我“重新设置”它以使用组合密钥,该组合密钥应使用为HistoryEventKey.class 提供的 serde 在本地存储。据我了解,这将导致使用新密钥创建一个额外的主题(可以在 kafka 容器中的主题列表中看到)。这很好。
现在的问题是应用程序无法启动,即使在干净的环境中只有一个主题中的文档:
org.apache.kafka.streams.errors.StreamsException: Exception caught in process. taskId=0_0, processor=KSTREAM-SOURCE-0000000000, topic=flow-event-stream-file-service-test-instance, partition=0, offset=0
Caused by: org.apache.kafka.streams.errors.StreamsException: A serializer (key: org.apache.kafka.common.serialization.StringSerializer / value: HistoryEventSerializer) is not compatible to the actual key or value type (key type: HistoryEventKey / value type: HistoryEvent). Change the default Serdes in StreamConfig or provide correct Serdes via method parameters.
从消息中很难判断问题到底出在哪里。它在我的基本主题中说,但这是不可能的,因为密钥不是HistoryEventKey 类型。由于我在reduce 中为HistoryEventKey 提供了一个serde,它也不能与本地商店一起使用。
对我来说唯一有意义的是它与导致重新排列和新主题的selectKey 操作有关。但是,我无法弄清楚如何为该操作提供 serde。我不想将其设置为默认值,因为它不是默认键 serde。
【问题讨论】:
标签: java apache-kafka apache-kafka-streams