【发布时间】:2019-05-09 13:14:02
【问题描述】:
我有一个 Kafka Streams 应用程序 (V 2.1.1),它生成记录并以键值格式输入输出主题。
key 是窗口时间 serde,我期望 key 和句柄到窗口开始/结束时间。
示例 -
.to(kafkaOutPutTopic, Produced.with(windowedSerde, jsonSerde));
示例 - [KEY@1551807076000/1551807077000] 其中 KEY 是键,开始时间 - 1551807076000 和结束时间 - 1551807077000
WindowedSerde 在哪里
StringSerializer stringSerializer = new StringSerializer();
final TimeWindowedSerializer<String> windowedSerializer = new TimeWindowedSerializer(stringSerializer);
final TimeWindowedDeserializer<String> windowedDeSerializer = new TimeWindowedDeserializer();
final Serde<Windowed<String>> windowedSerde = Serdes.serdeFrom(windowedSerializer, windowedDeSerializer);
还有一个名为 kafka 消费者的组件,它尝试使用来自主题的消息并通过自定义类的反序列化来获取键和窗口开始/结束时间。
kafka 属性:
kafkaConsumerProperties.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, TimeWindowedDeserializer.class.getName());
我正在使用附加链接中的 TimeWindowedDeserializer.java - https://gist.github.com/nfo/eaf350afb5667a3516593da4d48e757a
但启用获取窗口结束时间,并且消费者由于反序列化而无法使用它。
【问题讨论】:
标签: kafka-consumer-api apache-kafka-streams