【问题标题】:How to query the last value of a record in apache kafka topic如何在apache kafka主题中查询记录的最后一个值
【发布时间】:2019-07-07 00:19:27
【问题描述】:

我是 Kafka 的新手,所以可能很容易。但我无法从我现在面临的问题中看到任何解决方案。我有一个 Kafka 主题 metric_32,我想找到键 user_abc 的最新值。这在 Kafka 中是如何实现的。

我尝试使用KStream,但它会在有新事件涉及该主题时进行订阅。但我想要的是查询一个键的最后一个值。任何示例都会有所帮助。

【问题讨论】:

  • 如果您不知道该消息可能存在多长时间,则需要阅读整个主题(分区)。

标签: java apache-kafka


【解决方案1】:

您可以使用状态存储(如果您使用的是 Kafka 流),然后向其中添加一个处理器,该处理器会在每次向主题推送新值时更新状态存储。

builder.addGlobalStore(storeBuilder, topic, Consumed.with(keySerde, valueSerde), return new Processor<K,V>() {
    private KeyValueStore<K,V> store;

    public void init(ProcessorContext context) {
        store=(KeyValueStore<K,V>) context.getStateStore("statestorename");
    }

    public void process(K key, V value) {
        store.put(key,value);
    }

    public void close() {}
});

然后你就可以使用了

readOnlyStore=streams.store("statestorename", QueryableStoreTypes.keyValueStore());
readOnlyStore.get("key");

【讨论】:

    【解决方案2】:

    如果对于某个主题,您总是对特定键的最后一个值感兴趣,您可以设置 log.cleanup.policy=compact 。这样,您将始终最终每个键只有一条记录。如果您生成 5 条具有相同 id 的消息,则只有最后一条将保留在 kafka 中。这样,如果您有很多具有相同密钥的消息,您可以提高很多磁盘使用率。你可以在这里阅读更多:https://dzone.com/articles/kafka-architecture-log-compaction

    【讨论】:

    • 不幸的是,这不是日志压缩的作用。日志压缩保证至少存储给定键的最后一条记录。但是,这意味着您可能还会收到较旧的记录。 Kafka docs 明确表示:“日志压缩为我们提供了更细粒度的保留机制,因此我们保证至少保留每个主键的最后一次更新
    猜你喜欢
    • 2016-06-13
    • 1970-01-01
    • 1970-01-01
    • 2018-07-15
    • 1970-01-01
    • 1970-01-01
    • 2016-07-26
    • 2023-03-19
    • 2013-06-13
    相关资源
    最近更新 更多