【问题标题】:Kafka Consumer is not able to deserialize timewindowed key which has start and end timeKafka Consumer 无法反序列化具有开始和结束时间的时间窗口键
【发布时间】: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


    【解决方案1】:

    看来你在打https://issues.apache.org/jira/browse/KAFKA-7110

    在 2.2.0 中已修复,允许您将窗口大小传递给构造函数或TimeWindows

    请注意,窗口结束时间戳和窗口大小都不会存储在数据中。这是一种存储优化,因为对于TimeWindows,所有窗口的大小都是相同的,因此可以根据开始时间戳加上窗口大小来计算结束时间戳。

    【讨论】:

    • 谢谢您,与您的建议类似,我已经根据窗口开始时间 + 窗口大小实现了结束时间。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2022-12-08
    • 2012-05-05
    • 2011-11-27
    • 1970-01-01
    • 2016-09-24
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多